mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(orchestration): keep a stopping worker from wedging dispatch and the relay
`updateTaskStatus` ran its live-supervised-worker guard only for terminal statuses, so `task-update --status dispatched` re-opened a Task under a worker that was already stopping. The worker's own report then had nowhere to land: settlement transitioned the worker from a hardcoded `ready`, so it threw `lifecycle_conflict` out of the federated relay import (rolling back the batch and freezing `to_home_imported_sequence`, which the contiguity check then used to refuse every later sequence forever) and out of the coordinator's message sweep (marking the run failed with its batch unread). Three edges, one rule each: - every status the Task guard lets past the active-Dispatch check now clears the same live-worker check, so a genuine re-open cannot happen under a live worker; - settlement reads the worker's actual state and asks the lifecycle graph, then rejects with `worker_not_settleable` instead of throwing, so already-wedged databases drain rather than stall; - `beginWorkerStop` accepts a re-issue from `stopping`, giving the operator a legal exit whose outcome is the honest `stop_unknown` rather than an asserted exit. The state-space test drives the real public operations and proves no reachable (Dispatch, Task, worker) triple can make a worker report throw.
This commit is contained in:
@@ -0,0 +1,237 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { OrchestrationDb } from './db'
|
||||
|
||||
const PANE = 'tab_w:11111111-1111-4111-8111-111111111111'
|
||||
|
||||
/**
|
||||
* A worker report arrives inside a federated relay import and inside the coordinator's message
|
||||
* sweep, and a throw rolls both back — the relay cursor never advances again and the run is
|
||||
* marked failed with its batch unread (#16904). So `settleWorkerReport` must answer every
|
||||
* reachable (Dispatch, Task, worker) triple with a settlement or a structured rejection.
|
||||
*
|
||||
* The search drives the real public operations and memoizes on the resulting triple, so it stays
|
||||
* a state-space proof rather than a list of cases someone remembered to write down.
|
||||
*/
|
||||
const OPERATIONS = [
|
||||
'prepareAuthority',
|
||||
'markReady',
|
||||
'markStartUnknown',
|
||||
'failStart',
|
||||
'beginStop',
|
||||
'settleStop',
|
||||
'markStopUnknown',
|
||||
'resumeFedRelay',
|
||||
'abandon',
|
||||
'reconcileMissing',
|
||||
'fedStart:ready',
|
||||
'fedStart:start_unknown',
|
||||
'fedStart:failed',
|
||||
'fedStart:stopped',
|
||||
'fedStop',
|
||||
'failDispatch',
|
||||
'failDispatch:exited',
|
||||
'completeDispatch',
|
||||
'taskUpdate:ready',
|
||||
'taskUpdate:blocked',
|
||||
'taskUpdate:dispatched',
|
||||
'taskUpdate:completed',
|
||||
'taskUpdate:failed',
|
||||
'mintCap',
|
||||
'revokeCap',
|
||||
'report:succeeded',
|
||||
'report:failed'
|
||||
] as const
|
||||
|
||||
type Probe = {
|
||||
db: OrchestrationDb
|
||||
taskId: string
|
||||
dispatchId: string
|
||||
}
|
||||
|
||||
function replay(sequence: readonly string[]): Probe {
|
||||
const db = new OrchestrationDb(':memory:')
|
||||
const task = db.createTask({ spec: 'state space' })
|
||||
const { dispatch } = db.createStartingWorkerDispatch({
|
||||
taskId: task.id,
|
||||
startOptions: {},
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: 99,
|
||||
federation: {
|
||||
environmentId: 'e',
|
||||
environmentName: 'n',
|
||||
peerFingerprint: 'p',
|
||||
protocolVersion: 12
|
||||
}
|
||||
})
|
||||
const id = dispatch.id
|
||||
const apply = (name: string): void => {
|
||||
switch (name) {
|
||||
case 'prepareAuthority':
|
||||
db.prepareStartingWorkerAuthority({
|
||||
dispatchId: id,
|
||||
handle: 'term_w',
|
||||
paneKey: PANE,
|
||||
processIncarnation: 'inc1',
|
||||
worktreeId: 'wt',
|
||||
effects: [],
|
||||
setupState: 'not_configured'
|
||||
})
|
||||
return
|
||||
case 'markReady':
|
||||
db.markWorkerDispatchReady(id)
|
||||
return
|
||||
case 'markStartUnknown':
|
||||
db.markWorkerStartUnknown(id, 's', 'w')
|
||||
return
|
||||
case 'failStart':
|
||||
db.failWorkerStart(id, 's', 'w')
|
||||
return
|
||||
case 'beginStop':
|
||||
db.beginWorkerStop(id, 'ep')
|
||||
return
|
||||
case 'settleStop':
|
||||
db.settleWorkerStop(id)
|
||||
return
|
||||
case 'markStopUnknown':
|
||||
db.markWorkerStopUnknown(id, 'w')
|
||||
return
|
||||
case 'resumeFedRelay':
|
||||
db.resumeFederatedWorkerForTerminalRelay(id)
|
||||
return
|
||||
case 'abandon':
|
||||
db.abandonWorkerDispatch(id)
|
||||
return
|
||||
case 'reconcileMissing':
|
||||
db.reconcileMissingWorkerTerminal(id, 'gone')
|
||||
return
|
||||
case 'fedStart:ready':
|
||||
db.reconcileFederatedWorkerStart({ dispatchId: id, state: 'ready', stage: 's' })
|
||||
return
|
||||
case 'fedStart:start_unknown':
|
||||
db.reconcileFederatedWorkerStart({ dispatchId: id, state: 'start_unknown', stage: 's' })
|
||||
return
|
||||
case 'fedStart:failed':
|
||||
db.reconcileFederatedWorkerStart({ dispatchId: id, state: 'failed', stage: 's' })
|
||||
return
|
||||
case 'fedStart:stopped':
|
||||
db.reconcileFederatedWorkerStart({ dispatchId: id, state: 'stopped', stage: 's' })
|
||||
return
|
||||
case 'fedStop':
|
||||
db.reconcileFederatedWorkerStop(id)
|
||||
return
|
||||
case 'failDispatch':
|
||||
db.failDispatch(id, 'x')
|
||||
return
|
||||
case 'failDispatch:exited':
|
||||
db.failDispatch(id, 'x', { workerProcessExited: true })
|
||||
return
|
||||
case 'completeDispatch':
|
||||
db.completeDispatch(id)
|
||||
return
|
||||
case 'taskUpdate:ready':
|
||||
db.updateTaskStatus(task.id, 'ready')
|
||||
return
|
||||
case 'taskUpdate:blocked':
|
||||
db.updateTaskStatus(task.id, 'blocked')
|
||||
return
|
||||
case 'taskUpdate:dispatched':
|
||||
db.updateTaskStatus(task.id, 'dispatched')
|
||||
return
|
||||
case 'taskUpdate:completed':
|
||||
db.updateTaskStatus(task.id, 'completed', 'r')
|
||||
return
|
||||
case 'taskUpdate:failed':
|
||||
db.updateTaskStatus(task.id, 'failed', 'r')
|
||||
return
|
||||
case 'mintCap':
|
||||
db.mintDispatchCapability({ dispatchId: id, paneKey: PANE, processIncarnation: 'inc2' })
|
||||
return
|
||||
case 'revokeCap':
|
||||
db.revokeDispatchCapability(id)
|
||||
return
|
||||
case 'report:succeeded':
|
||||
db.settleWorkerReport({
|
||||
taskId: task.id,
|
||||
dispatchId: id,
|
||||
outcome: 'succeeded',
|
||||
result: 'r'
|
||||
})
|
||||
return
|
||||
case 'report:failed':
|
||||
db.settleWorkerReport({ taskId: task.id, dispatchId: id, outcome: 'failed', result: 'r' })
|
||||
}
|
||||
}
|
||||
for (const name of sequence) {
|
||||
try {
|
||||
apply(name)
|
||||
} catch {
|
||||
// A refused operation is a legal outcome; the state it did not reach is simply not explored.
|
||||
}
|
||||
}
|
||||
return { db, taskId: task.id, dispatchId: id }
|
||||
}
|
||||
|
||||
function tripleOf(probe: Probe): string {
|
||||
const dispatch = probe.db.getDispatchContextById(probe.dispatchId)
|
||||
const task = probe.db.getTask(probe.taskId)
|
||||
const worker = probe.db.getWorkerDispatch(probe.dispatchId)
|
||||
return `${dispatch?.status}/${task?.status}/${worker?.state}`
|
||||
}
|
||||
|
||||
function reachableTriples(maxDepth: number): Map<string, string[]> {
|
||||
const seen = new Map<string, string[]>()
|
||||
const initial = replay([])
|
||||
seen.set(tripleOf(initial), [])
|
||||
initial.db.close()
|
||||
let frontier: string[][] = [[]]
|
||||
for (let depth = 0; depth < maxDepth && frontier.length > 0; depth += 1) {
|
||||
const next: string[][] = []
|
||||
for (const sequence of frontier) {
|
||||
for (const operation of OPERATIONS) {
|
||||
const probe = replay([...sequence, operation])
|
||||
const triple = tripleOf(probe)
|
||||
probe.db.close()
|
||||
if (seen.has(triple)) {
|
||||
continue
|
||||
}
|
||||
seen.set(triple, [...sequence, operation])
|
||||
next.push([...sequence, operation])
|
||||
}
|
||||
}
|
||||
frontier = next
|
||||
}
|
||||
return seen
|
||||
}
|
||||
|
||||
describe('worker report settlement over the reachable lifecycle state space', () => {
|
||||
it('answers every reachable state without throwing', () => {
|
||||
// Five is where the triple set saturates: a sixth round of all 27 operations adds none.
|
||||
const reachable = reachableTriples(5)
|
||||
const throwing: string[] = []
|
||||
|
||||
for (const [triple, sequence] of reachable) {
|
||||
for (const outcome of ['succeeded', 'failed'] as const) {
|
||||
const probe = replay(sequence)
|
||||
try {
|
||||
const settlement = probe.db.settleWorkerReport({
|
||||
taskId: probe.taskId,
|
||||
dispatchId: probe.dispatchId,
|
||||
outcome,
|
||||
result: 'r'
|
||||
})
|
||||
expect(settlement.action === 'settled' || settlement.action === 'rejected').toBe(true)
|
||||
} catch (error) {
|
||||
throwing.push(
|
||||
`${triple}/${outcome} after ${sequence.join(' -> ') || '(initial)'}: ${(error as Error).message}`
|
||||
)
|
||||
} finally {
|
||||
probe.db.close()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
expect(throwing).toEqual([])
|
||||
// Guards the search itself: a harness that stopped exploring would also report zero throws.
|
||||
expect(reachable.size).toBeGreaterThan(50)
|
||||
}, 300_000)
|
||||
})
|
||||
@@ -0,0 +1,254 @@
|
||||
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
|
||||
import { OrchestrationDb } from './db'
|
||||
import { Coordinator } from './coordinator'
|
||||
import type { CoordinatorRuntime } from './coordinator-runtime-contract'
|
||||
import { reconcileLifecycleMessage } from './lifecycle-reconciliation'
|
||||
import { createRootDispatch } from './db/root-dispatch-test-fixture'
|
||||
import type { MessagePriority, MessageType } from './types'
|
||||
|
||||
const PANE_W = 'tab_w:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'
|
||||
|
||||
describe('a Task whose supervised worker is stopping', () => {
|
||||
let db: OrchestrationDb
|
||||
beforeEach(() => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
})
|
||||
afterEach(() => db.close())
|
||||
|
||||
function localWorker() {
|
||||
const task = db.createTask({ spec: 'local work' })
|
||||
const { dispatch } = db.createStartingWorkerDispatch({
|
||||
taskId: task.id,
|
||||
startOptions: {},
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: 9
|
||||
})
|
||||
db.prepareStartingWorkerAuthority({
|
||||
dispatchId: dispatch.id,
|
||||
handle: 'term_w',
|
||||
paneKey: PANE_W,
|
||||
processIncarnation: 'inc1',
|
||||
worktreeId: 'wt',
|
||||
effects: [],
|
||||
setupState: 'not_configured'
|
||||
})
|
||||
db.markWorkerDispatchReady(dispatch.id)
|
||||
return { task, dispatch }
|
||||
}
|
||||
|
||||
function federatedWorker() {
|
||||
const task = db.createTask({ spec: 'remote work' })
|
||||
const { dispatch } = db.createStartingWorkerDispatch({
|
||||
taskId: task.id,
|
||||
startOptions: {},
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: 99,
|
||||
federation: {
|
||||
environmentId: 'env_remote',
|
||||
environmentName: 'remote',
|
||||
peerFingerprint: 'peer_fp',
|
||||
protocolVersion: 12
|
||||
}
|
||||
})
|
||||
db.reconcileFederatedWorkerStart({
|
||||
dispatchId: dispatch.id,
|
||||
state: 'ready',
|
||||
stage: 'input_accepted',
|
||||
worktreeId: 'wt_remote',
|
||||
terminalHandle: 'term_remote'
|
||||
})
|
||||
return { task, dispatch }
|
||||
}
|
||||
|
||||
function relayItem(
|
||||
taskId: string,
|
||||
runId: string,
|
||||
dispatchId: string,
|
||||
sequence: number,
|
||||
kind: 'status' | 'done'
|
||||
) {
|
||||
return {
|
||||
dispatchId,
|
||||
sequence,
|
||||
message: {
|
||||
id: `msg_${sequence}`,
|
||||
runId,
|
||||
from: `dispatch:${dispatchId}`,
|
||||
to: `run:${runId}`,
|
||||
subject: kind === 'done' ? 'Worker done' : `Update ${sequence}`,
|
||||
body: 'b',
|
||||
type: 'status' as MessageType,
|
||||
priority: 'normal' as MessagePriority
|
||||
},
|
||||
lifecycle:
|
||||
kind === 'done'
|
||||
? {
|
||||
kind: 'worker_report' as const,
|
||||
taskId,
|
||||
outcome: 'succeeded' as const,
|
||||
result: 'the real answer'
|
||||
}
|
||||
: { kind: 'none' as const }
|
||||
}
|
||||
}
|
||||
|
||||
describe('task-update', () => {
|
||||
it('refuses to re-open the Task while the worker is stopping', () => {
|
||||
const { task, dispatch } = localWorker()
|
||||
db.beginWorkerStop(dispatch.id, 'epoch_home')
|
||||
expect(db.getTask(task.id)?.status).toBe('blocked')
|
||||
|
||||
expect(() => db.updateTaskStatus(task.id, 'dispatched')).toThrowError(
|
||||
expect.objectContaining({
|
||||
code: 'task_not_startable',
|
||||
data: { taskId: task.id, dispatchId: dispatch.id }
|
||||
})
|
||||
)
|
||||
expect(db.getTask(task.id)?.status).toBe('blocked')
|
||||
})
|
||||
|
||||
it('refuses to re-open the Task while the stop outcome is unknown', () => {
|
||||
const { task, dispatch } = localWorker()
|
||||
db.beginWorkerStop(dispatch.id, 'epoch_home')
|
||||
db.markWorkerStopUnknown(dispatch.id, 'the execution host did not answer')
|
||||
|
||||
expect(() => db.updateTaskStatus(task.id, 'dispatched')).toThrowError(
|
||||
expect.objectContaining({ code: 'task_not_startable' })
|
||||
)
|
||||
expect(db.getTask(task.id)?.status).toBe('blocked')
|
||||
})
|
||||
|
||||
it('control: still accepts dispatched for an active Dispatch with no supervised worker', () => {
|
||||
const task = db.createTask({ spec: 'unsupervised work' })
|
||||
createRootDispatch(db, task.id, 'term_worker')
|
||||
|
||||
expect(db.updateTaskStatus(task.id, 'dispatched')?.status).toBe('dispatched')
|
||||
})
|
||||
|
||||
it('control: still accepts dispatched while the supervised worker is ready', () => {
|
||||
const { task } = localWorker()
|
||||
|
||||
expect(db.updateTaskStatus(task.id, 'dispatched')?.status).toBe('dispatched')
|
||||
})
|
||||
})
|
||||
|
||||
describe('a worker report the lifecycle graph cannot settle', () => {
|
||||
// Rows an older binary already wrote: it let task-update re-open the Task under a stopping
|
||||
// worker, so a shipped database can hold this triple even though nothing can reach it now.
|
||||
function wedgeTaskDispatchedUnderStoppingWorker(taskId: string, dispatchId: string): void {
|
||||
db.beginWorkerStop(dispatchId, 'epoch_home')
|
||||
db.db.prepare("UPDATE tasks SET status = 'dispatched' WHERE id = ?").run(taskId)
|
||||
}
|
||||
|
||||
it('rejects the federated report and still advances the relay cursor', () => {
|
||||
const { task, dispatch } = federatedWorker()
|
||||
wedgeTaskDispatchedUnderStoppingWorker(task.id, dispatch.id)
|
||||
|
||||
const imported = db.importFederatedRelayItem(
|
||||
relayItem(task.id, task.run_id, dispatch.id, 1, 'done')
|
||||
)
|
||||
|
||||
expect(imported.lifecycle).toMatchObject({
|
||||
action: 'rejected',
|
||||
code: 'worker_not_settleable'
|
||||
})
|
||||
expect(db.getFederatedDispatch(dispatch.id)?.to_home_imported_sequence).toBe(1)
|
||||
expect(db.getMessageById('msg_1')).toBeDefined()
|
||||
// The stream is not wedged behind the report it could not apply.
|
||||
expect(
|
||||
db.importFederatedRelayItem(relayItem(task.id, task.run_id, dispatch.id, 2, 'status'))
|
||||
.message.id
|
||||
).toBe('msg_2')
|
||||
expect(db.getFederatedDispatch(dispatch.id)?.to_home_imported_sequence).toBe(2)
|
||||
})
|
||||
|
||||
it('rejects the local report instead of throwing out of reconciliation', () => {
|
||||
const { task, dispatch } = localWorker()
|
||||
wedgeTaskDispatchedUnderStoppingWorker(task.id, dispatch.id)
|
||||
const msg = workerDoneMessage(task.id, task.run_id, dispatch.id)
|
||||
|
||||
expect(reconcileLifecycleMessage(db, msg)).toMatchObject({
|
||||
action: 'rejected',
|
||||
code: 'worker_not_settleable'
|
||||
})
|
||||
expect(db.getWorkerDispatch(dispatch.id)?.state).toBe('stopping')
|
||||
})
|
||||
|
||||
it('lets the coordinator loop finish its batch and mark the report read', async () => {
|
||||
const { task, dispatch } = localWorker()
|
||||
wedgeTaskDispatchedUnderStoppingWorker(task.id, dispatch.id)
|
||||
workerDoneMessage(task.id, task.run_id, dispatch.id)
|
||||
const coordinator = new Coordinator(db, stubCoordinatorRuntime(), {
|
||||
spec: 'stop-wedge',
|
||||
coordinatorHandle: `run:${task.run_id}`,
|
||||
pollIntervalMs: 0,
|
||||
onLog: (line) => {
|
||||
if (line.includes('rejected')) {
|
||||
coordinator.stop()
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
const run = await coordinator.runFromExistingRun(task.run_id)
|
||||
|
||||
expect(run.failedTasks).toEqual([])
|
||||
expect(db.getUnreadMessages(`run:${task.run_id}`)).toEqual([])
|
||||
})
|
||||
|
||||
it('control: a ready worker still settles its Task', () => {
|
||||
const { task, dispatch } = localWorker()
|
||||
|
||||
expect(
|
||||
db.settleWorkerReport({
|
||||
taskId: task.id,
|
||||
dispatchId: dispatch.id,
|
||||
outcome: 'succeeded',
|
||||
result: 'done'
|
||||
})
|
||||
).toMatchObject({ action: 'settled', outcome: 'succeeded' })
|
||||
expect(db.getTask(task.id)?.status).toBe('completed')
|
||||
expect(db.getWorkerDispatch(dispatch.id)?.state).toBe('succeeded')
|
||||
})
|
||||
})
|
||||
|
||||
describe('operator escape', () => {
|
||||
it('accepts a re-issued worker-stop and reaches an honest stop_unknown outcome', () => {
|
||||
const { task, dispatch } = localWorker()
|
||||
db.beginWorkerStop(dispatch.id, 'epoch_dead_runtime')
|
||||
|
||||
// The runtime that owned the first stop died mid-flight; the re-issue is the way out.
|
||||
const reissued = db.beginWorkerStop(dispatch.id, 'epoch_new_runtime')
|
||||
expect(reissued).toMatchObject({ disposition: 'stopping' })
|
||||
expect(db.getWorkerDispatch(dispatch.id)?.runtime_epoch).toBe('epoch_new_runtime')
|
||||
|
||||
db.markWorkerStopUnknown(dispatch.id, 'the execution host did not answer')
|
||||
expect(db.abandonWorkerDispatch(dispatch.id)).toMatchObject({ disposition: 'abandoned' })
|
||||
expect(db.getTask(task.id)?.status).toBe('blocked')
|
||||
})
|
||||
})
|
||||
|
||||
function workerDoneMessage(taskId: string, runId: string, dispatchId: string) {
|
||||
return db.insertMessage({
|
||||
runId,
|
||||
from: 'term_w',
|
||||
to: `run:${runId}`,
|
||||
subject: 'Worker done',
|
||||
body: 'finished the work',
|
||||
type: 'worker_done',
|
||||
priority: 'normal',
|
||||
senderPaneKey: PANE_W,
|
||||
payload: JSON.stringify({ taskId, dispatchId, outcome: 'succeeded' })
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
/** The coordinator's message sweep touches none of these; a dispatched Task creates no worker. */
|
||||
function stubCoordinatorRuntime(): CoordinatorRuntime {
|
||||
return {
|
||||
sendTerminalAgentPrompt: async () => ({}),
|
||||
listTerminals: async () => ({ terminals: [] }),
|
||||
createTerminal: async () => ({ handle: 'term_unused', worktreeId: 'wt' }),
|
||||
waitForTerminal: async (handle: string) => ({ handle, condition: 'idle' }),
|
||||
probeWorktreeDrift: async () => null
|
||||
}
|
||||
}
|
||||
@@ -3,7 +3,7 @@ import type { OrchestrationDb } from '../orchestration-db'
|
||||
import { AGENT_PROMPT_STALLED_ERROR } from '../../../agent-prompt-submission-verification'
|
||||
import { settleActiveDispatchesForTask } from './dispatch-completion'
|
||||
import { getActiveDispatchForTask } from './task-dispatch-reconciliation'
|
||||
import { transitionLifecycleWithDb } from '../lifecycle-transition'
|
||||
import { isLegalLifecycleTransition, transitionLifecycleWithDb } from '../lifecycle-transition'
|
||||
import { runLifecycleWriteTransaction } from '../lifecycle-write-transaction-runner'
|
||||
|
||||
type WorkerReportObservation = {
|
||||
@@ -159,6 +159,24 @@ export function settleWorkerReportInTransaction(
|
||||
)
|
||||
.all(params.taskId, params.dispatchId) as { id: string }[]
|
||||
|
||||
const settledWorkerState = params.outcome === 'succeeded' ? 'succeeded' : 'failed'
|
||||
// Why: this settlement rides a federated relay import and the coordinator's message sweep, and
|
||||
// both roll the whole batch back on a throw — a relay that can never advance its cursor and a
|
||||
// run marked failed with its messages left unread (#16904). A worker whose own state the
|
||||
// lifecycle graph will not settle is a rejection, which those callers already record and move on.
|
||||
if (
|
||||
reportingWorker &&
|
||||
!settledByUnobservedPrompt &&
|
||||
!reconnectingStart &&
|
||||
!isLegalLifecycleTransition('worker', reportingWorker.state, settledWorkerState)
|
||||
) {
|
||||
return {
|
||||
action: 'rejected',
|
||||
code: 'worker_not_settleable',
|
||||
reason: `Dispatch ${params.dispatchId} worker is ${reportingWorker.state}; it cannot settle as ${settledWorkerState}.`
|
||||
}
|
||||
}
|
||||
|
||||
this.db.exec('SAVEPOINT settle_worker_report')
|
||||
let dispatchUpdate: { changes: number }
|
||||
let taskUpdate: { changes: number }
|
||||
@@ -250,11 +268,12 @@ export function settleWorkerReportInTransaction(
|
||||
})
|
||||
} else if (reportingWorker) {
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
// The guard above already proved the graph accepts this edge from the state the worker is
|
||||
// actually in; a hardcoded `from` would instead throw for every other live state.
|
||||
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'],
|
||||
to: params.outcome === 'succeeded' ? 'succeeded' : 'failed',
|
||||
from: reportingWorker.state,
|
||||
to: settledWorkerState,
|
||||
projection: { stage: 'settled', updated_at: new Date().toISOString() }
|
||||
})
|
||||
}
|
||||
|
||||
@@ -117,6 +117,15 @@ const PROJECTION_COLUMNS = new Set([
|
||||
'runtime_epoch'
|
||||
])
|
||||
|
||||
/** The lifecycle graph asked without writing, for callers that must reject rather than throw. */
|
||||
export function isLegalLifecycleTransition(
|
||||
entity: LifecycleEntity,
|
||||
from: string,
|
||||
to: string
|
||||
): boolean {
|
||||
return (LEGAL_TRANSITIONS[entity][from] ?? []).includes(to)
|
||||
}
|
||||
|
||||
export function transitionLifecycle(
|
||||
this: OrchestrationDb,
|
||||
params: LifecycleTransitionParams
|
||||
@@ -151,14 +160,16 @@ export function transitionLifecycleWithDb(
|
||||
{ entity: params.entity, id: params.id, state: current.state }
|
||||
)
|
||||
}
|
||||
const legal = LEGAL_TRANSITIONS[params.entity][current.state] ?? []
|
||||
const promptReportCorrection =
|
||||
params.correction === 'unobserved_prompt_report' &&
|
||||
current.state === 'failed' &&
|
||||
((params.entity === 'task' && params.to === 'completed') ||
|
||||
(params.entity === 'dispatch' && params.to === 'completed') ||
|
||||
(params.entity === 'worker' && params.to === 'succeeded'))
|
||||
if (!legal.includes(params.to) && !promptReportCorrection) {
|
||||
if (
|
||||
!isLegalLifecycleTransition(params.entity, current.state, params.to) &&
|
||||
!promptReportCorrection
|
||||
) {
|
||||
throw new OrchestrationError(
|
||||
'lifecycle_conflict',
|
||||
`${params.entity} ${params.id} cannot transition from ${current.state} to ${params.to}.`,
|
||||
|
||||
@@ -35,18 +35,24 @@ export function updateTaskStatus(
|
||||
ORDER BY rowid DESC LIMIT 1`
|
||||
)
|
||||
.get(id) as { id: string } | undefined
|
||||
const activeWorker = terminalStatus
|
||||
? (this.db
|
||||
.prepare(
|
||||
`SELECT active.id
|
||||
// Why: a supervised worker owns its Task for as long as it is alive. Every status this
|
||||
// function lets past the active-Dispatch check must clear the same worker check, or the Task
|
||||
// re-opens under a worker whose own lifecycle can no longer settle it (#16904 relay wedge).
|
||||
// A no-op re-assert of `dispatched` re-opens nothing and stays legal.
|
||||
const reopensUnderWorker = requiresActiveDispatch && task.status !== 'dispatched'
|
||||
const activeWorker =
|
||||
terminalStatus || reopensUnderWorker
|
||||
? (this.db
|
||||
.prepare(
|
||||
`SELECT active.id
|
||||
FROM dispatch_contexts active
|
||||
JOIN worker_dispatches worker ON worker.dispatch_id = active.id
|
||||
WHERE active.task_id = ? AND active.status IN ('pending', 'dispatched')
|
||||
AND worker.state NOT IN ('failed', 'succeeded', 'stopped', 'abandoned')
|
||||
ORDER BY active.rowid DESC LIMIT 1`
|
||||
)
|
||||
.get(id) as { id: string } | undefined)
|
||||
: undefined
|
||||
)
|
||||
.get(id) as { id: string } | undefined)
|
||||
: undefined
|
||||
if (activeWorker) {
|
||||
throw new OrchestrationError(
|
||||
'task_not_startable',
|
||||
|
||||
@@ -59,7 +59,11 @@ export function beginWorkerStop(
|
||||
this.db.exec('COMMIT')
|
||||
return { disposition: 'already_settled', worker, dispatch }
|
||||
}
|
||||
if (!['ready', 'start_unknown'].includes(worker.state)) {
|
||||
// Why `stopping` is accepted: a stop whose runtime died mid-flight leaves the row here
|
||||
// forever, and refusing the re-issue was the only operator escape (#16904). Re-running the
|
||||
// stop is what earns the honest outcome — settled, or `stop_unknown`, from which the worker
|
||||
// can be abandoned. It never asserts an exit the runtime did not observe.
|
||||
if (!['ready', 'start_unknown', 'stopping'].includes(worker.state)) {
|
||||
throw new OrchestrationError(
|
||||
'dispatch_inactive',
|
||||
`Dispatch ${dispatchId} cannot stop from ${worker.state}.`
|
||||
|
||||
@@ -51,6 +51,7 @@ export type LifecycleRejectionCode =
|
||||
| 'task_dispatch_mismatch'
|
||||
| 'inactive_dispatch'
|
||||
| 'stale_dispatch'
|
||||
| 'worker_not_settleable'
|
||||
|
||||
export type LifecycleRejectionResult = {
|
||||
action: 'rejected'
|
||||
|
||||
@@ -33,6 +33,7 @@ export type WorkerReportSettlement =
|
||||
| 'task_dispatch_mismatch'
|
||||
| 'inactive_dispatch'
|
||||
| 'stale_dispatch'
|
||||
| 'worker_not_settleable'
|
||||
reason: string
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user