mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(orchestration): stop a structured worker-start reporting a preamble it never delivered
Two ways a structured `worker-start` handed the coordinator a receipt that did not describe the worker it got. `sendStructuredWorkerPreamble` threw only on a refusal and on `rejected`, so a submission that settled `unknown` fell through as success: the start pushed `dispatch_input: accepted` and marked the dispatch ready. `unknown` is not rare — `dispatchSafely` converts ANY thrown adapter call (provider child gone, transport dropped, ack window missed) into it, and `performSend` still returns ok. The worker then has no task spec while its coordinator blocks in `check --wait --types worker_done` until timeout. This PR's own mail lane already states the rule — "`pending` is not yet an acknowledgement; only `accepted` may consume mail" — so the preamble now applies it too, and raises `operation_unknown` for the states that prove neither delivery nor failure, which is the code `failWorkerStartWithReceipt` turns into the `outcome_unknown` receipt whose nextCommands send the coordinator to look. `rejected` stays a proven failure. `--structured` also accepted `--model` / `--effort` and dropped them: structured session creation takes no launch preferences, while `launch.receipt.effective` echoes whatever was requested either way, so `--model opus` ran on the workspace default and the receipt still said `opus`. Refused now, for the same reason `--terminal` refuses them, and the spec note records that refusal along with the new-child/new-top-level one it never mentioned. Tests: the refusal guard had no coverage at all, and `structured-mailbox-pointer-host` — where the full-timeline gate read lives — had none either; reinstating the bounded tail there left the whole repo green. Both are covered now, and the vacuous "never selects an exact provider session" case is re-pointed at the absent `ORCA_PANE_KEY` that actually keeps that selector shut.
This commit is contained in:
@@ -34,7 +34,7 @@ export const ORCHESTRATION_WORKER_COMMAND_SPECS: CommandSpec[] = [
|
||||
'--model supports Claude, Codex, and Cursor opaque provider model ids; --effort requires --model. Neither can combine with --terminal.',
|
||||
'New worktrees use agent-first creation and default --setup to run. Repository start-immediately runs setup beside the agent; wait-for-setup gates agent readiness and task input.',
|
||||
'Creation flags (--name, --repo, --base-branch, --display-name, --comment, --setup) are rejected for current/existing worktrees. Use exact --repo on the selected server; project/host convenience routing remains on worktree create.',
|
||||
'--structured starts the worker as a native structured chat session instead of a terminal agent. Local claude/codex only; it cannot combine with --terminal or a remote --on.',
|
||||
'--structured starts the worker as a native structured chat session instead of a terminal agent. Local claude/codex only, outside WSL; it cannot combine with --terminal, --model, --effort, a remote --on, or a new-child/new-top-level worktree.',
|
||||
'--on selects only the worker server; the Run and this command remain on the current Orca server.',
|
||||
'Remote current and new-child are invalid; discover an exact remote selector or use new-top-level.',
|
||||
'--retry-of links the replacement attempt but does not inherit placement; repeat the intended --on/worktree and --agent/terminal choices.',
|
||||
|
||||
@@ -0,0 +1,132 @@
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentJournalRenderItem } from '../../../shared/agent-session-journal-types'
|
||||
|
||||
const hostRef: { current: unknown } = { current: null }
|
||||
|
||||
vi.mock('../../native-chat/agent-session-wire/structured-agent-session-registry', () => ({
|
||||
getStructuredAgentSessionHost: () => hostRef.current
|
||||
}))
|
||||
|
||||
const { createStructuredMailboxPointerHost, structuredPointerCallerKey } =
|
||||
await import('./structured-mailbox-pointer-host')
|
||||
|
||||
function runningTurn(): AgentJournalRenderItem {
|
||||
return {
|
||||
itemId: 'lifecycle-1',
|
||||
revision: 1,
|
||||
body: { kind: 'status', text: 'working', turnLifecycle: { turnId: 'turn-1', state: 'running' } }
|
||||
} as unknown as AgentJournalRenderItem
|
||||
}
|
||||
|
||||
function transcript(count: number): AgentJournalRenderItem[] {
|
||||
return Array.from(
|
||||
{ length: count },
|
||||
(_unused, index) =>
|
||||
({
|
||||
itemId: `tool-${index}`,
|
||||
revision: 1,
|
||||
body: { kind: 'tool-call', name: 'Bash', input: {}, state: 'completed' }
|
||||
}) as unknown as AgentJournalRenderItem
|
||||
)
|
||||
}
|
||||
|
||||
describe('structured mailbox pointer host', () => {
|
||||
beforeEach(() => {
|
||||
hostRef.current = null
|
||||
})
|
||||
|
||||
it('reads the gate facts from the FULL timeline, never a bounded tail', () => {
|
||||
// The defect this pins: a running turn is announced by ONE lifecycle item, and settlement
|
||||
// tombstones it rather than rewriting it. A long tool-calling turn pushes that item arbitrarily
|
||||
// far from the tail, so any page-sized read reports a busy worker as idle — and the pointer is
|
||||
// then delivered mid-turn, which Codex answers with `turn already running` and Claude settles
|
||||
// `unknown` while the message is really queued.
|
||||
const items = [runningTurn(), ...transcript(500)]
|
||||
hostRef.current = { journalSnapshot: () => ({ items }) }
|
||||
expect(createStructuredMailboxPointerHost().readGateFacts('s1')).toEqual({
|
||||
turnRunning: true,
|
||||
awaitingHuman: false
|
||||
})
|
||||
})
|
||||
|
||||
it('answers null rather than idle when the session cannot be read', () => {
|
||||
// Null retains the pointer; `{turnRunning:false}` would deliver a nudge into a session this
|
||||
// runtime cannot see at all.
|
||||
expect(createStructuredMailboxPointerHost().readGateFacts('s1')).toBeNull()
|
||||
hostRef.current = {
|
||||
journalSnapshot: () => {
|
||||
throw new Error('agent_session_ownership_unknown')
|
||||
}
|
||||
}
|
||||
expect(createStructuredMailboxPointerHost().readGateFacts('s1')).toBeNull()
|
||||
})
|
||||
|
||||
it('reports an unattached host rather than a rejection when nothing can be sent', async () => {
|
||||
await expect(
|
||||
createStructuredMailboxPointerHost().send({
|
||||
sessionId: 's1',
|
||||
dispatchId: 'd1',
|
||||
operationId: 'op1',
|
||||
expectedRuntimeFence: 1,
|
||||
payloadFingerprint: 'fp',
|
||||
body: { kind: 'message', role: 'user', blocks: [] }
|
||||
} as never)
|
||||
).resolves.toEqual({ kind: 'unattached' })
|
||||
})
|
||||
|
||||
it.each([
|
||||
['accepted', 'accepted'],
|
||||
['rejected', 'rejected'],
|
||||
// Neither is an acknowledgement, and only `accepted` may consume mail: both have to reach the
|
||||
// caller as `unknown` so the pointer is retained for the next journal edge.
|
||||
['pending', 'unknown'],
|
||||
['unknown', 'unknown']
|
||||
])('maps a %s submission to %s', async (dispatchState, expected) => {
|
||||
const send = vi.fn(
|
||||
async (_caller: { callerKey: string }, _payload: { retryUnknown?: boolean }) => ({
|
||||
ok: true,
|
||||
value: { submission: { dispatchState } }
|
||||
})
|
||||
)
|
||||
hostRef.current = { send }
|
||||
await expect(
|
||||
createStructuredMailboxPointerHost().send({
|
||||
sessionId: 's1',
|
||||
dispatchId: 'd1',
|
||||
operationId: 'op1',
|
||||
expectedRuntimeFence: 1,
|
||||
payloadFingerprint: 'fp',
|
||||
body: { kind: 'message', role: 'user', blocks: [] }
|
||||
} as never)
|
||||
).resolves.toEqual({ kind: 'sent', state: expected })
|
||||
// Per-dispatch, so one worker's nudges cannot exhaust the shared operation-ledger budget.
|
||||
expect(send.mock.calls[0]![0]).toEqual({ callerKey: structuredPointerCallerKey('d1') })
|
||||
expect(send.mock.calls[0]![1]!.retryUnknown).toBe(true)
|
||||
})
|
||||
|
||||
it('separates a not-attached refusal from a real one', async () => {
|
||||
for (const [code, expected] of [
|
||||
['agent_session_ownership_unknown', { kind: 'unattached' }],
|
||||
['agent_session_conflict', { kind: 'sent', state: 'rejected' }]
|
||||
] as const) {
|
||||
hostRef.current = { send: async () => ({ ok: false, refusal: { code, message: 'no' } }) }
|
||||
await expect(
|
||||
createStructuredMailboxPointerHost().send({
|
||||
sessionId: 's1',
|
||||
dispatchId: 'd1',
|
||||
operationId: 'op1',
|
||||
expectedRuntimeFence: 1,
|
||||
payloadFingerprint: 'fp',
|
||||
body: { kind: 'message', role: 'user', blocks: [] }
|
||||
} as never)
|
||||
).resolves.toEqual(expected)
|
||||
}
|
||||
})
|
||||
|
||||
it('reads the runtime fence off the durable record', () => {
|
||||
hostRef.current = { deps: { store: { getRecord: () => ({ lease: { runtimeFence: 9 } }) } } }
|
||||
expect(createStructuredMailboxPointerHost().currentFence('s1')).toBe(9)
|
||||
hostRef.current = { deps: { store: { getRecord: () => null } } }
|
||||
expect(createStructuredMailboxPointerHost().currentFence('s1')).toBeNull()
|
||||
})
|
||||
})
|
||||
@@ -10,8 +10,13 @@ vi.mock('./structured-agent-session-create', () => ({
|
||||
createStructuredAgentSessionForWorktree: (...args: unknown[]) => createSpy(...args)
|
||||
}))
|
||||
|
||||
const { createStructuredWorkerSession, releaseStructuredWorkerSession, structuredWorkerHoldId } =
|
||||
await import('./orchestration-structured-worker-session')
|
||||
const {
|
||||
createStructuredWorkerSession,
|
||||
releaseStructuredWorkerSession,
|
||||
sendStructuredWorkerPreamble,
|
||||
structuredWorkerHoldId
|
||||
} = await import('./orchestration-structured-worker-session')
|
||||
const { isUnknownWorkerStartOutcome } = await import('./orchestration-worker-topology')
|
||||
const { structuredWorkerIdentities } = await import('../../structured-worker-identity')
|
||||
const { structuredWorkerChildIdentityEnv } =
|
||||
await import('../../structured-worker-child-identity-env')
|
||||
@@ -216,3 +221,45 @@ describe('structured worker session hold', () => {
|
||||
).rejects.toThrow(/local execution host/)
|
||||
})
|
||||
})
|
||||
|
||||
describe('structured worker dispatch preamble', () => {
|
||||
function hostWithSubmission(submission: Record<string, unknown>) {
|
||||
return {
|
||||
deps: { store: { getRecord: () => ({ lease: { runtimeFence: 7 } }) } },
|
||||
send: async () => ({ ok: true, value: { clientMessageId: 'c1', submission } })
|
||||
} as never
|
||||
}
|
||||
|
||||
const send = (host: never) =>
|
||||
sendStructuredWorkerPreamble({ host, sessionId: 's1', dispatchId: 'd1', preamble: 'spec' })
|
||||
|
||||
it('reports the preamble delivered only on an accepted submission', async () => {
|
||||
await expect(
|
||||
send(hostWithSubmission({ dispatchState: 'accepted', reason: null }))
|
||||
).resolves.toBeUndefined()
|
||||
})
|
||||
|
||||
it('never claims delivery for a submission the provider never acknowledged', async () => {
|
||||
// `dispatchSafely` turns ANY thrown adapter call — provider child dead, transport dropped —
|
||||
// into `unknown`, and `performSend` still returns ok. Reporting that as `dispatch_input:
|
||||
// accepted` marks the worker ready with no task, and the coordinator blocks in
|
||||
// `check --wait --types worker_done` until it times out.
|
||||
for (const dispatchState of ['unknown', 'pending'] as const) {
|
||||
const error = await send(
|
||||
hostWithSubmission({ dispatchState, reason: 'provider child exited' })
|
||||
).catch((thrown: unknown) => thrown)
|
||||
expect((error as { code?: string }).code).toBe('operation_unknown')
|
||||
// The wiring, not just the throw: this is the code that makes the start receipt
|
||||
// `outcome_unknown` with the worker-show / worker-abandon recovery commands.
|
||||
expect(isUnknownWorkerStartOutcome(error, 'dispatch_input')).toBe(true)
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps a rejected preamble a proven failure rather than an unknown one', async () => {
|
||||
const error = await send(
|
||||
hostWithSubmission({ dispatchState: 'rejected', reason: 'fence moved' })
|
||||
).catch((thrown: unknown) => thrown)
|
||||
expect((error as Error).message).toMatch(/rejected: fence moved/)
|
||||
expect(isUnknownWorkerStartOutcome(error, 'dispatch_input')).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -235,11 +235,21 @@ export async function sendStructuredWorkerPreamble(args: {
|
||||
if (!result.ok) {
|
||||
throw new Error(`The dispatch preamble was refused: ${result.refusal.message}`)
|
||||
}
|
||||
if (result.value.submission.dispatchState === 'rejected') {
|
||||
throw new Error(
|
||||
`The dispatch preamble was rejected: ${result.value.submission.reason ?? 'no reason given'}`
|
||||
)
|
||||
const submission = result.value.submission
|
||||
if (submission.dispatchState === 'accepted') {
|
||||
return
|
||||
}
|
||||
if (submission.dispatchState === 'rejected') {
|
||||
throw new Error(`The dispatch preamble was rejected: ${submission.reason ?? 'no reason given'}`)
|
||||
}
|
||||
// Only `accepted` is an acknowledgement — the same rule the mail lane already applies. A thrown
|
||||
// adapter call settles as `unknown`, which is indistinguishable from a lost reply, so the start
|
||||
// may claim neither delivery nor failure: `operation_unknown` is what turns this into the
|
||||
// `outcome_unknown` receipt whose nextCommands send the coordinator to look.
|
||||
throw new OrchestrationError(
|
||||
'operation_unknown',
|
||||
`The dispatch preamble was submitted but not acknowledged (${submission.dispatchState}): ${submission.reason ?? 'no reason given'}.`
|
||||
)
|
||||
}
|
||||
|
||||
function requireInstalledHost(): StructuredAgentSessionHost {
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { assertStructuredWorkerStartSupported } from './orchestration-workers'
|
||||
|
||||
// Refusals rather than silent ignores. Each one names a start option `createStructuredWorkerSession`
|
||||
// cannot honour; accepting any of them hands the coordinator a worker that differs from the receipt
|
||||
// it is given back.
|
||||
describe('structured worker-start refusals', () => {
|
||||
it('leaves every option alone when --structured was not asked for', () => {
|
||||
expect(() =>
|
||||
assertStructuredWorkerStartSupported({
|
||||
on: 'server-1',
|
||||
terminal: 'term_1',
|
||||
worktree: 'new-child',
|
||||
model: 'opus',
|
||||
effort: 'high'
|
||||
})
|
||||
).not.toThrow()
|
||||
})
|
||||
|
||||
it('accepts the combinations a structured worker can actually honour', () => {
|
||||
expect(() =>
|
||||
assertStructuredWorkerStartSupported({ structured: true, worktree: 'current' })
|
||||
).not.toThrow()
|
||||
expect(() => assertStructuredWorkerStartSupported({ structured: true })).not.toThrow()
|
||||
})
|
||||
|
||||
it.each([
|
||||
['on', { on: 'server-1' }, /--on/],
|
||||
['terminal', { terminal: 'term_1' }, /--terminal/],
|
||||
['new-child', { worktree: 'new-child' }, /create the worktree first/],
|
||||
['new-top-level', { worktree: 'new-top-level' }, /create the worktree first/],
|
||||
// Launch preferences reach the terminal path as `launch.preferences` and the structured path
|
||||
// not at all, while `launch.receipt.effective` echoes whatever was requested either way. Left
|
||||
// accepted, `--model opus` runs on the workspace default and the receipt still says `opus`.
|
||||
['model', { model: 'opus' }, /--model and --effort/],
|
||||
['effort', { model: 'opus', effort: 'high' }, /--model and --effort/]
|
||||
])('refuses --structured with %s', (_name, params, message) => {
|
||||
expect(() => assertStructuredWorkerStartSupported({ structured: true, ...params })).toThrow(
|
||||
message
|
||||
)
|
||||
})
|
||||
})
|
||||
@@ -76,15 +76,26 @@ export const ORCHESTRATION_WORKER_START_METHODS: RpcMethod[] = [
|
||||
* Refused rather than ignored: silently starting a terminal worker under a structured request
|
||||
* would hand the coordinator a worker of a different kind than it asked for.
|
||||
*/
|
||||
function assertStructuredWorkerStartSupported(params: {
|
||||
export function assertStructuredWorkerStartSupported(params: {
|
||||
structured?: boolean
|
||||
on?: string
|
||||
terminal?: string
|
||||
worktree?: string
|
||||
model?: string
|
||||
effort?: string
|
||||
}): void {
|
||||
if (!params.structured) {
|
||||
return
|
||||
}
|
||||
if (params.model || params.effort) {
|
||||
// Structured session creation takes no launch preferences, so accepting these would run the
|
||||
// worker on the workspace default while the start receipt's `launch.effective` reported the
|
||||
// model that was asked for. Refused for the same reason `--terminal` refuses them.
|
||||
throw new OrchestrationError(
|
||||
'invalid_argument',
|
||||
'--model and --effort cannot be applied to a structured worker; its session uses the workspace default.'
|
||||
)
|
||||
}
|
||||
if (params.on) {
|
||||
throw new OrchestrationError(
|
||||
'invalid_argument',
|
||||
|
||||
@@ -2,6 +2,7 @@ import { describe, expect, it, beforeEach } from 'vitest'
|
||||
import { isTerminalLeafId, parsePaneKey } from '../../shared/stable-pane-id'
|
||||
import { structuredAgentSessionPaneKey } from '../../shared/structured-agent-session-projection'
|
||||
import { selectExactWorkerProviderSession } from './orchestration/worker-provider-session'
|
||||
import { structuredWorkerChildIdentityEnv } from './structured-worker-child-identity-env'
|
||||
import {
|
||||
StructuredWorkerIdentityRegistry,
|
||||
isStructuredWorkerHandle,
|
||||
@@ -10,6 +11,7 @@ import {
|
||||
sessionIdFromStructuredWorkerIncarnation,
|
||||
structuredWorkerHostScope,
|
||||
structuredWorkerPaneKeyBelongsToSession,
|
||||
structuredWorkerIdentities,
|
||||
structuredWorkerProcessIncarnation,
|
||||
structuredWorkerRecordIsCurrent
|
||||
} from './structured-worker-identity'
|
||||
@@ -185,9 +187,31 @@ describe('structured worker identity registry', () => {
|
||||
})
|
||||
|
||||
describe('structured workers stay outside the PTY-only fail-closed paths', () => {
|
||||
it('never selects an exact provider session', () => {
|
||||
// Fail-closed because a structured session emits no hook agent status. It stays closed only
|
||||
// while ORCA_PANE_KEY is absent from the structured child's environment.
|
||||
it('keeps the selector shut by never letting a structured pane key reach a hook status', () => {
|
||||
// The selector matches on pane key, so it is fail-closed for a structured worker only while
|
||||
// ORCA_PANE_KEY is absent from its child's environment. That absence IS the guard: put the key
|
||||
// back and the first assertion below is what an attacker gets.
|
||||
// The PROCESS registry, because that is the one the spawn path reads.
|
||||
const handle = mintStructuredWorkerHandle()
|
||||
structuredWorkerIdentities.register({
|
||||
handle,
|
||||
sessionId: SESSION_ID,
|
||||
agent: 'claude',
|
||||
paneKey: mintStructuredWorkerPaneKey(SESSION_ID),
|
||||
processIncarnation: structuredWorkerProcessIncarnation(SESSION_ID),
|
||||
worktreeId: 'wt_1',
|
||||
hostScope: { kind: 'local', hostId: 'local' }
|
||||
})
|
||||
try {
|
||||
const env = structuredWorkerChildIdentityEnv(SESSION_ID)
|
||||
// Registered, so this is a populated env — not the empty one an unregistered session gets,
|
||||
// which would satisfy the pane-key assertion for the wrong reason.
|
||||
expect(env.ORCA_TERMINAL_HANDLE).toBe(handle)
|
||||
expect(Object.keys(env)).not.toContain('ORCA_PANE_KEY')
|
||||
} finally {
|
||||
structuredWorkerIdentities.forget(handle)
|
||||
}
|
||||
|
||||
const paneKey = mintStructuredWorkerPaneKey(SESSION_ID)
|
||||
expect(
|
||||
selectExactWorkerProviderSession({
|
||||
|
||||
Reference in New Issue
Block a user