mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(orchestration): stop worker receipts from contradicting themselves
worker-show spread the raw worker row beside its parsed copies, so a reader got residual_resources (a JSON string) next to residualResources (an array), plus host_scope as JSON-inside-JSON and two authority hashes with no consumer. Parse once, emit camelCase once, and withhold the hashes. worker-show also published only PTY liveness, so an agent that died at a trust prompt read live there while worker-list called it unverifiable -- and worker-list's nextAction pointed back at worker-show. Both now publish the same fleet projection. worker-list's projection.resource restated fields the row already carried, and the unconfirmed-stop sentence doubled a terminator on an already-punctuated reason.
This commit is contained in:
@@ -15,22 +15,37 @@ import { formatWorkerRead, type LegacyWorkerReadResult } from './worker-output'
|
||||
export const ORCHESTRATION_WORKER_OBSERVATION_HANDLERS: Record<string, CommandHandler> = {
|
||||
'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')
|
||||
})
|
||||
},
|
||||
|
||||
|
||||
@@ -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' })
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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<ReturnType<typeof inspectWorkerTerminal>>) {
|
||||
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<DispatchContextRow> & { 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
|
||||
|
||||
@@ -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: {
|
||||
|
||||
@@ -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.'
|
||||
)
|
||||
})
|
||||
})
|
||||
@@ -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)}`
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user