fix: preserve worker authority through start observation

This commit is contained in:
Merge Sim
2026-09-07 18:44:56 -07:00
parent 325fa88b09
commit ba6a5f0318
8 changed files with 157 additions and 19 deletions
+3 -2
View File
@@ -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"
]
},
{
@@ -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({
@@ -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() }
})
@@ -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) {
@@ -41,6 +41,8 @@ const openDatabases: OrchestrationDb[] = []
const temporaryRoots: string[] = []
type PromptContractHarness = {
runtime: Awaited<ReturnType<typeof createAgentPromptSubmissionRuntime>>['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 ?? '[]'
@@ -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,
@@ -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)
})
})
@@ -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.'