mirror of
https://github.com/stablyai/orca.git
synced 2026-10-05 08:02:33 +00:00
fix(orchestration): tell a superseded worker it lost the Dispatch
After worker-abandon plus worker-start --retry-of onto another terminal, the old
worker's check found no active Dispatch and no bound Run, had a live pane so
stable_pane_required never fired, and fell through to its direct mailbox with
{count: 0}. That is byte-identical to "no mail", which the worker contract calls
a checkpoint rather than a failure, so the loser kept editing files the new
owner now owned until its worker_done was rejected.
A consuming check whose handle's latest Dispatch settled as failed or
circuit_broken now raises consumer_fenced. A completed Attempt is not fenced:
that terminal is free again and may still receive direct mail. Retries need no
separate successor lookup — every settle that makes an Attempt retry-eligible
(abandon, stop, fail) also drives its Dispatch row to failed or circuit_broken.
--peek and --all stay open so the terminal can still inspect its own inbox.
This commit is contained in:
@@ -6,7 +6,11 @@ import { checkRunMailbox } from './check-run'
|
||||
import { checkWorkerMailbox } from './check-worker'
|
||||
import { checkDirectMailbox } from './check-direct'
|
||||
import { orchestrationSkillRecoveryData } from '../../../../../../shared/orchestration-rpc-contract'
|
||||
import { callerHoldsDispatchPane, dispatchFenced } from './dispatch-mailbox-fence'
|
||||
import {
|
||||
callerHoldsDispatchPane,
|
||||
dispatchFenced,
|
||||
isSupersededDispatch
|
||||
} from './dispatch-mailbox-fence'
|
||||
|
||||
export const ORCHESTRATION_CHECK_METHODS: RpcMethod[] = [
|
||||
defineMethod({
|
||||
@@ -79,9 +83,16 @@ export const ORCHESTRATION_CHECK_METHODS: RpcMethod[] = [
|
||||
remoteAttachment
|
||||
})
|
||||
}
|
||||
// Why: an empty consuming check is the worker contract's "checkpoint, not a failure", so a
|
||||
// caller whose Attempt moved on has to be told rather than handed an empty direct mailbox.
|
||||
const consumingCheck = params.peek !== true && params.all !== true && params.unread !== false
|
||||
const settledDispatch = consumingCheck ? db.getLatestDispatchForTerminal(handle) : undefined
|
||||
if (settledDispatch && isSupersededDispatch(settledDispatch)) {
|
||||
throw dispatchFenced()
|
||||
}
|
||||
// Why: a consuming check on a handle with no live pane and no Dispatch can never see
|
||||
// Run mail, so an empty inbox would read as "nothing yet" instead of a stale caller.
|
||||
if (!paneKey && params.peek !== true && params.all !== true && params.unread !== false) {
|
||||
if (!paneKey && consumingCheck) {
|
||||
throw new OrchestrationError(
|
||||
'stable_pane_required',
|
||||
`Terminal ${handle} has no live pane bound to a Run, so this inbox can never receive Run mail. Rebind this terminal with orchestration run-use, or read the Run mailbox with --run <run_id>.`,
|
||||
|
||||
+122
@@ -0,0 +1,122 @@
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
import type { RpcContext } from '../../../core'
|
||||
import type { OrchestrationDb } from '../../../../orchestration/db'
|
||||
import { createRootDispatch } from '../../../../orchestration/db/root-dispatch-test-fixture'
|
||||
import { createOrchestrationRpcHarness } from '../rpc-test-harness'
|
||||
|
||||
const PANE_OLD = 'tab_old:cccccccc-cccc-4ccc-8ccc-cccccccccccc'
|
||||
const PANE_NEW = 'tab_new:dddddddd-dddd-4ddd-8ddd-dddddddddddd'
|
||||
|
||||
type CheckResult = { messages: { subject: string }[]; count: number }
|
||||
|
||||
/**
|
||||
* worker-abandon + worker-start --retry-of moves the Task to another terminal, but the old worker
|
||||
* keeps polling. Its check used to fall through to the direct mailbox and answer `count: 0`, which
|
||||
* the worker contract reads as "checkpoint, not a failure" — so it kept editing the new owner's files.
|
||||
*/
|
||||
describe('orchestration.check from a terminal whose Attempt was superseded', () => {
|
||||
const h = createOrchestrationRpcHarness()
|
||||
let db: OrchestrationDb
|
||||
let ctx: RpcContext
|
||||
|
||||
afterEach(() => {
|
||||
h.cleanup()
|
||||
})
|
||||
|
||||
function check(handle: string, paneKey: string, params: Record<string, unknown> = {}) {
|
||||
return h.call(
|
||||
'orchestration.check',
|
||||
{ terminal: handle, terminalPaneKey: paneKey, ...params },
|
||||
ctx
|
||||
) as Promise<CheckResult>
|
||||
}
|
||||
|
||||
function startWorker(taskId: string, handle: string, paneKey: string, retryOf?: string): string {
|
||||
const started = db.createStartingWorkerDispatch({
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: Number.MAX_SAFE_INTEGER,
|
||||
taskId,
|
||||
retryOf,
|
||||
startOptions: {}
|
||||
})
|
||||
db.prepareStartingWorkerAuthority({
|
||||
dispatchId: started.dispatch.id,
|
||||
handle,
|
||||
paneKey,
|
||||
processIncarnation: `runtime:${handle}:1`,
|
||||
worktreeId: 'repo::local',
|
||||
setupState: 'not_applicable',
|
||||
effects: []
|
||||
})
|
||||
return started.dispatch.id
|
||||
}
|
||||
|
||||
function retriedOntoAnotherTerminal(): string {
|
||||
;({ db, ctx } = h.setup())
|
||||
const task = db.createTask({ spec: 'work that moves terminals' })
|
||||
const abandoned = startWorker(task.id, 'term_old', PANE_OLD)
|
||||
db.abandonWorkerDispatch(abandoned)
|
||||
startWorker(task.id, 'term_new', PANE_NEW, abandoned)
|
||||
return abandoned
|
||||
}
|
||||
|
||||
it('tells the old worker it lost the Dispatch instead of answering "no mail"', async () => {
|
||||
retriedOntoAnotherTerminal()
|
||||
|
||||
await expect(check('term_old', PANE_OLD)).rejects.toMatchObject({
|
||||
code: 'consumer_fenced',
|
||||
message: expect.stringContaining('no longer owns its mailbox')
|
||||
})
|
||||
})
|
||||
|
||||
// The direct mailbox is the old terminal's own, so inspection stays open; only the consuming
|
||||
// read that a worker treats as a checkpoint is refused.
|
||||
it('still lets the old worker inspect its direct mailbox with --peek and --all', async () => {
|
||||
retriedOntoAnotherTerminal()
|
||||
db.insertMessage({ from: 'term_coord', to: 'term_old', subject: 'stand down' })
|
||||
|
||||
const peeked = await check('term_old', PANE_OLD, { peek: true })
|
||||
const history = await check('term_old', PANE_OLD, { all: true })
|
||||
|
||||
expect(peeked.count).toBe(1)
|
||||
expect(history.count).toBe(1)
|
||||
expect(db.getUnreadMessages('term_old')).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('fences a terminal whose Attempt failed with no successor', async () => {
|
||||
;({ db, ctx } = h.setup())
|
||||
const task = db.createTask({ spec: 'work that failed outright' })
|
||||
const dispatch = createRootDispatch(db, task.id, 'term_old', PANE_OLD)
|
||||
db.failDispatch(dispatch.id, 'worker terminal closed')
|
||||
|
||||
await expect(check('term_old', PANE_OLD)).rejects.toMatchObject({ code: 'consumer_fenced' })
|
||||
})
|
||||
|
||||
it('keeps serving direct mail to a terminal whose Attempt completed normally', async () => {
|
||||
;({ db, ctx } = h.setup())
|
||||
const task = db.createTask({ spec: 'work that finished' })
|
||||
const dispatch = createRootDispatch(db, task.id, 'term_old', PANE_OLD)
|
||||
db.completeDispatch(dispatch.id)
|
||||
db.insertMessage({ from: 'term_coord', to: 'term_old', subject: 'one more thing' })
|
||||
|
||||
const result = await check('term_old', PANE_OLD)
|
||||
|
||||
expect(result.messages.map((message) => message.subject)).toEqual(['one more thing'])
|
||||
expect(db.getUnreadMessages('term_old')).toEqual([])
|
||||
})
|
||||
|
||||
it('serves the new owner its Dispatch mailbox as usual', async () => {
|
||||
const abandoned = retriedOntoAnotherTerminal()
|
||||
const current = db.getDispatchContext(db.getDispatchContextById(abandoned)!.task_id)!
|
||||
db.insertMessage({
|
||||
from: 'term_coord',
|
||||
to: `dispatch:${current.id}`,
|
||||
subject: 'carry on',
|
||||
runId: current.run_id
|
||||
})
|
||||
|
||||
const result = await check('term_new', PANE_NEW)
|
||||
|
||||
expect(result.messages.map((message) => message.subject)).toEqual(['carry on'])
|
||||
})
|
||||
})
|
||||
@@ -1,3 +1,4 @@
|
||||
import type { DispatchContextRow } from '../../../../orchestration/types'
|
||||
import { OrchestrationError } from '../../../../orchestration/orchestration-error'
|
||||
import { isEquivalentPaneKey } from '../../../../orchestration/db/pane-key-match'
|
||||
|
||||
@@ -27,3 +28,14 @@ export function callerHoldsDispatchPane(
|
||||
isEquivalentPaneKey(dispatch.assignee_pane_key, paneKey)
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* A terminal whose last Attempt was abandoned, stopped or failed must not read its direct mailbox:
|
||||
* an empty result is the worker contract's "checkpoint, not a failure", so the loser would keep
|
||||
* working on a Task another terminal now owns. A `completed` Attempt is not fenced — that terminal
|
||||
* is free again and may legitimately receive direct mail. Retries need no separate test: every
|
||||
* settle that makes an Attempt retry-eligible also drives its Dispatch to failed/circuit_broken.
|
||||
*/
|
||||
export function isSupersededDispatch(dispatch: DispatchContextRow): boolean {
|
||||
return dispatch.status === 'failed' || dispatch.status === 'circuit_broken'
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user