diff --git a/config/reliability-gates.jsonc b/config/reliability-gates.jsonc index 899914339c7..0c08481c7e1 100644 --- a/config/reliability-gates.jsonc +++ b/config/reliability-gates.jsonc @@ -13025,7 +13025,7 @@ "https://github.com/stablyai/orca/issues/13821", "https://github.com/stablyai/orca/issues/14347" ], - "invariant": "Injected orchestration task prompts for recognized agent CLIs must send the prompt body inside one bracketed-paste frame, sanitize embedded ESC bytes, preserve chunk boundaries without losing the frame, and submit exactly once only after the agent can accept Enter. A successful orchestration.workerStart must durably record exactly one accepted and started turn; a swallowed Enter must fail with agent_prompt_stalled and never trigger a blind rescue Enter. Claude and Codex must emit a post-paste composer marker and then settle, or reach the bounded fallback first; every other agent retains the platform delay.", + "invariant": "Injected orchestration task prompts for recognized agent CLIs must send the prompt body inside one bracketed-paste frame, sanitize embedded ESC bytes, preserve chunk boundaries without losing the frame, and submit exactly once only after the agent can accept Enter. Local worker-start with supported observation must preserve an unobserved turn as start_unknown without revoking authority, closing questions, or triggering a rescue Enter; a worker report during observation must settle normally. Claude and Codex must emit a post-paste composer marker and then settle, or reach the bounded fallback first; every other agent retains the platform delay.", "oracle": "Runtime tests assert the exact PTY write sequence, failure cleanup, Claude/Codex marker-gated multi-frame renders, and the legacy platform delay for every other configured agent. The candidate resets settlement on later frames, gives a late marker a fresh bounded window, and still submits once at the hard deadline if output never settles. The worker-start contract drives the production RPC through a delayed fake Codex composer and independently checks exact turn/Enter counts plus reopened SQLite Task, Dispatch, worker receipt, and mutation receipt state for accepted and swallowed outcomes. Other orchestration tests assert dispatch/coordinator use the agent prompt path; the live CLI harness covers long Codex-like framing.", "commands": [ "pnpm exec vitest run --config config/vitest.config.ts src/shared/agent-prompt-injection.test.ts src/main/runtime/orca-runtime.test.ts src/main/runtime/rpc/methods/orchestration/runs/tasks-dispatch.test.ts src/main/runtime/orchestration/coordinator.test.ts", @@ -13075,7 +13075,8 @@ "file": "src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts", "assertions": [ "delayed composer readiness produces exactly one submitted and started turn with no premature Enter and durable ready receipts", - "a swallowed Enter records agent_prompt_stalled across Task, Dispatch, worker, and mutation receipts without a rescue Enter" + "a swallowed Enter durably records start_unknown without a rescue Enter or capability revocation", + "early worker reports settle during observation, and outstanding questions survive observation uncertainty" ] }, { diff --git a/src/cli/handlers/orchestration-worker-cli.test.ts b/src/cli/handlers/orchestration-worker-cli.test.ts index cddd36a4cf7..931c4477855 100644 --- a/src/cli/handlers/orchestration-worker-cli.test.ts +++ b/src/cli/handlers/orchestration-worker-cli.test.ts @@ -104,6 +104,34 @@ describe('orchestration worker-start CLI contract', () => { expect(process.exitCode).toBeUndefined() }) + it.each(['succeeded', 'failed'])( + 'accepts a successful start whose task already %s', + async (workerOutcome) => { + const receipt = { + taskId: 'task_1', + dispatchId: 'ctx_1', + state: 'ready', + stage: 'settled', + workerOutcome, + effects: [], + residualResources: [] + } + callMock.mockResolvedValue({ result: receipt }) + await invokeWorkerStart( + new Map([ + ['task', 'task_1'], + ['from', 'term_coord'] + ]) + ) + expect(process.exitCode).toBeUndefined() + expect(printResult).toHaveBeenCalledWith( + expect.objectContaining({ result: receipt }), + true, + expect.any(Function) + ) + } + ) + it('capability-gates and forwards per-invocation launch preferences', async () => { callMock .mockResolvedValueOnce({ diff --git a/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts b/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts index 3584a2e59a6..ddb0e383143 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts @@ -109,9 +109,13 @@ export function settleWorkerReportInTransaction( (dispatch.status === 'pending' || dispatch.status === 'dispatched') && task.status === 'blocked' && reportingWorker?.state === 'start_unknown' + const reportingStart = + dispatch.status === 'pending' && + task.status === 'dispatched' && + reportingWorker?.state === 'starting' const previousDispatchStatus = settledByUnobservedPrompt ? 'failed' - : reconnectingStart + : reconnectingStart || reportingStart ? dispatch.status : 'dispatched' const previousTaskStatus = settledByUnobservedPrompt @@ -198,7 +202,7 @@ export function settleWorkerReportInTransaction( const dispatchTransition = transitionLifecycleWithDb(this.db, { entity: 'dispatch', id: params.dispatchId, - from: reconnectingStart ? ['pending', 'dispatched'] : 'dispatched', + from: reconnectingStart || reportingStart ? ['pending', 'dispatched'] : 'dispatched', to: expectedDispatchStatus, projection: { completed_at: new Date().toISOString(), @@ -234,11 +238,11 @@ export function settleWorkerReportInTransaction( projection: { stage: 'settled', updated_at: new Date().toISOString() }, correction: 'unobserved_prompt_report' }) - } else if (reconnectingStart && params.outcome === 'succeeded') { + } else if ((reconnectingStart || reportingStart) && params.outcome === 'succeeded') { transitionLifecycleWithDb(this.db, { entity: 'worker', id: params.dispatchId, - from: 'start_unknown', + from: reportingStart ? 'starting' : 'start_unknown', to: 'ready' }) transitionLifecycleWithDb(this.db, { @@ -253,7 +257,7 @@ export function settleWorkerReportInTransaction( entity: 'worker', id: params.dispatchId, // A start_unknown success report reconnects through 'ready' above; only failure settles here. - from: params.outcome === 'succeeded' ? 'ready' : ['ready', 'start_unknown'], + from: params.outcome === 'succeeded' ? 'ready' : ['ready', 'start_unknown', 'starting'], to: params.outcome === 'succeeded' ? 'succeeded' : 'failed', projection: { stage: 'settled', updated_at: new Date().toISOString() } }) diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts index f16e207c9e8..c2806aef647 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts @@ -155,7 +155,7 @@ export function markWorkerStartUnknown( from: 'dispatched', to: 'blocked' }) - this.closeQuestionsForDispatch(dispatchId) + // Authority survives uncertainty, so its outstanding questions must remain answerable. this.db.exec('COMMIT') return this.getWorkerDispatch(dispatchId) as WorkerDispatchRow } catch (error) { diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts index b860223f818..593d097580f 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts @@ -41,6 +41,8 @@ const openDatabases: OrchestrationDb[] = [] const temporaryRoots: string[] = [] type PromptContractHarness = { + runtime: Awaited>['runtime'] + handle: string db: OrchestrationDb dbPath: string dispatcher: RpcDispatcher @@ -131,6 +133,8 @@ async function createPromptContractHarness( vi.spyOn(runtime, 'getTerminalOrchestrationCliCommand').mockReturnValue('orca') return { + runtime, + handle, db, dbPath, dispatcher: new RpcDispatcher({ runtime, methods: ORCHESTRATION_METHODS }), @@ -227,6 +231,80 @@ describe('orchestration worker-start prompt contract', () => { }) }) + it.each([ + ['succeeded', false], + ['succeeded', true], + ['failed', false], + ['failed', true] + ] as const)('preserves an early %s report with turn evidence=%s', async (outcome, observed) => { + vi.useFakeTimers() + const harness = await createPromptContractHarness('swallowed') + vi.spyOn(harness.runtime, 'observeTerminalAgentPrompt').mockImplementation( + async (_handle, prompt) => { + const dispatch = harness.db.findActiveDispatchForAssignee(harness.handle) + expect(dispatch).toBeDefined() + expect( + harness.db.settleWorkerReport({ + taskId: harness.taskId, + dispatchId: dispatch!.id, + outcome, + result: 'Finished before the hook arrived' + }) + ).toMatchObject({ action: 'settled', outcome }) + return observed ? { ...prompt, stages: ['input_accepted', 'turn_started'] } : prompt + } + ) + const pending = harness.dispatcher.dispatch(harness.request) + await vi.runAllTimersAsync() + expect(await pending).toMatchObject({ + ok: true, + result: { state: 'ready', stage: 'settled', workerOutcome: outcome } + }) + expect(harness.db.getTask(harness.taskId)?.status).toBe( + outcome === 'succeeded' ? 'completed' : 'failed' + ) + }) + + it('retains accepted authority when the observation binding becomes stale', async () => { + vi.useFakeTimers() + const harness = await createPromptContractHarness('swallowed') + vi.spyOn(harness.runtime, 'observeTerminalAgentPrompt').mockRejectedValue( + new Error('terminal_handle_stale') + ) + const pending = harness.dispatcher.dispatch(harness.request) + await vi.runAllTimersAsync() + expect(await pending).toMatchObject({ ok: true, result: { state: 'outcome_unknown' } }) + expect(harness.db.findActiveDispatchForAssignee(harness.handle)).toMatchObject({ + status: 'pending', + capability_hash: expect.any(String), + capability_revoked_at: null + }) + expect(harness.submittedTurns()).toBe(1) + expect(vi.getTimerCount()).toBe(0) + }) + + it('keeps a worker question answerable after turn observation times out', async () => { + vi.useFakeTimers() + const harness = await createPromptContractHarness('swallowed') + let questionId = '' + vi.spyOn(harness.runtime, 'observeTerminalAgentPrompt').mockImplementation( + async (_handle, prompt) => { + const dispatch = harness.db.findActiveDispatchForAssignee(harness.handle)! + questionId = harness.db.createQuestion({ + runId: dispatch.run_id!, + dispatchId: dispatch.id, + askerHandle: harness.handle, + question: 'Which target should I use?' + }).question.message_id + return prompt + } + ) + const pending = harness.dispatcher.dispatch(harness.request) + await vi.runAllTimersAsync() + expect(await pending).toMatchObject({ ok: true, result: { state: 'outcome_unknown' } }) + expect(harness.db.getQuestion(questionId)?.status).toBe('pending') + }) + it('reports a swallowed Enter as start_unknown while keeping the worker and its capability', async () => { vi.useFakeTimers() const harness = await createPromptContractHarness('swallowed') @@ -273,7 +351,7 @@ describe('orchestration worker-start prompt contract', () => { expect(persisted.getWorkerDispatch(dispatchId)).toMatchObject({ state: 'start_unknown', stage: 'turn_start_unobserved', - last_error: expect.stringContaining('never started a turn') + last_error: expect.stringContaining('turn start could not be verified') }) const persistedEffects = JSON.parse( persisted.getWorkerDispatch(dispatchId)?.effects ?? '[]' diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-readiness-settlement.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-readiness-settlement.ts index 904b9263ea3..a3c57a357d2 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-readiness-settlement.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-readiness-settlement.ts @@ -85,7 +85,10 @@ export async function deliverAndSettleWorkerStartReadiness(args: { setupReceipt: args.setupReceipt, effects }) - if (turnStart.verdict === 'unobserved') { + // A worker report can settle the dispatch while turn observation is outstanding. + const currentWorker = db.getWorkerDispatch(args.dispatchId) + const alreadySettled = currentWorker && currentWorker.state !== 'starting' + if (turnStart.verdict === 'unobserved' && !alreadySettled) { // 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. @@ -125,13 +128,21 @@ export async function deliverAndSettleWorkerStartReadiness(args: { ...(args.terminalRevealWarning ? { warning: args.terminalRevealWarning } : {}) } } - const worker = db.markWorkerDispatchReady(args.dispatchId, effects) + const worker = alreadySettled + ? currentWorker + : db.markWorkerDispatchReady(args.dispatchId, effects) + // A completed task proves start succeeded; older callers use only 'ready' as start success. + const reportedOutcome = + worker.stage === 'settled' && (worker.state === 'succeeded' || worker.state === 'failed') + ? worker.state + : undefined return { runId: run.id, taskId: task.id, dispatchId: args.dispatchId, - state: worker.state, + state: reportedOutcome ? 'ready' : worker.state, stage: worker.stage, + ...(reportedOutcome ? { workerOutcome: reportedOutcome } : {}), turnStart: turnStart.verdict, setup: args.setupReceipt, launch: args.launchReceipt, diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.test.ts index 9e290f7ecf4..f1992eb9b9c 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.test.ts @@ -107,4 +107,14 @@ describe('observeWorkerTurnStart', () => { ).resolves.toEqual({ verdict: 'unsupported', prompt }) expect(observe).not.toHaveBeenCalled() }) + + it('preserves uncertainty when observation loses the terminal binding', async () => { + const prompt = delivery() + const { runtime, observe } = runtimeObserving(prompt) + observe.mockRejectedValue(new Error('terminal_handle_stale')) + await expect( + observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt }) + ).resolves.toEqual({ verdict: 'unobserved', prompt }) + expect(observe).toHaveBeenCalledTimes(1) + }) }) diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.ts index 7c812ffabf3..d1c9e59adfa 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.ts @@ -56,11 +56,17 @@ export async function observeWorkerTurnStart(args: { 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 - ) + let observed: RuntimeTerminalPromptDelivery + try { + observed = await args.runtime.observeTerminalAgentPrompt( + args.terminalHandle, + args.prompt, + args.timeoutMs ?? AGENT_PROMPT_EFFECT_TIMEOUT_MS + ) + } catch { + // Observation failure cannot revoke authority for input that was already accepted. + return { verdict: 'unobserved', prompt: args.prompt } + } if (observed.observation === 'incarnation_replaced') { // The PTY under this handle changed mid-observation; the accepted write is unproven. return { verdict: 'unobserved', prompt: observed } @@ -71,8 +77,8 @@ export async function observeWorkerTurnStart(args: { 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 ` + + `Dispatch input was written and submitted, but ${name}'s turn start could not be verified ` + + `during observation (up to ${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.'