fix(orchestration): worker-start settles readiness on observed turn start, not write acceptance

A dispatched PTY worker whose agent wedged at startup (six codex workers on
2026-09-07) was reported 'ok: true, state: ready, stage: input_accepted': the
preamble write was acknowledged with observationTimeoutMs: 0 and nothing ever
verified a turn began. The corpse and the healthy worker produced identical
receipts.

worker-start now runs the existing second-stage prompt observer
(observeTerminalAgentPrompt) after acceptance, inside the 30s window the
client RPC grace already budgets for (orchestration-worker-start-prompt-budget):

- turn observed (or provider ack for structured sessions) -> ready
- permission prompt -> ready; positive liveness, surfaced in the receipt
- provider without a turn-start signal -> ready; observation: unsupported
- observation supported and nothing started -> worker state start_unknown,
  response state outcome_unknown with nextCommands. Honest 'unverifiable',
  never a death claim: the capability and terminal are kept, and
  worker-report settlement already reconnects a start_unknown worker that
  recovers and reports.

Also fixes the effect-verb lie that misdirected the first diagnosis of this
incident: agent-first worktree creation labeled its own brand-new agent
terminal 'reused_agent_terminal' (a role test picking a lifecycle verb) on
both the local and federation paths. It now says 'created'; readers keep
accepting the retired verb for rows persisted before the rename.
This commit is contained in:
Merge Sim
2026-09-07 18:01:50 -07:00
parent 588240043e
commit 325fa88b09
12 changed files with 399 additions and 55 deletions
@@ -15,9 +15,9 @@ export function prepareStartingWorkerAuthority(
effects: unknown[]
setupState: string
hostScope?: string | null
// 'created': this worker-start operation created the agent terminal (including agent-first
// worktree creation, whose effects receipt says 'reused_agent_terminal'). 'external': an
// explicit --terminal reuse; ownership transfers only from an exact owned settled resource.
// 'created': this worker-start operation created the agent terminal (agent-first worktree
// creation included; its pre-rename effects rows said 'reused_agent_terminal'). 'external':
// an explicit --terminal reuse; ownership transfers only from an exact owned settled resource.
terminalOwnership?: 'created' | 'external'
}
): string {
@@ -121,7 +121,8 @@ export function markWorkerStartUnknown(
this: OrchestrationDb,
dispatchId: string,
stage: string,
reason: string
reason: string,
effects?: unknown[]
): WorkerDispatchRow {
this.db.exec('BEGIN IMMEDIATE')
try {
@@ -135,7 +136,12 @@ export function markWorkerStartUnknown(
id: dispatchId,
from: 'starting',
to: 'start_unknown',
projection: { stage, last_error: reason, updated_at: new Date().toISOString() }
projection: {
stage,
last_error: reason,
updated_at: new Date().toISOString(),
...(effects ? { effects: JSON.stringify(effects) } : {})
}
})
transitionLifecycleWithDb(this.db, {
entity: 'dispatch',
@@ -29,7 +29,9 @@ export function appendFederationTerminalEffects(
: terminal.handle === setupHandle
? 'setup'
: 'configured_tab',
action: terminal.handle === agentHandle ? 'reused_agent_terminal' : 'created',
// Remote agent-first worktree creation made every listed terminal, the agent one included;
// the verb is a lifecycle fact, not a role marker.
action: 'created',
id: terminal.handle,
tabId: terminal.tabId,
leafId: terminal.leafId
@@ -572,7 +572,7 @@ describe('orchestration RPC methods', () => {
)
expect(result.effects).toEqual(
expect.arrayContaining([
expect.objectContaining({ role: 'agent', action: 'reused_agent_terminal' }),
expect.objectContaining({ role: 'agent', action: 'created' }),
expect.objectContaining({ role: 'setup', action: 'created' }),
expect.objectContaining({ role: 'configured_tab', action: 'created' })
])
@@ -18,18 +18,17 @@ import {
import { failWorkerStartWithReceipt } from './worker-start-receipt'
import { parseTaskDeps } from './task-deps-argument'
import { assertExplicitWorkerTerminalUsable } from './explicit-worker-terminal-validation'
import { deliverWorkerDispatchPreamble } from './deliver-worker-dispatch-preamble'
import { tearDownFailedWorkerStart } from './failed-worker-start-teardown'
import {
createExistingWorktreeWorkerTerminal,
createStructuredWorkerSessionForWorktree,
createWorkerWorktree,
monitorWorkerSetup,
requireWorkerAuthority,
type WorkerEffect,
type WorkerSetupReceipt
} from './worker-topology'
import { prepareLocalWorkerStart } from './worker-start-validation'
import { deliverAndSettleWorkerStartReadiness } from './worker-start-readiness-settlement'
type WorkerStartMutation = {
callerFingerprint: string
@@ -240,50 +239,30 @@ export async function startLocalWorker(args: {
terminalOwnership: params.terminal ? 'external' : 'created'
})
failedStage = 'dispatch_input'
const promptDelivery = await deliverWorkerDispatchPreamble({
return await deliverAndSettleWorkerStartReadiness({
runtime,
structuredSession,
terminalHandle,
db,
run,
task,
dispatchId: started.dispatch.id,
dispatchDepth: started.dispatch.depth,
taskId: task.id,
taskSpec: task.spec,
structuredSession,
terminalHandle,
coordinatorHandle: params.from,
dispatchCapability: capability,
devMode: params.devMode,
requestId: orchestrationMutation?.requestId ?? started.dispatch.id
})
effects.push({
kind: 'dispatch_input',
role: 'agent',
id: terminalHandle,
state: 'accepted'
})
const worker = db.markWorkerDispatchReady(started.dispatch.id, effects)
monitorWorkerSetup({
runtime,
db,
runId: run.id,
dispatchId: started.dispatch.id,
requestId: orchestrationMutation?.requestId ?? started.dispatch.id,
agent: agent ?? null,
setupReceipt,
effects
})
return {
runId: run.id,
taskId: task.id,
dispatchId: started.dispatch.id,
state: worker.state,
stage: worker.stage,
setup: setupReceipt,
launch: launch.receipt,
launchReceipt: launch.receipt,
mode,
timeoutMs: params.timeoutMs ?? 60_000,
effects,
...(promptDelivery ? { prompt: promptDelivery } : {}),
residualResources: [],
...(terminalRevealWarning ? { warning: terminalRevealWarning } : {})
}
terminalRevealWarning,
onStage: (stage) => {
failedStage = stage
}
})
} catch (error) {
const residualAgentTerminal = await tearDownFailedWorkerStart({
runtime,
@@ -6,6 +6,8 @@ import {
} from './worker-topology'
function residualWorkerEffects(effects: WorkerEffect[]): WorkerEffect[] {
// 'reused_agent_terminal' is the retired verb agent-first creation used for its own agent
// terminal; rows persisted before the rename still carry it.
return effects.filter(
(effect) => effect.action?.startsWith('created') || effect.action === 'reused_agent_terminal'
)
@@ -227,22 +227,26 @@ describe('orchestration worker-start prompt contract', () => {
})
})
it('keeps a swallowed Enter queued without revoking the worker or retrying input', async () => {
it('reports a swallowed Enter as start_unknown while keeping the worker and its capability', async () => {
vi.useFakeTimers()
const harness = await createPromptContractHarness('swallowed')
const pending = harness.dispatcher.dispatch(harness.request)
await vi.runAllTimersAsync()
const response = await pending
// Codex supports turn-start observation and no turn started, so ready would be a lie: the
// paste can sit unsent in the composer while the receipt looks like a healthy dispatch.
expect(response).toMatchObject({
ok: true,
result: {
state: 'ready',
stage: 'input_accepted',
state: 'outcome_unknown',
stage: 'turn_start_unobserved',
turnStart: 'unobserved',
prompt: {
requestId: harness.requestId,
stages: ['input_accepted']
},
nextCommands: expect.arrayContaining([expect.stringContaining('worker-show')]),
mutation: { requestId: harness.requestId, replayed: false }
}
})
@@ -251,29 +255,42 @@ describe('orchestration worker-start prompt contract', () => {
}
const dispatchId = (response.result as { dispatchId: string }).dispatchId
await vi.advanceTimersByTimeAsync(20_000)
// Unverifiable is not failure: exactly one submit, no blind retry, nothing torn down.
expect(harness.submittedTurns()).toBe(1)
expect(harness.startedTurns()).toBe(0)
expect(harness.prematureSubmits()).toBe(0)
expect(harness.writes.filter((data) => data === '\r')).toHaveLength(1)
const persisted = reopenPromptContractDb(harness)
expect(persisted.getTask(harness.taskId)?.status).toBe('dispatched')
expect(persisted.getTask(harness.taskId)?.status).toBe('blocked')
expect(persisted.getDispatchContextById(dispatchId)).toMatchObject({
status: 'dispatched',
status: 'pending',
last_failure: null,
// The capability survives so a worker that recovers can still report; worker-report
// settlement reconnects a start_unknown worker through 'ready'.
capability_hash: expect.any(String),
capability_revoked_at: null
})
expect(persisted.getWorkerDispatch(dispatchId)).toMatchObject({
state: 'ready',
stage: 'input_accepted',
last_error: null
state: 'start_unknown',
stage: 'turn_start_unobserved',
last_error: expect.stringContaining('never started a turn')
})
const persistedEffects = JSON.parse(
persisted.getWorkerDispatch(dispatchId)?.effects ?? '[]'
) as { kind?: string; state?: string }[]
expect(persistedEffects).toEqual(
expect.arrayContaining([
expect.objectContaining({ kind: 'dispatch_input', state: 'accepted' }),
expect.objectContaining({ kind: 'dispatch_input', state: 'turn_unobserved' })
])
)
const callerFingerprint = persisted.getOrCreateLocalMutationCallerFingerprint()
const receipt = persisted.getMutationReceipt(callerFingerprint, harness.requestId)
expect(receipt).toMatchObject({ state: 'completed' })
expect(JSON.parse(receipt?.receipt ?? 'null')).toMatchObject({
dispatchId,
state: 'ready',
stage: 'input_accepted',
state: 'outcome_unknown',
stage: 'turn_start_unobserved',
prompt: {
requestId: harness.requestId,
stages: ['input_accepted']
@@ -0,0 +1,145 @@
import type { OrcaRuntimeService } from '../../../../orca-runtime'
import type { OrchestrationDb } from '../../../../orchestration/db'
import type { RunRow, TaskRow } from '../../../../orchestration/types'
import type { WorkerStartModeReceipt } from '../../orchestration-worker-start-mode'
import { deliverWorkerDispatchPreamble } from './deliver-worker-dispatch-preamble'
import type { OrchestrationWorkerLaunchReceipt } from './worker-launch-preferences'
import {
describeUnobservedWorkerTurnStart,
observeWorkerTurnStart,
type WorkerTurnStartObservation
} from './worker-start-turn-observation'
import {
monitorWorkerSetup,
type createStructuredWorkerSessionForWorktree,
type WorkerEffect,
type WorkerSetupReceipt
} from './worker-topology'
/**
* Delivers the dispatch preamble and settles the worker's start state on the strongest
* evidence available: `ready` only with a positive turn-start (or a provider that cannot
* prove one), `start_unknown` when observation is supported and nothing started.
*/
export async function deliverAndSettleWorkerStartReadiness(args: {
runtime: OrcaRuntimeService
db: OrchestrationDb
run: RunRow
task: TaskRow
dispatchId: string
dispatchDepth: number
structuredSession: Awaited<ReturnType<typeof createStructuredWorkerSessionForWorktree>> | null
terminalHandle: string
coordinatorHandle: string
dispatchCapability: string
devMode: boolean | undefined
requestId: string
agent: string | null
setupReceipt: WorkerSetupReceipt
launchReceipt: OrchestrationWorkerLaunchReceipt
mode: WorkerStartModeReceipt
timeoutMs: number
effects: WorkerEffect[]
terminalRevealWarning: string | undefined
/** Keeps the caller's failure receipt naming the stage that actually failed. */
onStage: (stage: 'dispatch_input' | 'turn_observation') => void
}): Promise<unknown> {
const { runtime, db, run, task, structuredSession, terminalHandle, effects } = args
args.onStage('dispatch_input')
const promptDelivery = await deliverWorkerDispatchPreamble({
runtime,
structuredSession,
terminalHandle,
dispatchId: args.dispatchId,
dispatchDepth: args.dispatchDepth,
taskId: task.id,
taskSpec: task.spec,
coordinatorHandle: args.coordinatorHandle,
dispatchCapability: args.dispatchCapability,
devMode: args.devMode,
requestId: args.requestId
})
effects.push({
kind: 'dispatch_input',
role: 'agent',
id: terminalHandle,
state: 'accepted'
})
args.onStage('turn_observation')
// The write above was accepted without waiting on provider hooks; now demand the positive
// evidence the receipt claims is observable. A worker whose turn never starts must not be
// reported ready — a wedged agent and a working one looked identical before this gate.
// A structured preamble send is acknowledged by the provider or throws, so it is already
// positive evidence.
const turnStart: WorkerTurnStartObservation = structuredSession
? { verdict: 'observed' }
: await observeWorkerTurnStart({ runtime, terminalHandle, prompt: promptDelivery })
const deliveredPrompt = turnStart.prompt ?? promptDelivery
monitorWorkerSetup({
runtime,
db,
runId: run.id,
dispatchId: args.dispatchId,
setupReceipt: args.setupReceipt,
effects
})
if (turnStart.verdict === 'unobserved') {
// Honest `unverifiable`: keep the dispatch capability and the terminal — the worker may
// still recover and report (worker-report settlement reconnects a start_unknown worker) —
// but never claim ready for a turn nobody observed.
effects.push({
kind: 'dispatch_input',
role: 'agent',
id: terminalHandle,
state: 'turn_unobserved'
})
const reason = describeUnobservedWorkerTurnStart(args.agent)
const worker = db.markWorkerStartUnknown(
args.dispatchId,
'turn_start_unobserved',
reason,
effects
)
return {
runId: run.id,
taskId: task.id,
dispatchId: args.dispatchId,
state: 'outcome_unknown',
stage: worker.stage,
turnStart: turnStart.verdict,
lastError: reason,
setup: args.setupReceipt,
launch: args.launchReceipt,
mode: args.mode,
timeoutMs: args.timeoutMs,
effects,
...(deliveredPrompt ? { prompt: deliveredPrompt } : {}),
residualResources: JSON.parse(worker.residual_resources) as unknown[],
nextCommands: [
`orca orchestration worker-show --dispatch ${args.dispatchId} --json`,
`orca terminal read --terminal ${terminalHandle} --screen`,
`orca orchestration worker-abandon --dispatch ${args.dispatchId} --json`
],
...(args.terminalRevealWarning ? { warning: args.terminalRevealWarning } : {})
}
}
const worker = db.markWorkerDispatchReady(args.dispatchId, effects)
return {
runId: run.id,
taskId: task.id,
dispatchId: args.dispatchId,
state: worker.state,
stage: worker.stage,
turnStart: turnStart.verdict,
setup: args.setupReceipt,
launch: args.launchReceipt,
mode: args.mode,
timeoutMs: args.timeoutMs,
effects,
...(deliveredPrompt ? { prompt: deliveredPrompt } : {}),
residualResources: [],
...(args.terminalRevealWarning ? { warning: args.terminalRevealWarning } : {})
}
}
@@ -0,0 +1,110 @@
import { describe, expect, it, vi } from 'vitest'
import type { RuntimeTerminalPromptDelivery } from '../../../../../../shared/runtime-terminal-contracts'
import type { OrcaRuntimeService } from '../../../../orca-runtime'
import { observeWorkerTurnStart } from './worker-start-turn-observation'
function delivery(
overrides: Partial<RuntimeTerminalPromptDelivery> = {}
): RuntimeTerminalPromptDelivery {
return {
requestId: 'req-1',
stages: ['input_accepted'],
provider: 'codex',
observation: 'supported',
processIncarnation: 'inc-1',
generation: 1,
baselineWorkingSequence: 0,
...overrides
}
}
function runtimeObserving(result: RuntimeTerminalPromptDelivery): {
runtime: OrcaRuntimeService
observe: ReturnType<typeof vi.fn>
} {
const observe = vi.fn().mockResolvedValue(result)
return {
runtime: { observeTerminalAgentPrompt: observe } as unknown as OrcaRuntimeService,
observe
}
}
describe('observeWorkerTurnStart', () => {
it('treats a missing receipt as unsupported observation, never as failure', async () => {
const { runtime, observe } = runtimeObserving(delivery())
await expect(
observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt: undefined })
).resolves.toEqual({ verdict: 'unsupported' })
expect(observe).not.toHaveBeenCalled()
})
it('accepts a first-stage turn_started without a second observation pass', async () => {
const prompt = delivery({ stages: ['input_accepted', 'turn_started'] })
const { runtime, observe } = runtimeObserving(prompt)
await expect(
observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt })
).resolves.toEqual({ verdict: 'observed', prompt })
expect(observe).not.toHaveBeenCalled()
})
it('reports observed when the second-stage observer sees the turn start', async () => {
const observed = delivery({ stages: ['input_accepted', 'turn_started'] })
const { runtime, observe } = runtimeObserving(observed)
await expect(
observeWorkerTurnStart({
runtime,
terminalHandle: 'term_w',
prompt: delivery(),
timeoutMs: 5
})
).resolves.toEqual({ verdict: 'observed', prompt: observed })
expect(observe).toHaveBeenCalledWith('term_w', delivery(), 5)
})
it('reports unobserved — not dead — when a supported observation stalls', async () => {
const { runtime } = runtimeObserving(delivery())
await expect(
observeWorkerTurnStart({
runtime,
terminalHandle: 'term_w',
prompt: delivery(),
timeoutMs: 5
})
).resolves.toMatchObject({ verdict: 'unobserved' })
})
it('reports permission as positive liveness', async () => {
const observed = delivery({ observation: 'permission' })
const { runtime } = runtimeObserving(observed)
await expect(
observeWorkerTurnStart({
runtime,
terminalHandle: 'term_w',
prompt: delivery(),
timeoutMs: 5
})
).resolves.toEqual({ verdict: 'permission', prompt: observed })
})
it('treats a replaced incarnation as unobserved rather than unsupported', async () => {
const observed = delivery({ observation: 'incarnation_replaced' })
const { runtime } = runtimeObserving(observed)
await expect(
observeWorkerTurnStart({
runtime,
terminalHandle: 'term_w',
prompt: delivery(),
timeoutMs: 5
})
).resolves.toEqual({ verdict: 'unobserved', prompt: observed })
})
it('leaves an unsupported provider on the accepted receipt', async () => {
const prompt = delivery({ provider: 'unsupported', observation: 'unsupported' })
const { runtime, observe } = runtimeObserving(prompt)
await expect(
observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt })
).resolves.toEqual({ verdict: 'unsupported', prompt })
expect(observe).not.toHaveBeenCalled()
})
})
@@ -0,0 +1,80 @@
import { AGENT_PROMPT_EFFECT_TIMEOUT_MS } from '../../../../../../shared/orchestration-timing-budgets'
import type { RuntimeTerminalPromptDelivery } from '../../../../../../shared/runtime-terminal-contracts'
import type { OrcaRuntimeService } from '../../../../orca-runtime'
/**
* Turn-start verdict for a dispatched worker prompt, in the execution-boundary vocabulary:
*
* - 'observed': the provider proved a turn started for this request. Positive liveness.
* - 'permission': the agent rendered an approval prompt after the write. Positive liveness,
* but the turn is blocked on a human.
* - 'unsupported': this provider exposes no turn-start signal; the accepted write is the
* strongest receipt that can exist. Never treated as failure.
* - 'unobserved': observation IS supported and no turn started within the window. This is
* `unverifiable`, never evidence of death — the bytes were written, but the agent may be
* wedged at startup or holding the spec unsent in its composer.
*/
export type WorkerTurnStartVerdict = 'observed' | 'permission' | 'unsupported' | 'unobserved'
export type WorkerTurnStartObservation = {
verdict: WorkerTurnStartVerdict
prompt?: RuntimeTerminalPromptDelivery
}
function classifyPromptDelivery(prompt: RuntimeTerminalPromptDelivery): WorkerTurnStartVerdict {
if (prompt.stages.includes('turn_started')) {
return 'observed'
}
if (prompt.observation === 'permission') {
return 'permission'
}
if (prompt.observation === 'supported') {
return 'unobserved'
}
// 'unsupported' (and an old host's missing observation) leaves acceptance as the best receipt.
return 'unsupported'
}
/**
* Second-stage turn-start observation for a worker prompt that was accepted without waiting on
* provider hooks. Reuses the same observer that terminal.send receipts replay through, so the
* evidence rules (lifecycle edge or hook turn-start, never output bytes) stay in one place.
*
* The observation window is `AGENT_PROMPT_EFFECT_TIMEOUT_MS`, which worker-start's client RPC
* grace already budgets for (see orchestration-worker-start-prompt-budget.ts).
*/
export async function observeWorkerTurnStart(args: {
runtime: OrcaRuntimeService
terminalHandle: string
prompt: RuntimeTerminalPromptDelivery | undefined
timeoutMs?: number
}): Promise<WorkerTurnStartObservation> {
if (!args.prompt) {
return { verdict: 'unsupported' }
}
const verdict = classifyPromptDelivery(args.prompt)
if (verdict !== 'unobserved') {
return { verdict, prompt: args.prompt }
}
const observed = await args.runtime.observeTerminalAgentPrompt(
args.terminalHandle,
args.prompt,
args.timeoutMs ?? AGENT_PROMPT_EFFECT_TIMEOUT_MS
)
if (observed.observation === 'incarnation_replaced') {
// The PTY under this handle changed mid-observation; the accepted write is unproven.
return { verdict: 'unobserved', prompt: observed }
}
return { verdict: classifyPromptDelivery(observed), prompt: observed }
}
export function describeUnobservedWorkerTurnStart(agent: string | null): string {
const name = agent ?? 'the agent'
return (
`Dispatch input was written and submitted, but ${name} never started a turn within ` +
`${Math.round(AGENT_PROMPT_EFFECT_TIMEOUT_MS / 1000)}s. This is unverifiable, not proof the ` +
'worker is dead: the agent may still be starting, may be wedged (for example waiting on ' +
'network), or may be holding the task unsent in its composer. If the worker recovers and ' +
'reports, this Dispatch settles normally.'
)
}
@@ -228,7 +228,10 @@ export async function createWorkerWorktree(args: {
: terminal.handle === setupTerminalHandle
? 'setup'
: 'configured_tab',
action: terminal.handle === terminalHandle ? 'reused_agent_terminal' : 'created',
// Every terminal listed here — the agent terminal included — was created by this call's
// agent-first worktree creation. The old 'reused_agent_terminal' verb on the agent row
// conflated the role test with a lifecycle claim and misdiagnosed at least one incident.
action: 'created',
id: terminal.handle,
tabId: terminal.tabId,
leafId: terminal.leafId
@@ -157,7 +157,7 @@ describe('orchestration new-worktree workers', () => {
expect.objectContaining({
kind: 'terminal',
role: 'agent',
action: 'reused_agent_terminal',
action: 'created',
id: 'term_worker'
})
])