diff --git a/src/cli/handlers/orchestration/worker-observation-handlers.ts b/src/cli/handlers/orchestration/worker-observation-handlers.ts index 29841f7390f..018683fc285 100644 --- a/src/cli/handlers/orchestration/worker-observation-handlers.ts +++ b/src/cli/handlers/orchestration/worker-observation-handlers.ts @@ -15,22 +15,37 @@ import { formatWorkerRead, type LegacyWorkerReadResult } from './worker-output' export const ORCHESTRATION_WORKER_OBSERVATION_HANDLERS: Record = { 'orchestration worker-show': async ({ flags, client, json }) => { const result = await client.call<{ - dispatch: { id: string; task_id: string; status: string } - worker: { state: string; stage: string; agent_terminal_handle: string | null } + dispatch: { id: string; task_id: string; status: string } | null + worker: { state: string; stage: string; agentTerminalHandle: string | null } + projection?: { liveness: { verdict: string }; nextAction: { argv: string[] } } | null observation?: { agentWait?: { source: string; reason?: string } | null } }>('orchestration.workerShow', { dispatch: getRequiredStringFlag(flags, 'dispatch') }) printResult(result, json, (value) => { - const base = `${value.dispatch.id} task=${value.dispatch.task_id} [${value.worker.state}] stage=${value.worker.stage}` + const lines = [ + `${value.dispatch?.id ?? 'unknown'} task=${value.dispatch?.task_id ?? 'unknown'} [${value.worker.state}] stage=${value.worker.stage}` + ] + // Why: PTY status alone read `live` for an agent that died at a trust prompt, so the + // fleet verdict and its next action print beside it rather than in another command. + if (value.projection) { + lines.push( + `Agent liveness: ${value.projection.liveness.verdict}`, + `Next action: ${value.projection.nextAction.argv.join(' ') || 'none'}` + ) + } // Why: absent means unknown on older runtimes, distinct from an evaluated null wait. if (value.observation === undefined || !('agentWait' in value.observation)) { - return `${base}\nInteractive wait: unknown (not evaluated)` + lines.push('Interactive wait: unknown (not evaluated)') + } else if (value.observation.agentWait) { + const wait = value.observation.agentWait + lines.push( + `Waiting on a human: ${wait.reason ?? 'interactive prompt'} (via ${wait.source})` + ) + } else { + lines.push('Interactive wait: none') } - const wait = value.observation.agentWait - return wait - ? `${base}\nWaiting on a human: ${wait.reason ?? 'interactive prompt'} (via ${wait.source})` - : `${base}\nInteractive wait: none` + return lines.join('\n') }) }, diff --git a/src/main/runtime/rpc/methods/orchestration-federation-setup.test.ts b/src/main/runtime/rpc/methods/orchestration-federation-setup.test.ts index c1d91c18479..e275760c80b 100644 --- a/src/main/runtime/rpc/methods/orchestration-federation-setup.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-federation-setup.test.ts @@ -187,7 +187,7 @@ describe('orchestration federated setup evidence', () => { worker: { state: 'ready', stage: 'input_accepted', - setup_state: 'failed', + setupState: 'failed', effects: expect.arrayContaining([ expect.objectContaining({ kind: 'setup', state: 'failed' }), expect.objectContaining({ kind: 'dispatch_input', state: 'accepted' }) diff --git a/src/main/runtime/rpc/methods/orchestration-manual-dispatch-observation.test.ts b/src/main/runtime/rpc/methods/orchestration-manual-dispatch-observation.test.ts index c99fcb1c328..409eaafa836 100644 --- a/src/main/runtime/rpc/methods/orchestration-manual-dispatch-observation.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-manual-dispatch-observation.test.ts @@ -156,6 +156,7 @@ describe('manual Dispatch observation', () => { workerState: string terminalState: string | null agentTerminalHandle: string | null + projection: { liveness: { verdict: string } } }[] } expect(workerList.workers).toEqual([ @@ -167,12 +168,18 @@ describe('manual Dispatch observation', () => { }) ]) - await expect( - call('orchestration.workerShow', { dispatch: dispatch.id }) - ).resolves.toMatchObject({ - worker: { state: 'unsupervised', stage: 'injected', agent_terminal_handle: 'term_worker' }, + const workerShow = (await call('orchestration.workerShow', { + dispatch: dispatch.id + })) as { projection: { liveness: { verdict: string } } | null } + expect(workerShow).toMatchObject({ + worker: { state: 'unsupervised', stage: 'injected', agentTerminalHandle: 'term_worker' }, observation: { status: 'live', exactWorker: true } }) + // Why: worker-show published only PTY liveness, so it read `live` for a dispatch that + // worker-list called `unverifiable` — and worker-list's nextAction sent you back here. + expect(workerShow.projection?.liveness.verdict).toBe( + workerList.workers[0].projection.liveness.verdict + ) await expect( call('orchestration.workerRead', { dispatch: dispatch.id, source: 'terminal' }) ).resolves.toMatchObject({ diff --git a/src/main/runtime/rpc/methods/orchestration-worker-control.ts b/src/main/runtime/rpc/methods/orchestration-worker-control.ts index 9689d5ca0cc..7538237518c 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-control.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-control.ts @@ -6,9 +6,12 @@ import { defineMethod, type RpcMethod } from '../core' import { OptionalFiniteNumber, requiredString } from '../schemas' import { callFederatedWorkerShow, + exposeDispatchContext, exposeFederatedWorkerObservation, + exposeObservation, exposeWorker, inspectWorkerTerminal, + projectFleetWorker, resolvePinnedFederatedServer, showContextOnlyWorker } from './orchestration-worker-observation' @@ -118,8 +121,9 @@ export const ORCHESTRATION_WORKER_CONTROL_METHODS: RpcMethod[] = [ ) } return { - dispatch: db.getDispatchContextById(params.dispatch), + dispatch: exposeDispatchContext(db.getDispatchContextById(params.dispatch) ?? dispatch), worker: exposeWorker(worker), + projection: projectFleetWorker(runtime, db, params.dispatch), server: { environmentId: server.environmentId, name: server.name }, remoteRuntimeEpoch: db.getFederatedDispatch(params.dispatch)?.remote_runtime_epoch ?? @@ -148,19 +152,12 @@ export const ORCHESTRATION_WORKER_CONTROL_METHODS: RpcMethod[] = [ const observation = await inspectWorkerTerminal(runtime, db, params.dispatch) const resource = db.getWorkerTerminalResourceByOwner(params.dispatch) return { - dispatch, + dispatch: exposeDispatchContext(dispatch), worker: exposeWorker(worker), + // Why: the fleet verdict, so worker-show and worker-list cannot disagree. + projection: projectFleetWorker(runtime, db, params.dispatch), terminal: observation.exact ? observation.terminal : null, - observation: { - status: observation.status, - exactWorker: observation.exact, - // Why: a bare `unverifiable` is not actionable without naming what we lost. - ...(observation.reason ? { reason: observation.reason } : {}), - // Why conditional: a present null must mean "looked, nothing waiting". An - // unattached, missing or identity-changed worker was never looked at, and saying - // null there is the false negative this field exists to remove. - ...(observation.agentWait !== undefined ? { agentWait: observation.agentWait } : {}) - }, + observation: exposeObservation(observation), terminalResource: resource ? exposeWorkerTerminalResource(resource) : null } } diff --git a/src/main/runtime/rpc/methods/orchestration-worker-list-method.ts b/src/main/runtime/rpc/methods/orchestration-worker-list-method.ts index da5ed905c0a..033f2610b68 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-list-method.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-list-method.ts @@ -272,7 +272,15 @@ async function projectWorkerListPageWithFilteredSnapshot( agentTerminalHandle: row.agentTerminalHandle, terminalState: row.terminalState, resource: row.resource ? exposeWorkerTerminalResource(row.resource) : null, - projection + // Why: `projection.resource` restated id/ownerDispatchId/releaseState/terminalState + // that the row already carries; only the derived ownership classification is new. + projection: { + ...projection, + resource: + projection.resource.state === 'absent' + ? projection.resource + : { state: projection.resource.state } + } } }) const counts = Object.fromEntries( diff --git a/src/main/runtime/rpc/methods/orchestration-worker-observation.test.ts b/src/main/runtime/rpc/methods/orchestration-worker-observation.test.ts index c53572be1ae..df8e24836a0 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-observation.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-observation.test.ts @@ -1,7 +1,12 @@ import { describe, expect, it, vi } from 'vitest' import type { OrcaRuntimeService } from '../../orca-runtime' import type { OrchestrationDb } from '../../orchestration/db' -import { inspectWorkerTerminal } from './orchestration-worker-observation' +import { + exposeDispatchContext, + exposeWorker, + inspectWorkerTerminal +} from './orchestration-worker-observation' +import type { DispatchContextRow, WorkerDispatchRow } from '../../orchestration/types' const DISPATCH_ID = 'ctx-worker' const TERMINAL_HANDLE = 'term-worker' @@ -63,3 +68,58 @@ describe('inspectWorkerTerminal missing liveness verdict', () => { }) }) }) + +describe('worker-show receipt shape', () => { + it('parses the JSON columns once and emits one casing', () => { + const exposed = exposeWorker({ + dispatch_id: DISPATCH_ID, + runtime_epoch: 'epoch-1', + state: 'ready', + stage: 'input_accepted', + worktree_id: 'repo::/tmp/wt', + agent_terminal_handle: TERMINAL_HANDLE, + setup_state: 'ran', + effects: '[{"kind":"setup"}]', + residual_resources: '["res-1"]', + start_options: '{"agent":"codex"}', + last_error: null, + created_at: 'now', + updated_at: 'now' + } as WorkerDispatchRow) + + expect(exposed).toEqual({ + dispatchId: DISPATCH_ID, + runtimeEpoch: 'epoch-1', + state: 'ready', + stage: 'input_accepted', + worktreeId: 'repo::/tmp/wt', + agentTerminalHandle: TERMINAL_HANDLE, + setupState: 'ran', + effects: [{ kind: 'setup' }], + residualResources: ['res-1'], + startOptions: { agent: 'codex' }, + lastError: null, + createdAt: 'now', + updatedAt: 'now' + }) + }) + + it('parses host_scope and withholds authority hashes from the dispatch row', () => { + const exposed = exposeDispatchContext({ + id: DISPATCH_ID, + run_id: 'run-1', + task_id: 'task-1', + launch_token_hash: 'launch-secret', + capability_hash: 'capability-secret', + host_scope: JSON.stringify({ kind: 'local', hostId: 'local' }) + } as DispatchContextRow) + + expect(exposed).toMatchObject({ + id: DISPATCH_ID, + hostScope: { kind: 'local', hostId: 'local' } + }) + expect(exposed).not.toHaveProperty('host_scope') + expect(exposed).not.toHaveProperty('launch_token_hash') + expect(exposed).not.toHaveProperty('capability_hash') + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration-worker-observation.ts b/src/main/runtime/rpc/methods/orchestration-worker-observation.ts index d10c5071b99..c7386d28eaf 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-observation.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-observation.ts @@ -3,6 +3,8 @@ import type { OrcaRuntimeService } from '../../orca-runtime' import type { OrchestrationDb } from '../../orchestration/db' import { OrchestrationError } from '../../orchestration/orchestration-error' import { parseWorkerTerminalHostScope } from '../../orchestration/worker-terminal-process-liveness' +import type { OrchestrationFleetWorker } from '../../../../shared/orchestration-fleet-projection' +import { projectWorkerFleet } from './orchestration-worker-list-projection' import type { DispatchContextRow, FederatedDispatchRow, @@ -83,24 +85,46 @@ export async function inspectWorkerTerminal( } } +/** Why conditional: a present `agentWait: null` must mean "looked, nothing waiting"; an + * unattached, missing or identity-changed worker was never looked at, and a bare + * `unverifiable` is not actionable without naming what contact was lost. */ +export function exposeObservation(observation: Awaited>) { + return { + status: observation.status, + exactWorker: observation.exact, + ...(observation.reason ? { reason: observation.reason } : {}), + ...(observation.agentWait !== undefined ? { agentWait: observation.agentWait } : {}) + } +} + export function exposeContextOnlyWorker(dispatch: DispatchContextRow) { return { - dispatch_id: dispatch.id, - runtime_epoch: null, + dispatchId: dispatch.id, + runtimeEpoch: null, state: 'unsupervised' as const, stage: dispatch.capability_hash ? 'injected' : 'context_only', - worktree_id: null, - agent_terminal_handle: dispatch.assignee_handle, - setup_state: 'not_applicable', - effects: [], - residualResources: [], - startOptions: {}, - last_error: dispatch.last_failure, - created_at: dispatch.created_at, - updated_at: dispatch.completed_at ?? dispatch.created_at + worktreeId: null, + agentTerminalHandle: dispatch.assignee_handle, + setupState: 'not_applicable', + effects: [] as unknown[], + residualResources: [] as unknown[], + startOptions: {} as unknown, + lastError: dispatch.last_failure, + createdAt: dispatch.created_at, + updatedAt: dispatch.completed_at ?? dispatch.created_at } } +// Why: `launch_token_hash` and `capability_hash` are authority material with no receipt +// consumer, and `host_scope` shipped as a JSON string inside JSON. One camelCase shape. +export function exposeDispatchContext(dispatch: DispatchContextRow) { + const exposed: Partial & { hostScope?: unknown } = { ...dispatch } + delete exposed.launch_token_hash + delete exposed.capability_hash + delete exposed.host_scope + return { ...exposed, hostScope: parseWorkerTerminalHostScope(dispatch.host_scope) } +} + export async function showContextOnlyWorker( runtime: OrcaRuntimeService, db: OrchestrationDb, @@ -108,28 +132,64 @@ export async function showContextOnlyWorker( ) { const observation = await inspectWorkerTerminal(runtime, db, dispatch.id) return { - dispatch, + dispatch: exposeDispatchContext(dispatch), worker: exposeContextOnlyWorker(dispatch), + projection: projectFleetWorker(runtime, db, dispatch.id), terminal: observation.exact ? observation.terminal : null, - observation: { - status: observation.status, - exactWorker: observation.exact, - ...(observation.reason ? { reason: observation.reason } : {}), - ...(observation.agentWait !== undefined ? { agentWait: observation.agentWait } : {}) - }, + observation: exposeObservation(observation), terminalResource: null } } +// Why: the row was spread verbatim beside its parsed copies, so a reader got +// `residual_resources` (a JSON string) next to `residualResources` (an array) and had to +// guess which was authoritative. Parse once, emit camelCase once. export function exposeWorker(worker: WorkerDispatchRow) { return { - ...worker, + dispatchId: worker.dispatch_id, + runtimeEpoch: worker.runtime_epoch, + state: worker.state, + stage: worker.stage, + worktreeId: worker.worktree_id, + agentTerminalHandle: worker.agent_terminal_handle, + setupState: worker.setup_state, effects: JSON.parse(worker.effects) as unknown[], residualResources: JSON.parse(worker.residual_resources) as unknown[], - startOptions: JSON.parse(worker.start_options) as unknown + startOptions: JSON.parse(worker.start_options) as unknown, + lastError: worker.last_error, + createdAt: worker.created_at, + updatedAt: worker.updated_at } } +/** + * The same fleet verdict `worker-list` publishes, for one Dispatch. + * + * Why worker-show needs it: `observation.status` is PTY liveness, so an agent that died + * at a trust prompt inside a live pane read `live` here and `unverifiable` from + * `worker-list` — and `worker-list`'s own `nextAction` pointed back at this command. + */ +export function projectFleetWorker( + runtime: OrcaRuntimeService, + db: OrchestrationDb, + dispatchId: string +): OrchestrationFleetWorker | null { + const rows = db.listWorkerTerminalResources({ dispatchIds: [dispatchId], limit: 1 }) + if (rows.length === 0) { + return null + } + const now = Date.now() + return ( + projectWorkerFleet({ + rows, + attentionFacts: db.getWorkerAttentionFactsForDispatches([dispatchId], now), + statuses: runtime.getOrchestrationFleetAgentStatusSnapshot(), + limit: 1, + now + }).workers[0] ?? null + ) +} + export function exposeFederatedWorkerObservation( observation: { status?: string; exactWorker: boolean; reason?: string }, projected: boolean diff --git a/src/main/runtime/rpc/methods/orchestration-workers-recovery.test.ts b/src/main/runtime/rpc/methods/orchestration-workers-recovery.test.ts index 8c37db516cc..361368ca2e2 100644 --- a/src/main/runtime/rpc/methods/orchestration-workers-recovery.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-workers-recovery.test.ts @@ -290,7 +290,7 @@ describe('orchestration worker recovery', () => { await expect( call('orchestration.workerShow', { dispatch: started.dispatch.id }) ).resolves.toMatchObject({ - worker: { state: 'stopped', stage: 'process_stopped', last_error: null }, + worker: { state: 'stopped', stage: 'process_stopped', lastError: null }, observation: { status: 'exited', exactWorker: true } }) expect(db.getTask(task.id)?.status).toBe('blocked') @@ -370,7 +370,7 @@ describe('orchestration worker recovery', () => { }) await expect(show).resolves.toMatchObject({ - worker: { stage: 'released', agent_terminal_handle: null }, + worker: { stage: 'released', agentTerminalHandle: null }, remoteRuntimeEpoch: 'windows_epoch_new', terminal: null, observation: { diff --git a/src/shared/pty-liveness-verdict.test.ts b/src/shared/pty-liveness-verdict.test.ts new file mode 100644 index 00000000000..48d7afe84c8 --- /dev/null +++ b/src/shared/pty-liveness-verdict.test.ts @@ -0,0 +1,28 @@ +import { describe, expect, it } from 'vitest' +import { describeUnconfirmedAgentStop, describeUnconfirmedStop } from './pty-liveness-verdict' + +describe('unconfirmed-stop sentences', () => { + it('terminates a reason that has no terminator', () => { + expect(describeUnconfirmedStop('its SSH provider is no longer registered')).toBe( + 'The PTY was not confirmed stopped: its SSH provider is no longer registered.' + ) + }) + + it('does not double the terminator on a reason that is already a sentence', () => { + // A relayed lifecycle_conflict message arrives punctuated and printed `...to failed..`. + expect( + describeUnconfirmedAgentStop({ + ptyStopVerdict: 'unverifiable', + ptyStopReason: 'worker w1 cannot transition from stopping to failed.' + }) + ).toBe( + 'The agent terminal was closed but its process could not be confirmed stopped: worker w1 cannot transition from stopping to failed.' + ) + }) + + it('still terminates the live-process wording', () => { + expect(describeUnconfirmedAgentStop({ ptyStopVerdict: 'live' })).toBe( + 'The agent terminal was closed but its process could not be confirmed stopped: it is live.' + ) + }) +}) diff --git a/src/shared/pty-liveness-verdict.ts b/src/shared/pty-liveness-verdict.ts index f16e756ca9a..0dd650161f3 100644 --- a/src/shared/pty-liveness-verdict.ts +++ b/src/shared/pty-liveness-verdict.ts @@ -16,9 +16,15 @@ export const NO_OBSERVING_PROVIDER_REASON = 'no registered provider can observe export const SSH_EXIT_UNCONFIRMED_REASON = 'the owning SSH host did not confirm the PTY exit' export const PTY_LIVE_NOTE = 'The PTY is live.' +// Why: reasons reach these sentences from verdicts, receipts and relayed errors, and +// some already end in a terminator — appending one blindly printed `...to failed..`. +function endSentence(detail: string): string { + return /[.!?]$/u.test(detail.trimEnd()) ? detail.trimEnd() : `${detail.trimEnd()}.` +} + /** The one sentence every surface uses to admit a stop was not confirmed. */ export function describeUnconfirmedStop(reason: string): string { - return `The PTY was not confirmed stopped: ${reason}.` + return `The PTY was not confirmed stopped: ${endSentence(reason)}` } /** Words a close whose PTY teardown was never confirmed, for a stop receipt. */ @@ -30,5 +36,5 @@ export function describeUnconfirmedAgentStop(close: { close.ptyStopVerdict === 'live' ? 'it is live' : (close.ptyStopReason ?? 'the stop outcome could not be verified') - return `The agent terminal was closed but its process could not be confirmed stopped: ${detail}.` + return `The agent terminal was closed but its process could not be confirmed stopped: ${endSentence(detail)}` }