mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
refactor(orchestration): delete the write-only lifecycle transition ledger
lifecycle_transition_receipts had no production reader: the append, the getter, the delete triggers, the reset paths and the bounded recovery retention all fed a table only tests read. The transition graph and its guards stay; v35 sheds the table and rebuilds the two delete triggers that survive it. Tests that read the ledger now assert the observable state change, and the atomicity tests inject their failure on the last real projection instead.
This commit is contained in:
@@ -45,10 +45,6 @@ export function releaseContextOnlyDispatch(
|
||||
last_failure: requestedState,
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString(),
|
||||
completed_at: dispatch.completed_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: {
|
||||
kind: `dispatch_context_only_${requestedState}`,
|
||||
details: { taskId: dispatch.task_id }
|
||||
}
|
||||
})
|
||||
const remaining = db
|
||||
@@ -67,11 +63,7 @@ export function releaseContextOnlyDispatch(
|
||||
entity: 'task',
|
||||
id: dispatch.task_id,
|
||||
from: 'dispatched',
|
||||
to: 'blocked',
|
||||
receipt: {
|
||||
kind: 'task_context_only_dispatch_released',
|
||||
details: { dispatchId: dispatch.id, requestedState }
|
||||
}
|
||||
to: 'blocked'
|
||||
}).changed
|
||||
}
|
||||
}
|
||||
|
||||
@@ -40,11 +40,6 @@ describe('Task/Dispatch invariant transactions', () => {
|
||||
expect(updated?.status).toBe(status)
|
||||
expect(db.getTask(task.id)?.status).toBe(status)
|
||||
expect(db.getTask(dependent.id)?.status).toBe(status === 'completed' ? 'ready' : 'pending')
|
||||
expect(db.getLifecycleTransitionReceipts('task', task.id)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({ from_state: 'pending', to_state: status })
|
||||
])
|
||||
)
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
@@ -194,33 +194,6 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
last_error: 'process exited'
|
||||
})
|
||||
expectCapability(database, worker, false)
|
||||
expect(database.getLifecycleTransitionReceipts('worker', worker.dispatchId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'stop_unknown',
|
||||
to_state: 'failed',
|
||||
kind: 'worker_process_exited'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(database.getLifecycleTransitionReceipts('dispatch', worker.dispatchId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'dispatched',
|
||||
to_state: 'failed',
|
||||
kind: 'dispatch_failed'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(database.getLifecycleTransitionReceipts('task', task.id)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'dispatched',
|
||||
to_state: 'blocked',
|
||||
kind: 'task_dispatch_interrupted'
|
||||
})
|
||||
])
|
||||
)
|
||||
})
|
||||
|
||||
it('keeps a Task dispatched when missing-terminal recovery leaves another worker active', () => {
|
||||
@@ -340,10 +313,10 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
expect(database.getDispatchContextById(started.dispatch.id)?.status).toBe('pending')
|
||||
expect(database.getTask(task.id)?.status).toBe('blocked')
|
||||
expect(database.getQuestion(question.message.id)?.status).toBe('pending')
|
||||
const receiptCounts = {
|
||||
worker: database.getLifecycleTransitionReceipts('worker', started.dispatch.id).length,
|
||||
dispatch: database.getLifecycleTransitionReceipts('dispatch', started.dispatch.id).length,
|
||||
task: database.getLifecycleTransitionReceipts('task', task.id).length
|
||||
const settled = {
|
||||
worker: database.getWorkerDispatch(started.dispatch.id),
|
||||
dispatch: database.getDispatchContextById(started.dispatch.id),
|
||||
task: database.getTask(task.id)
|
||||
}
|
||||
|
||||
database.reconcileFederatedWorkerStart({
|
||||
@@ -353,15 +326,10 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
lastError: 'worker server restarted'
|
||||
})
|
||||
|
||||
expect(database.getLifecycleTransitionReceipts('worker', started.dispatch.id)).toHaveLength(
|
||||
receiptCounts.worker
|
||||
)
|
||||
expect(database.getLifecycleTransitionReceipts('dispatch', started.dispatch.id)).toHaveLength(
|
||||
receiptCounts.dispatch
|
||||
)
|
||||
expect(database.getLifecycleTransitionReceipts('task', task.id)).toHaveLength(
|
||||
receiptCounts.task
|
||||
)
|
||||
// A repeated report of the same uncertainty must not re-project any of the three entities.
|
||||
expect(database.getWorkerDispatch(started.dispatch.id)).toEqual(settled.worker)
|
||||
expect(database.getDispatchContextById(started.dispatch.id)).toEqual(settled.dispatch)
|
||||
expect(database.getTask(task.id)).toEqual(settled.task)
|
||||
const answered = database.answerQuestion({
|
||||
messageId: question.message.id,
|
||||
runId: run.id,
|
||||
@@ -372,7 +340,7 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
expect(answered.message.body).toBe('yes')
|
||||
})
|
||||
|
||||
it('rolls back federated start uncertainty when the Task receipt cannot commit', () => {
|
||||
it('rolls back federated start uncertainty when the Task transition cannot commit', () => {
|
||||
const database = createDatabase()
|
||||
const task = database.createTask({ spec: 'atomic federated uncertainty' })
|
||||
const started = database.createStartingWorkerDispatch({
|
||||
@@ -388,10 +356,10 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
}
|
||||
})
|
||||
sqliteFor(database).exec(`
|
||||
CREATE TRIGGER reject_federated_unknown_task_receipt
|
||||
BEFORE INSERT ON lifecycle_transition_receipts
|
||||
WHEN NEW.kind = 'federated_task_start_unknown'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced federated uncertainty receipt failure'); END;
|
||||
CREATE TRIGGER reject_federated_unknown_task_block
|
||||
BEFORE UPDATE ON tasks
|
||||
WHEN NEW.status = 'blocked'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced federated uncertainty task block failure'); END;
|
||||
`)
|
||||
|
||||
expect(() =>
|
||||
@@ -401,7 +369,7 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
stage: 'remote_attach',
|
||||
lastError: 'worker server restarted'
|
||||
})
|
||||
).toThrow('forced federated uncertainty receipt failure')
|
||||
).toThrow('forced federated uncertainty task block failure')
|
||||
expect(database.getWorkerDispatch(started.dispatch.id)).toMatchObject({
|
||||
state: 'starting',
|
||||
stage: 'accepted',
|
||||
@@ -409,11 +377,6 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
})
|
||||
expect(database.getDispatchContextById(started.dispatch.id)?.status).toBe('pending')
|
||||
expect(database.getTask(task.id)?.status).toBe('dispatched')
|
||||
expect(
|
||||
database
|
||||
.getLifecycleTransitionReceipts('worker', started.dispatch.id)
|
||||
.some((receipt) => receipt.kind === 'federated_worker_start_unknown')
|
||||
).toBe(false)
|
||||
})
|
||||
|
||||
it.each(['stop', 'abandon'] as const)(
|
||||
@@ -472,40 +435,27 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
alreadySettled: false,
|
||||
releasedCurrentTask: true
|
||||
})
|
||||
expect(database.getLifecycleTransitionReceipts('dispatch', contextOnly.id)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'dispatched',
|
||||
to_state: 'failed',
|
||||
kind: `dispatch_context_only_${operation === 'stop' ? 'stopped' : 'abandoned'}`
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(database.getLifecycleTransitionReceipts('task', task.id)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'dispatched',
|
||||
to_state: 'blocked',
|
||||
kind: 'task_context_only_dispatch_released'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(database.getDispatchContextById(contextOnly.id)).toMatchObject({
|
||||
status: 'failed',
|
||||
last_failure: operation === 'stop' ? 'stopped' : 'abandoned'
|
||||
})
|
||||
expect(database.getTask(task.id)?.status).toBe('blocked')
|
||||
}
|
||||
)
|
||||
|
||||
it('rolls back both context-only projections when receipt append fails', () => {
|
||||
it('rolls back both context-only projections when the Task transition fails', () => {
|
||||
const database = createDatabase()
|
||||
const task = database.createTask({ spec: 'context-only atomic receipt' })
|
||||
const contextOnly = createRootDispatch(database, task.id, 'term_context')
|
||||
sqliteFor(database).exec(`
|
||||
CREATE TRIGGER reject_context_release_receipt
|
||||
BEFORE INSERT ON lifecycle_transition_receipts
|
||||
WHEN NEW.entity = 'task'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced context release receipt failure'); END;
|
||||
CREATE TRIGGER reject_context_release_task_block
|
||||
BEFORE UPDATE ON tasks
|
||||
WHEN NEW.status = 'blocked'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced context release task block failure'); END;
|
||||
`)
|
||||
|
||||
expect(() => database.beginWorkerStop(contextOnly.id, 'runtime_test')).toThrow(
|
||||
'forced context release receipt failure'
|
||||
'forced context release task block failure'
|
||||
)
|
||||
expect(database.getTask(task.id)?.status).toBe('dispatched')
|
||||
expect(database.getDispatchContextById(contextOnly.id)).toMatchObject({
|
||||
@@ -514,21 +464,6 @@ describe('Task/Dispatch lifecycle guards', () => {
|
||||
completed_at: null,
|
||||
capability_revoked_at: null
|
||||
})
|
||||
expect(database.getLifecycleTransitionReceipts('dispatch', contextOnly.id)).toEqual([])
|
||||
expect(database.getLifecycleTransitionReceipts('task', task.id)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'ready',
|
||||
to_state: 'dispatched',
|
||||
kind: 'task_dispatched'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(
|
||||
database
|
||||
.getLifecycleTransitionReceipts('task', task.id)
|
||||
.some((receipt) => receipt.kind === 'task_context_only_dispatch_released')
|
||||
).toBe(false)
|
||||
})
|
||||
|
||||
it.each(['stop', 'abandon'] as const)(
|
||||
|
||||
@@ -160,11 +160,10 @@ describe('Task/Dispatch concurrency', () => {
|
||||
result: 'completed concurrently'
|
||||
})
|
||||
expect(first.db.getWorkerDispatch(started.dispatch.id)?.state).toBe('succeeded')
|
||||
expect(
|
||||
first.db
|
||||
.getLifecycleTransitionReceipts('dispatch', started.dispatch.id)
|
||||
.map((receipt) => receipt.kind)
|
||||
).not.toContain('dispatch_failed')
|
||||
expect(first.db.getDispatchContextById(started.dispatch.id)).toMatchObject({
|
||||
status: 'completed',
|
||||
last_failure: null
|
||||
})
|
||||
expect(
|
||||
first.db.verifyDispatchCapability({
|
||||
dispatchId: started.dispatch.id,
|
||||
@@ -184,9 +183,7 @@ describe('Task/Dispatch concurrency', () => {
|
||||
sqlite.exec('BEGIN IMMEDIATE')
|
||||
expect(db.failDispatch(dispatch.id, 'nested failure')).toMatchObject({ status: 'failed' })
|
||||
expect(sqlite.isTransaction).toBe(true)
|
||||
expect(db.getLifecycleTransitionReceipts('dispatch', dispatch.id)).toEqual([
|
||||
expect.objectContaining({ kind: 'dispatch_failed' })
|
||||
])
|
||||
expect(db.getDispatchContextById(dispatch.id)?.status).toBe('failed')
|
||||
sqlite.exec('ROLLBACK')
|
||||
|
||||
expect(db.getTask(task.id)?.status).toBe('dispatched')
|
||||
@@ -195,7 +192,6 @@ describe('Task/Dispatch concurrency', () => {
|
||||
failure_count: 0,
|
||||
last_failure: null
|
||||
})
|
||||
expect(db.getLifecycleTransitionReceipts('dispatch', dispatch.id)).toEqual([])
|
||||
})
|
||||
|
||||
it('serializes reminted-pane worker authority claims', () => {
|
||||
|
||||
@@ -7,42 +7,33 @@ describe('decision-gate lifecycle transitions', () => {
|
||||
|
||||
afterEach(() => db?.close())
|
||||
|
||||
it('records the guarded Task transition when creating a gate', () => {
|
||||
it('blocks the dispatched Task when creating a gate', () => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
const task = db.createTask({ spec: 'gate receipt' })
|
||||
const task = db.createTask({ spec: 'gate blocks task' })
|
||||
createRootDispatch(db, task.id, 'term_gate')
|
||||
expect(db.getTask(task.id)?.status).toBe('dispatched')
|
||||
|
||||
db.createGate({ taskId: task.id, question: 'Proceed?' })
|
||||
|
||||
expect(db.getLifecycleTransitionReceipts('task', task.id)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'dispatched',
|
||||
to_state: 'blocked',
|
||||
kind: 'task_gate_created'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(db.getTask(task.id)?.status).toBe('blocked')
|
||||
})
|
||||
|
||||
it('rolls back gate and dispatch projections when receipt append fails', () => {
|
||||
it('rolls back the gate row when the Task transition cannot commit', () => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
const task = db.createTask({ spec: 'atomic gate creation' })
|
||||
const dispatch = createRootDispatch(db, task.id, 'term_gate')
|
||||
const taskReceiptsBefore = db.getLifecycleTransitionReceipts('task', task.id)
|
||||
db.db.exec(`
|
||||
CREATE TRIGGER reject_gate_lifecycle_receipt
|
||||
BEFORE INSERT ON lifecycle_transition_receipts
|
||||
WHEN NEW.kind = 'task_gate_created'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced gate receipt failure'); END;
|
||||
CREATE TRIGGER reject_gate_task_block
|
||||
BEFORE UPDATE ON tasks
|
||||
WHEN NEW.status = 'blocked'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced gate task block failure'); END;
|
||||
`)
|
||||
|
||||
expect(() => db!.createGate({ taskId: task.id, question: 'Proceed?' })).toThrow(
|
||||
'forced gate receipt failure'
|
||||
'forced gate task block failure'
|
||||
)
|
||||
expect(db.listGates({ taskId: task.id })).toHaveLength(0)
|
||||
expect(db.getTask(task.id)?.status).toBe('dispatched')
|
||||
expect(db.getDispatchContextById(dispatch.id)?.status).toBe('dispatched')
|
||||
expect(db.getLifecycleTransitionReceipts('task', task.id)).toEqual(taskReceiptsBefore)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -85,8 +85,7 @@ export function createGate(
|
||||
entity: 'task',
|
||||
id: gate.taskId,
|
||||
from: task.status,
|
||||
to: 'blocked',
|
||||
receipt: { kind: 'task_gate_created', details: { gateId: id } }
|
||||
to: 'blocked'
|
||||
})
|
||||
const created = this.db.prepare('SELECT * FROM decision_gates WHERE id = ?').get(id) as
|
||||
| DecisionGateRow
|
||||
|
||||
@@ -27,8 +27,7 @@ export function completeDispatch(this: OrchestrationDb, ctxId: string): void {
|
||||
projection: {
|
||||
completed_at: new Date().toISOString(),
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'dispatch_completed' }
|
||||
}
|
||||
})
|
||||
this.db.exec('RELEASE complete_dispatch_transition')
|
||||
} catch (error) {
|
||||
@@ -62,8 +61,7 @@ export function settleActiveDispatchesForTask(
|
||||
? (failure ?? row.last_failure ?? 'Task marked failed')
|
||||
: row.last_failure,
|
||||
capability_revoked_at: row.capability_revoked_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: `dispatch_${status}`, details: { taskId } }
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -160,8 +158,7 @@ export function failDispatch(
|
||||
termination_reason: options.terminationReason ?? before.termination_reason,
|
||||
completed_at: before.completed_at ?? new Date().toISOString(),
|
||||
capability_revoked_at: before.capability_revoked_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'dispatch_failed', details: { error } }
|
||||
}
|
||||
})
|
||||
const ctx = this.db.prepare('SELECT * FROM dispatch_contexts WHERE id = ?').get(ctxId) as
|
||||
| DispatchContextRow
|
||||
@@ -181,8 +178,7 @@ export function failDispatch(
|
||||
stage: 'process_exited',
|
||||
last_error: error,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'worker_process_exited' }
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -203,8 +199,7 @@ export function failDispatch(
|
||||
id: ctx.task_id,
|
||||
from: 'dispatched',
|
||||
to: taskStatus,
|
||||
projection: { completed_at: taskStatus === 'failed' ? new Date().toISOString() : null },
|
||||
receipt: { kind: 'task_dispatch_failed', details: { dispatchId: ctxId } }
|
||||
projection: { completed_at: taskStatus === 'failed' ? new Date().toISOString() : null }
|
||||
})
|
||||
}
|
||||
const updated = this.db.prepare('SELECT * FROM dispatch_contexts WHERE id = ?').get(ctxId) as
|
||||
|
||||
@@ -84,8 +84,7 @@ export function createDispatchContext(
|
||||
entity: 'task',
|
||||
id: taskId,
|
||||
from: 'ready',
|
||||
to: 'dispatched',
|
||||
receipt: { kind: 'task_dispatched', details: { dispatchId: id } }
|
||||
to: 'dispatched'
|
||||
})
|
||||
const dispatch = this.db
|
||||
.prepare('SELECT * FROM dispatch_contexts WHERE id = ?')
|
||||
|
||||
@@ -36,7 +36,6 @@ export function reconcileTaskAfterDispatchInterruption(
|
||||
entity: 'task',
|
||||
id: taskId,
|
||||
from: task.status,
|
||||
to: next,
|
||||
receipt: { kind: 'task_dispatch_interrupted', details: { dispatchId } }
|
||||
to: next
|
||||
})
|
||||
}
|
||||
|
||||
@@ -174,10 +174,6 @@ export function settleWorkerReportInTransaction(
|
||||
last_failure: params.outcome === 'failed' ? params.result : dispatch.last_failure,
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? now
|
||||
},
|
||||
receipt: {
|
||||
kind: 'dispatch_prompt_stall_corrected',
|
||||
details: { outcome: params.outcome }
|
||||
},
|
||||
correction: 'unobserved_prompt_report'
|
||||
})
|
||||
const taskTransition = transitionLifecycleWithDb(this.db, {
|
||||
@@ -186,10 +182,6 @@ export function settleWorkerReportInTransaction(
|
||||
from: 'failed',
|
||||
to: expectedTaskStatus,
|
||||
projection: { result: params.result, completed_at: now },
|
||||
receipt: {
|
||||
kind: 'task_prompt_stall_corrected',
|
||||
details: { dispatchId: params.dispatchId, outcome: params.outcome }
|
||||
},
|
||||
correction: 'unobserved_prompt_report'
|
||||
})
|
||||
dispatchUpdate = { changes: dispatchTransition.changed ? 1 : 0 }
|
||||
@@ -200,11 +192,7 @@ export function settleWorkerReportInTransaction(
|
||||
entity: 'task',
|
||||
id: params.taskId,
|
||||
from: 'blocked',
|
||||
to: 'dispatched',
|
||||
receipt: {
|
||||
kind: 'task_start_unknown_report_reconnected',
|
||||
details: { dispatchId: params.dispatchId }
|
||||
}
|
||||
to: 'dispatched'
|
||||
})
|
||||
}
|
||||
const dispatchTransition = transitionLifecycleWithDb(this.db, {
|
||||
@@ -216,16 +204,14 @@ export function settleWorkerReportInTransaction(
|
||||
completed_at: new Date().toISOString(),
|
||||
last_failure: params.outcome === 'failed' ? params.result : dispatch.last_failure,
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'dispatch_report_settled', details: { outcome: params.outcome } }
|
||||
}
|
||||
})
|
||||
const taskTransition = transitionLifecycleWithDb(this.db, {
|
||||
entity: 'task',
|
||||
id: params.taskId,
|
||||
from: 'dispatched',
|
||||
to: expectedTaskStatus,
|
||||
projection: { result: params.result, completed_at: new Date().toISOString() },
|
||||
receipt: { kind: 'task_report_settled', details: { dispatchId: params.dispatchId } }
|
||||
projection: { result: params.result, completed_at: new Date().toISOString() }
|
||||
})
|
||||
dispatchUpdate = { changes: dispatchTransition.changed ? 1 : 0 }
|
||||
taskUpdate = { changes: taskTransition.changed ? 1 : 0 }
|
||||
@@ -246,10 +232,6 @@ export function settleWorkerReportInTransaction(
|
||||
from: 'failed',
|
||||
to: params.outcome === 'succeeded' ? 'succeeded' : 'failed',
|
||||
projection: { stage: 'settled', updated_at: new Date().toISOString() },
|
||||
receipt: {
|
||||
kind: 'worker_prompt_stall_corrected',
|
||||
details: { outcome: params.outcome }
|
||||
},
|
||||
correction: 'unobserved_prompt_report'
|
||||
})
|
||||
} else if (reconnectingStart && params.outcome === 'succeeded') {
|
||||
@@ -257,16 +239,14 @@ export function settleWorkerReportInTransaction(
|
||||
entity: 'worker',
|
||||
id: params.dispatchId,
|
||||
from: 'start_unknown',
|
||||
to: 'ready',
|
||||
receipt: { kind: 'worker_start_unknown_report_reconnected' }
|
||||
to: 'ready'
|
||||
})
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'worker',
|
||||
id: params.dispatchId,
|
||||
from: 'ready',
|
||||
to: 'succeeded',
|
||||
projection: { stage: 'settled', updated_at: new Date().toISOString() },
|
||||
receipt: { kind: 'worker_report_settled', details: { outcome: params.outcome } }
|
||||
projection: { stage: 'settled', updated_at: new Date().toISOString() }
|
||||
})
|
||||
} else if (reportingWorker) {
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
@@ -275,8 +255,7 @@ export function settleWorkerReportInTransaction(
|
||||
// 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',
|
||||
projection: { stage: 'settled', updated_at: new Date().toISOString() },
|
||||
receipt: { kind: 'worker_report_settled', details: { outcome: params.outcome } }
|
||||
projection: { stage: 'settled', updated_at: new Date().toISOString() }
|
||||
})
|
||||
}
|
||||
settleActiveDispatchesForTask(
|
||||
|
||||
+1
-2
@@ -54,8 +54,7 @@ describe('federated Dispatch observation fence', () => {
|
||||
id: started.dispatch.id,
|
||||
from: 'ready',
|
||||
to: 'ready',
|
||||
projection: { stage: 'released', agent_terminal_handle: null },
|
||||
receipt: { kind: 'test_release' }
|
||||
projection: { stage: 'released', agent_terminal_handle: null }
|
||||
})
|
||||
database.db
|
||||
.prepare(
|
||||
|
||||
@@ -6,7 +6,7 @@ describe('guarded lifecycle transitions', () => {
|
||||
|
||||
afterEach(() => db?.close())
|
||||
|
||||
it('rejects a stale prior state without writing a receipt', () => {
|
||||
it('rejects a stale prior state without changing the projection', () => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
const task = db.createTask({ spec: 'guarded transition' })
|
||||
|
||||
@@ -18,32 +18,27 @@ describe('guarded lifecycle transitions', () => {
|
||||
to: 'completed'
|
||||
})
|
||||
).toThrow(/expected pending/)
|
||||
expect(db.getLifecycleTransitionReceipts('task', task.id)).toEqual([])
|
||||
expect(db.getTask(task.id)?.status).toBe('ready')
|
||||
})
|
||||
|
||||
it('rolls back the legacy projection when receipt append fails', () => {
|
||||
it('composes its projection into the caller-owned transaction', () => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
const task = db.createTask({ spec: 'atomic receipt' })
|
||||
db.db.exec(`
|
||||
CREATE TRIGGER reject_lifecycle_receipt
|
||||
BEFORE INSERT ON lifecycle_transition_receipts
|
||||
BEGIN SELECT RAISE(ABORT, 'forced receipt failure'); END;
|
||||
`)
|
||||
const task = db.createTask({ spec: 'caller-owned rollback' })
|
||||
|
||||
db.db.exec('SAVEPOINT lifecycle_test')
|
||||
expect(() =>
|
||||
db!.transitionLifecycle({
|
||||
expect(
|
||||
db.transitionLifecycle({
|
||||
entity: 'task',
|
||||
id: task.id,
|
||||
from: 'ready',
|
||||
to: 'completed',
|
||||
receipt: { kind: 'test' }
|
||||
projection: { result: 'uncommitted' }
|
||||
})
|
||||
).toThrow('forced receipt failure')
|
||||
).toEqual({ changed: true })
|
||||
expect(db.getTask(task.id)?.status).toBe('completed')
|
||||
db.db.exec('ROLLBACK TO lifecycle_test')
|
||||
db.db.exec('RELEASE lifecycle_test')
|
||||
|
||||
expect(db.getTask(task.id)?.status).toBe('ready')
|
||||
expect(db.getLifecycleTransitionReceipts('task', task.id)).toEqual([])
|
||||
expect(db.getTask(task.id)).toMatchObject({ status: 'ready', result: null })
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,15 +1,14 @@
|
||||
import type Database from '../../../sqlite/sync-database'
|
||||
import { OrchestrationError } from '../orchestration-error'
|
||||
import type { OrchestrationDb } from './orchestration-db'
|
||||
import { generateId } from './generated-id'
|
||||
|
||||
/**
|
||||
* The single write boundary for Task, Dispatch, and supervised worker state.
|
||||
*
|
||||
* This function deliberately does not open or commit a transaction. Callers
|
||||
* often compose several projections (and a mailbox effect) in one transaction;
|
||||
* keeping the boundary neutral makes the receipt and projection atomic with
|
||||
* that caller-owned transaction.
|
||||
* keeping the boundary neutral makes every projection atomic with that
|
||||
* caller-owned transaction.
|
||||
*/
|
||||
export type LifecycleEntity = 'task' | 'dispatch' | 'worker'
|
||||
|
||||
@@ -55,8 +54,6 @@ export type LifecycleTransitionParams = {
|
||||
to: string
|
||||
/** Additional legacy projection columns written with the state change. */
|
||||
projection?: Record<string, string | number | null>
|
||||
/** Stable event name and optional details retained for replay/audit. */
|
||||
receipt?: { kind?: string; details?: unknown }
|
||||
/** Narrow exception for a worker report correcting an unobserved prompt start. */
|
||||
correction?: 'unobserved_prompt_report'
|
||||
}
|
||||
@@ -121,21 +118,10 @@ const PROJECTION_COLUMNS = new Set([
|
||||
'runtime_epoch'
|
||||
])
|
||||
|
||||
export type LifecycleTransitionReceipt = {
|
||||
id: string
|
||||
entity: LifecycleEntity
|
||||
entity_id: string
|
||||
from_state: string
|
||||
to_state: string
|
||||
kind: string
|
||||
details: string | null
|
||||
created_at: string
|
||||
}
|
||||
|
||||
export function transitionLifecycle(
|
||||
this: OrchestrationDb,
|
||||
params: LifecycleTransitionParams
|
||||
): { changed: boolean; receipt?: LifecycleTransitionReceipt } {
|
||||
): { changed: boolean } {
|
||||
return transitionLifecycleWithDb(this.db, params)
|
||||
}
|
||||
|
||||
@@ -143,7 +129,7 @@ export function transitionLifecycle(
|
||||
export function transitionLifecycleWithDb(
|
||||
db: Database.Database,
|
||||
params: LifecycleTransitionParams
|
||||
): { changed: boolean; receipt?: LifecycleTransitionReceipt } {
|
||||
): { changed: boolean } {
|
||||
const entity = ENTITY_TABLE[params.entity]
|
||||
const allowed = Array.isArray(params.from) ? params.from : [params.from]
|
||||
const current = db
|
||||
@@ -207,82 +193,13 @@ export function transitionLifecycleWithDb(
|
||||
)
|
||||
}
|
||||
|
||||
const receipt = appendLifecycleTransitionReceipt(db, {
|
||||
entity: params.entity,
|
||||
entityId: params.id,
|
||||
fromState: current.state,
|
||||
toState: params.to,
|
||||
kind: params.receipt?.kind,
|
||||
details: params.receipt?.details
|
||||
})
|
||||
return { changed: true, receipt }
|
||||
}
|
||||
|
||||
function appendLifecycleTransitionReceipt(
|
||||
db: Database.Database,
|
||||
params: {
|
||||
entity: LifecycleEntity
|
||||
entityId: string
|
||||
fromState: string
|
||||
toState: string
|
||||
kind?: string
|
||||
details?: unknown
|
||||
}
|
||||
): LifecycleTransitionReceipt {
|
||||
const receiptId = generateId('lcr')
|
||||
db.prepare(
|
||||
`INSERT INTO lifecycle_transition_receipts
|
||||
(id, entity, entity_id, from_state, to_state, kind, details)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)`
|
||||
).run(
|
||||
receiptId,
|
||||
params.entity,
|
||||
params.entityId,
|
||||
params.fromState,
|
||||
params.toState,
|
||||
params.kind ?? 'lifecycle_transition',
|
||||
params.details === undefined ? null : JSON.stringify(params.details)
|
||||
)
|
||||
return db
|
||||
.prepare('SELECT * FROM lifecycle_transition_receipts WHERE id = ?')
|
||||
.get(receiptId) as LifecycleTransitionReceipt
|
||||
}
|
||||
|
||||
export function getLifecycleTransitionReceipts(
|
||||
this: OrchestrationDb,
|
||||
entity: LifecycleEntity,
|
||||
entityId: string
|
||||
): LifecycleTransitionReceipt[] {
|
||||
const receipts = this.db
|
||||
.prepare(
|
||||
`SELECT * FROM lifecycle_transition_receipts
|
||||
WHERE entity = ? AND entity_id = ? ORDER BY rowid`
|
||||
)
|
||||
.all(entity, entityId) as LifecycleTransitionReceipt[]
|
||||
if (receipts.length > 0 || entity !== 'worker') {
|
||||
return receipts
|
||||
}
|
||||
// Recovery receipts are keyed by owning Dispatch. Accept the historical Resource-ID query
|
||||
// shape so callers inspecting older releases continue to see the same ledger entries.
|
||||
const owner = this.db
|
||||
.prepare('SELECT owner_dispatch_id FROM worker_terminal_resources WHERE id = ?')
|
||||
.get(entityId) as { owner_dispatch_id: string } | undefined
|
||||
if (!owner) {
|
||||
return receipts
|
||||
}
|
||||
return this.db
|
||||
.prepare(
|
||||
`SELECT * FROM lifecycle_transition_receipts
|
||||
WHERE entity = ? AND entity_id = ? AND kind = 'worker_terminal_recovery' ORDER BY rowid`
|
||||
)
|
||||
.all(entity, owner.owner_dispatch_id) as LifecycleTransitionReceipt[]
|
||||
return { changed: true }
|
||||
}
|
||||
|
||||
export type LifecycleTransitionMethods = {
|
||||
transitionLifecycle: typeof transitionLifecycle
|
||||
getLifecycleTransitionReceipts: typeof getLifecycleTransitionReceipts
|
||||
}
|
||||
|
||||
export function attachLifecycleTransition(ctor: { prototype: object }): void {
|
||||
Object.assign(ctor.prototype, { transitionLifecycle, getLifecycleTransitionReceipts })
|
||||
Object.assign(ctor.prototype, { transitionLifecycle })
|
||||
}
|
||||
|
||||
@@ -36,7 +36,6 @@ export function resetAll(this: OrchestrationDb): void {
|
||||
DELETE FROM worker_terminal_archives;
|
||||
DELETE FROM worker_terminal_resources;
|
||||
DELETE FROM attempt_observation_facts;
|
||||
DELETE FROM lifecycle_transition_receipts;
|
||||
DELETE FROM worker_dispatches;
|
||||
DELETE FROM dispatch_contexts;
|
||||
DELETE FROM tasks;
|
||||
@@ -69,7 +68,6 @@ export function resetTasks(this: OrchestrationDb): void {
|
||||
DELETE FROM worker_terminal_archives;
|
||||
DELETE FROM worker_terminal_resources;
|
||||
DELETE FROM attempt_observation_facts;
|
||||
DELETE FROM lifecycle_transition_receipts;
|
||||
DELETE FROM worker_dispatches;
|
||||
DELETE FROM dispatch_contexts;
|
||||
DELETE FROM tasks;
|
||||
|
||||
@@ -111,22 +111,6 @@ CREATE TABLE IF NOT EXISTS mutation_caller_identities (
|
||||
caller_fingerprint TEXT NOT NULL UNIQUE
|
||||
);
|
||||
|
||||
-- Append-only lifecycle audit facts. Legacy status enums remain unchanged;
|
||||
-- richer observation/projection work can consume this table in a later slice.
|
||||
CREATE TABLE IF NOT EXISTS lifecycle_transition_receipts (
|
||||
id TEXT PRIMARY KEY,
|
||||
entity TEXT NOT NULL CHECK(entity IN ('task', 'dispatch', 'worker')),
|
||||
entity_id TEXT NOT NULL,
|
||||
from_state TEXT NOT NULL,
|
||||
to_state TEXT NOT NULL,
|
||||
kind TEXT NOT NULL DEFAULT 'lifecycle_transition',
|
||||
details TEXT,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_lifecycle_transition_receipts_entity
|
||||
ON lifecycle_transition_receipts(entity, entity_id, created_at);
|
||||
|
||||
-- Attempt evidence stays additive so old Task/Dispatch/worker CHECK enums remain wire-compatible.
|
||||
CREATE TABLE IF NOT EXISTS attempt_observation_facts (
|
||||
id TEXT PRIMARY KEY,
|
||||
|
||||
@@ -5,6 +5,21 @@ import {
|
||||
} from '../pane-key-match'
|
||||
import { potentiallyLiveRemoteAttachmentSql } from '../federation/remote-attachment-liveness'
|
||||
|
||||
// Additive tables outlive v30 writers, so legacy parent deletes must clean their rows too.
|
||||
export const ADDITIVE_LIFECYCLE_DELETE_TRIGGERS_SQL = `
|
||||
CREATE TRIGGER IF NOT EXISTS trg_tasks_delete_additive_lifecycle
|
||||
AFTER DELETE ON tasks
|
||||
BEGIN
|
||||
DELETE FROM attempt_observation_facts WHERE task_id = OLD.id;
|
||||
END;
|
||||
|
||||
CREATE TRIGGER IF NOT EXISTS trg_dispatches_delete_additive_lifecycle
|
||||
AFTER DELETE ON dispatch_contexts
|
||||
BEGIN
|
||||
DELETE FROM attempt_observation_facts WHERE dispatch_id = OLD.id;
|
||||
END;
|
||||
`
|
||||
|
||||
export function createGraphTablesSql(): string {
|
||||
return `
|
||||
CREATE TABLE IF NOT EXISTS federated_dispatches (
|
||||
@@ -157,29 +172,7 @@ CREATE INDEX IF NOT EXISTS idx_dispatch_task ON dispatch_contexts(task_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_dispatch_status ON dispatch_contexts(status);
|
||||
CREATE INDEX IF NOT EXISTS idx_dispatch_assignee_handle ON dispatch_contexts(assignee_handle);
|
||||
|
||||
-- Additive tables outlive v30 writers, so legacy parent deletes must clean their rows too.
|
||||
CREATE TRIGGER IF NOT EXISTS trg_tasks_delete_additive_lifecycle
|
||||
AFTER DELETE ON tasks
|
||||
BEGIN
|
||||
DELETE FROM lifecycle_transition_receipts
|
||||
WHERE entity = 'task' AND entity_id = OLD.id;
|
||||
DELETE FROM attempt_observation_facts WHERE task_id = OLD.id;
|
||||
END;
|
||||
|
||||
CREATE TRIGGER IF NOT EXISTS trg_dispatches_delete_additive_lifecycle
|
||||
AFTER DELETE ON dispatch_contexts
|
||||
BEGIN
|
||||
DELETE FROM lifecycle_transition_receipts
|
||||
WHERE entity IN ('dispatch', 'worker') AND entity_id = OLD.id;
|
||||
DELETE FROM attempt_observation_facts WHERE dispatch_id = OLD.id;
|
||||
END;
|
||||
|
||||
CREATE TRIGGER IF NOT EXISTS trg_workers_delete_additive_lifecycle
|
||||
AFTER DELETE ON worker_dispatches
|
||||
BEGIN
|
||||
DELETE FROM lifecycle_transition_receipts
|
||||
WHERE entity = 'worker' AND entity_id = OLD.dispatch_id;
|
||||
END;
|
||||
${ADDITIVE_LIFECYCLE_DELETE_TRIGGERS_SQL}
|
||||
|
||||
CREATE TABLE IF NOT EXISTS decision_gates (
|
||||
id TEXT PRIMARY KEY,
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import type { OrchestrationDb } from '../orchestration-db'
|
||||
import { ADDITIVE_LIFECYCLE_DELETE_TRIGGERS_SQL } from './create-graph-tables-sql'
|
||||
|
||||
const ONE_OUTSTANDING_INDEX_SQL = `
|
||||
CREATE UNIQUE INDEX idx_deliveries_one_outstanding
|
||||
@@ -19,6 +20,16 @@ export function migrateMailboxDeliveryRepairV35(this: OrchestrationDb, current:
|
||||
if (current >= 35) {
|
||||
return
|
||||
}
|
||||
// The write-only lifecycle ledger is gone. Old delete triggers still reference it, and
|
||||
// CREATE TRIGGER IF NOT EXISTS cannot replace a body, so drop all three and rebuild the two
|
||||
// that survive.
|
||||
this.db.exec(`
|
||||
DROP TRIGGER IF EXISTS trg_tasks_delete_additive_lifecycle;
|
||||
DROP TRIGGER IF EXISTS trg_dispatches_delete_additive_lifecycle;
|
||||
DROP TRIGGER IF EXISTS trg_workers_delete_additive_lifecycle;
|
||||
DROP TABLE IF EXISTS lifecycle_transition_receipts;
|
||||
${ADDITIVE_LIFECYCLE_DELETE_TRIGGERS_SQL}
|
||||
`)
|
||||
rebuildDeliveriesWithMailboxDefault.call(this)
|
||||
recreateIndexMissingPredicate.call(
|
||||
this,
|
||||
|
||||
@@ -81,8 +81,7 @@ export function updateTaskStatus(
|
||||
projection: {
|
||||
result: result ?? task.result,
|
||||
completed_at: terminalStatus ? new Date().toISOString() : task.completed_at
|
||||
},
|
||||
receipt: { kind: 'task_status', details: { result: result ?? null } }
|
||||
}
|
||||
})
|
||||
} catch (error) {
|
||||
if (!(error instanceof OrchestrationError) || error.code !== 'lifecycle_conflict') {
|
||||
|
||||
@@ -214,8 +214,7 @@ export function promoteReadyTasks(this: OrchestrationDb, completedTaskId: string
|
||||
entity: 'task',
|
||||
id: task.id,
|
||||
from: 'pending',
|
||||
to: 'ready',
|
||||
receipt: { kind: 'task_ready', details: { dependency: completedTaskId } }
|
||||
to: 'ready'
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
+9
-18
@@ -55,16 +55,14 @@ export function reconcileFederatedWorkerStart(
|
||||
: worker.residual_resources,
|
||||
last_error: null,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'federated_worker_ready' }
|
||||
}
|
||||
})
|
||||
if (dispatch.status === 'pending') {
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'dispatch',
|
||||
id: params.dispatchId,
|
||||
from: 'pending',
|
||||
to: 'dispatched',
|
||||
receipt: { kind: 'federated_dispatch_ready' }
|
||||
to: 'dispatched'
|
||||
})
|
||||
}
|
||||
const task = this.getTask(dispatch.task_id)
|
||||
@@ -74,8 +72,7 @@ export function reconcileFederatedWorkerStart(
|
||||
id: dispatch.task_id,
|
||||
from: 'blocked',
|
||||
to: 'dispatched',
|
||||
projection: { completed_at: null },
|
||||
receipt: { kind: 'federated_task_ready' }
|
||||
projection: { completed_at: null }
|
||||
})
|
||||
}
|
||||
} else if (params.state === 'start_unknown') {
|
||||
@@ -90,15 +87,13 @@ export function reconcileFederatedWorkerStart(
|
||||
stage: params.stage,
|
||||
last_error: reason,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'federated_worker_start_unknown', details: { reason } }
|
||||
}
|
||||
})
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'dispatch',
|
||||
id: params.dispatchId,
|
||||
from: dispatch.status,
|
||||
to: dispatch.status,
|
||||
receipt: { kind: 'federated_dispatch_start_unknown' }
|
||||
to: dispatch.status
|
||||
})
|
||||
}
|
||||
const task = this.getTask(dispatch.task_id)
|
||||
@@ -107,8 +102,7 @@ export function reconcileFederatedWorkerStart(
|
||||
entity: 'task',
|
||||
id: dispatch.task_id,
|
||||
from: 'dispatched',
|
||||
to: 'blocked',
|
||||
receipt: { kind: 'federated_task_start_unknown', details: { reason } }
|
||||
to: 'blocked'
|
||||
})
|
||||
}
|
||||
} else {
|
||||
@@ -122,8 +116,7 @@ export function reconcileFederatedWorkerStart(
|
||||
stage: params.stage,
|
||||
last_error: reason,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: `federated_worker_${params.state}`, details: { reason } }
|
||||
}
|
||||
})
|
||||
if (['pending', 'dispatched'].includes(dispatch.status)) {
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
@@ -135,8 +128,7 @@ export function reconcileFederatedWorkerStart(
|
||||
last_failure: reason,
|
||||
completed_at: new Date().toISOString(),
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'federated_dispatch_failed', details: { reason } }
|
||||
}
|
||||
})
|
||||
}
|
||||
reconcileTaskAfterDispatchInterruption(this, dispatch.task_id, params.dispatchId)
|
||||
@@ -155,8 +147,7 @@ export function reconcileFederatedWorkerStart(
|
||||
id: dispatch.task_id,
|
||||
from: task.status,
|
||||
to: 'failed',
|
||||
projection: { completed_at: new Date().toISOString() },
|
||||
receipt: { kind: 'federated_task_failed' }
|
||||
projection: { completed_at: new Date().toISOString() }
|
||||
})
|
||||
}
|
||||
this.closeQuestionsForDispatch(params.dispatchId)
|
||||
|
||||
@@ -61,8 +61,7 @@ export function abandonWorkerDispatch(
|
||||
id: dispatchId,
|
||||
from: worker.state,
|
||||
to: 'abandoned',
|
||||
projection: { stage: 'abandoned', updated_at: new Date().toISOString() },
|
||||
receipt: { kind: 'worker_abandoned' }
|
||||
projection: { stage: 'abandoned', updated_at: new Date().toISOString() }
|
||||
})
|
||||
if (['pending', 'dispatched'].includes(dispatch.status)) {
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
@@ -73,8 +72,7 @@ export function abandonWorkerDispatch(
|
||||
projection: {
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString(),
|
||||
completed_at: dispatch.completed_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'dispatch_abandoned' }
|
||||
}
|
||||
})
|
||||
}
|
||||
reconcileTaskAfterDispatchInterruption(this, dispatch.task_id, dispatchId)
|
||||
|
||||
@@ -19,8 +19,7 @@ export function markWorkerDispatchReady(
|
||||
entity: 'dispatch',
|
||||
id: dispatchId,
|
||||
from: 'pending',
|
||||
to: 'dispatched',
|
||||
receipt: { kind: 'dispatch_ready' }
|
||||
to: 'dispatched'
|
||||
})
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'worker',
|
||||
@@ -30,8 +29,7 @@ export function markWorkerDispatchReady(
|
||||
projection: {
|
||||
stage: 'input_accepted',
|
||||
effects: effects ? JSON.stringify(effects) : worker.effects
|
||||
},
|
||||
receipt: { kind: 'worker_ready' }
|
||||
}
|
||||
})
|
||||
this.db.exec('COMMIT')
|
||||
return this.getWorkerDispatch(dispatchId) as WorkerDispatchRow
|
||||
@@ -70,16 +68,14 @@ export function failWorkerStart(
|
||||
capability_revoked_at: options.retainCapability
|
||||
? dispatch.capability_revoked_at
|
||||
: (dispatch.capability_revoked_at ?? now)
|
||||
},
|
||||
receipt: { kind: 'dispatch_start_failed', details: { reason } }
|
||||
}
|
||||
})
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'worker',
|
||||
id: dispatchId,
|
||||
from: 'starting',
|
||||
to: 'failed',
|
||||
projection: { stage, last_error: reason, updated_at: now },
|
||||
receipt: { kind: 'worker_start_failed', details: { reason } }
|
||||
projection: { stage, last_error: reason, updated_at: now }
|
||||
})
|
||||
const hasActiveDispatch = Boolean(
|
||||
this.db
|
||||
@@ -96,8 +92,7 @@ export function failWorkerStart(
|
||||
id: dispatch.task_id,
|
||||
from: task.status,
|
||||
to: 'failed',
|
||||
projection: { completed_at: now },
|
||||
receipt: { kind: 'task_start_failed', details: { reason } }
|
||||
projection: { completed_at: now }
|
||||
})
|
||||
}
|
||||
this.closeQuestionsForDispatch(dispatchId)
|
||||
@@ -127,22 +122,19 @@ export function markWorkerStartUnknown(
|
||||
id: dispatchId,
|
||||
from: 'starting',
|
||||
to: 'start_unknown',
|
||||
projection: { stage, last_error: reason, updated_at: new Date().toISOString() },
|
||||
receipt: { kind: 'worker_start_unknown', details: { reason } }
|
||||
projection: { stage, last_error: reason, updated_at: new Date().toISOString() }
|
||||
})
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'dispatch',
|
||||
id: dispatchId,
|
||||
from: dispatch.status,
|
||||
to: dispatch.status,
|
||||
receipt: { kind: 'dispatch_start_unknown' }
|
||||
to: dispatch.status
|
||||
})
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'task',
|
||||
id: dispatch.task_id,
|
||||
from: 'dispatched',
|
||||
to: 'blocked',
|
||||
receipt: { kind: 'task_start_unknown', details: { reason } }
|
||||
to: 'blocked'
|
||||
})
|
||||
this.closeQuestionsForDispatch(dispatchId)
|
||||
this.db.exec('COMMIT')
|
||||
|
||||
@@ -42,8 +42,7 @@ export function recordWorkerStage(
|
||||
: current.residual_resources,
|
||||
last_error: params.lastError ?? current.last_error,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'worker_stage', details: { stage: params.stage } }
|
||||
}
|
||||
})
|
||||
this.db.exec('RELEASE worker_stage_transition')
|
||||
} catch (error) {
|
||||
@@ -84,8 +83,7 @@ export function updateWorkerSetupEvidence(
|
||||
setup_state: params.setupState,
|
||||
effects,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'worker_setup_evidence' }
|
||||
}
|
||||
})
|
||||
this.db.exec('RELEASE worker_setup_transition')
|
||||
} catch (error) {
|
||||
|
||||
@@ -164,8 +164,7 @@ export function createStartingWorkerDispatch(
|
||||
id: task.id,
|
||||
from: params.retryOf ? ['failed', 'blocked'] : 'ready',
|
||||
to: 'dispatched',
|
||||
projection: { result: null, completed_at: null },
|
||||
receipt: { kind: 'task_dispatched', details: { dispatchId: id } }
|
||||
projection: { result: null, completed_at: null }
|
||||
})
|
||||
this.db.exec('COMMIT')
|
||||
this.hasAnyDispatchContextsCache = true
|
||||
|
||||
@@ -74,8 +74,7 @@ export function beginWorkerStop(
|
||||
stage: 'stop_requested',
|
||||
runtime_epoch: runtimeEpoch,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'worker_stop_requested' }
|
||||
}
|
||||
})
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
entity: 'dispatch',
|
||||
@@ -84,8 +83,7 @@ export function beginWorkerStop(
|
||||
to: dispatch.status,
|
||||
projection: {
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'dispatch_capability_revoked' }
|
||||
}
|
||||
})
|
||||
reconcileTaskAfterDispatchInterruption(this, dispatch.task_id, dispatchId)
|
||||
this.closeQuestionsForDispatch(dispatchId)
|
||||
@@ -114,8 +112,7 @@ export function settleWorkerStop(this: OrchestrationDb, dispatchId: string): Wor
|
||||
id: dispatchId,
|
||||
from: 'stopping',
|
||||
to: 'stopped',
|
||||
projection: { stage: 'process_stopped', updated_at: new Date().toISOString() },
|
||||
receipt: { kind: 'worker_stopped' }
|
||||
projection: { stage: 'process_stopped', updated_at: new Date().toISOString() }
|
||||
})
|
||||
if (['pending', 'dispatched'].includes(dispatch.status)) {
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
@@ -123,8 +120,7 @@ export function settleWorkerStop(this: OrchestrationDb, dispatchId: string): Wor
|
||||
id: dispatchId,
|
||||
from: dispatch.status,
|
||||
to: 'failed',
|
||||
projection: { completed_at: new Date().toISOString(), last_failure: 'stopped' },
|
||||
receipt: { kind: 'dispatch_stopped' }
|
||||
projection: { completed_at: new Date().toISOString(), last_failure: 'stopped' }
|
||||
})
|
||||
}
|
||||
reconcileTaskAfterDispatchInterruption(this, dispatch.task_id, dispatchId)
|
||||
@@ -169,8 +165,7 @@ export function reconcileFederatedWorkerStop(
|
||||
stage: 'process_stopped',
|
||||
last_error: null,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'federated_worker_stopped' }
|
||||
}
|
||||
})
|
||||
if (['pending', 'dispatched'].includes(dispatch.status)) {
|
||||
transitionLifecycleWithDb(this.db, {
|
||||
@@ -181,8 +176,7 @@ export function reconcileFederatedWorkerStop(
|
||||
projection: {
|
||||
completed_at: dispatch.completed_at ?? new Date().toISOString(),
|
||||
last_failure: 'stopped'
|
||||
},
|
||||
receipt: { kind: 'federated_dispatch_stopped' }
|
||||
}
|
||||
})
|
||||
}
|
||||
reconcileTaskAfterDispatchInterruption(this, dispatch.task_id, dispatchId)
|
||||
@@ -210,8 +204,7 @@ export function resumeFederatedWorkerForTerminalRelay(
|
||||
id: dispatchId,
|
||||
from: 'stopping',
|
||||
to: 'ready',
|
||||
projection: { stage: 'remote_report_pending', updated_at: new Date().toISOString() },
|
||||
receipt: { kind: 'worker_stop_relay_resumed' }
|
||||
projection: { stage: 'remote_report_pending', updated_at: new Date().toISOString() }
|
||||
})
|
||||
const task = this.getTask(dispatch.task_id)
|
||||
if (task?.status === 'blocked') {
|
||||
@@ -219,8 +212,7 @@ export function resumeFederatedWorkerForTerminalRelay(
|
||||
entity: 'task',
|
||||
id: dispatch.task_id,
|
||||
from: 'blocked',
|
||||
to: 'dispatched',
|
||||
receipt: { kind: 'task_stop_relay_resumed' }
|
||||
to: 'dispatched'
|
||||
})
|
||||
}
|
||||
this.db.exec('COMMIT')
|
||||
@@ -251,8 +243,7 @@ export function markWorkerStopUnknown(
|
||||
stage: 'stop_outcome_unknown',
|
||||
last_error: reason,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'worker_stop_unknown', details: { reason } }
|
||||
}
|
||||
})
|
||||
this.db.exec('RELEASE mark_worker_stop_unknown')
|
||||
return this.getWorkerDispatch(dispatchId) as WorkerDispatchRow
|
||||
|
||||
@@ -70,8 +70,7 @@ export function reconcileMissingWorkerTerminal(
|
||||
last_failure: reason,
|
||||
completed_at: new Date().toISOString(),
|
||||
capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'dispatch_terminal_missing', details: { reason } }
|
||||
}
|
||||
})
|
||||
if (!stopWasPending) {
|
||||
const taskStatus: TaskStatus = dispatchStatus === 'circuit_broken' ? 'failed' : 'ready'
|
||||
@@ -91,8 +90,7 @@ export function reconcileMissingWorkerTerminal(
|
||||
id: dispatch.task_id,
|
||||
from: task.status,
|
||||
to: taskStatus,
|
||||
projection: { completed_at: taskStatus === 'failed' ? new Date().toISOString() : null },
|
||||
receipt: { kind: 'task_terminal_missing' }
|
||||
projection: { completed_at: taskStatus === 'failed' ? new Date().toISOString() : null }
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -107,11 +105,6 @@ export function reconcileMissingWorkerTerminal(
|
||||
stage: 'terminal_missing',
|
||||
last_error: reason,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: {
|
||||
kind: stopWasPending
|
||||
? 'worker_terminal_missing_stopped'
|
||||
: 'worker_terminal_missing_abandoned'
|
||||
}
|
||||
})
|
||||
this.db.exec('COMMIT')
|
||||
|
||||
+2
-32
@@ -126,8 +126,7 @@ export function getWorkerTerminalResourceFormerlyOwnedBy(
|
||||
/** Records bounded recovery bookkeeping without changing ownership or release intent. */
|
||||
export function recordWorkerTerminalRecoveryAttempt(
|
||||
this: OrchestrationDb,
|
||||
resourceId: string,
|
||||
outcome: 'released' | 'pending' | 'unknown' | 'retained'
|
||||
resourceId: string
|
||||
): WorkerTerminalResourceRow | undefined {
|
||||
this.db
|
||||
.prepare(
|
||||
@@ -137,36 +136,7 @@ export function recordWorkerTerminalRecoveryAttempt(
|
||||
WHERE id = ?`
|
||||
)
|
||||
.run(resourceId)
|
||||
const resource = this.getWorkerTerminalResource(resourceId)
|
||||
if (resource) {
|
||||
// Keep the receipt in the worker lifecycle ledger keyed by its owning Dispatch, not the
|
||||
// terminal Resource ID. Resource-ID lookups remain compatible in getLifecycleTransitionReceipts.
|
||||
this.db
|
||||
.prepare(
|
||||
`INSERT INTO lifecycle_transition_receipts
|
||||
(id, entity, entity_id, from_state, to_state, kind, details)
|
||||
VALUES (?, 'worker', ?, ?, ?, 'worker_terminal_recovery', ?)`
|
||||
)
|
||||
.run(
|
||||
generateId('wrr'),
|
||||
resource.owner_dispatch_id,
|
||||
resource.release_state,
|
||||
outcome,
|
||||
JSON.stringify({ attempt: resource.recovery_attempt_count })
|
||||
)
|
||||
this.db
|
||||
.prepare(
|
||||
`DELETE FROM lifecycle_transition_receipts
|
||||
WHERE kind = 'worker_terminal_recovery' AND entity_id = ?
|
||||
AND id NOT IN (
|
||||
SELECT id FROM lifecycle_transition_receipts
|
||||
WHERE kind = 'worker_terminal_recovery' AND entity_id = ?
|
||||
ORDER BY created_at DESC LIMIT 32
|
||||
)`
|
||||
)
|
||||
.run(resource.owner_dispatch_id, resource.owner_dispatch_id)
|
||||
}
|
||||
return resource
|
||||
return this.getWorkerTerminalResource(resourceId)
|
||||
}
|
||||
|
||||
// Reusable exact settled terminal: transfers cleanup ownership to the new Dispatch and fences
|
||||
|
||||
@@ -102,12 +102,6 @@ describe('lifecycle reconciliation', () => {
|
||||
expect(db.getTask(task.id)?.status).toBe('completed')
|
||||
expect(db.getDispatchContextById(started.dispatch.id)?.status).toBe('completed')
|
||||
expect(db.getWorkerDispatch(started.dispatch.id)?.state).toBe('succeeded')
|
||||
expect(db.getLifecycleTransitionReceipts('task', task.id)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({ from_state: 'blocked', to_state: 'dispatched' }),
|
||||
expect.objectContaining({ from_state: 'dispatched', to_state: 'completed' })
|
||||
])
|
||||
)
|
||||
})
|
||||
|
||||
it('fails both the dispatch and task from an authenticated failed worker report', () => {
|
||||
|
||||
@@ -58,7 +58,6 @@ describe('OrchestrationDb reset scopes', () => {
|
||||
|
||||
it('resetAll clears Runs, worker/federation state, and messages', () => {
|
||||
const state = createState()
|
||||
expect(db!.getLifecycleTransitionReceipts('task', state.task.id)).not.toEqual([])
|
||||
|
||||
db!.resetAll()
|
||||
|
||||
@@ -69,7 +68,6 @@ describe('OrchestrationDb reset scopes', () => {
|
||||
const sqlite = (db as unknown as { db: { prepare: (sql: string) => { all: () => unknown[] } } })
|
||||
.db
|
||||
expect(sqlite.prepare('SELECT * FROM run_coordinator_handles').all()).toEqual([])
|
||||
expect(sqlite.prepare('SELECT * FROM lifecycle_transition_receipts').all()).toEqual([])
|
||||
// The ledger survives so a lost reset response cannot replay as a new mutation.
|
||||
expect(db!.getMutationReceipt('caller_1', 'request_1')).toBeDefined()
|
||||
expect(db!.getInbox()).toEqual([])
|
||||
@@ -84,7 +82,6 @@ describe('OrchestrationDb reset scopes', () => {
|
||||
|
||||
it('resetTasks preserves Runs and messages while clearing every worker attachment', () => {
|
||||
const state = createState()
|
||||
expect(db!.getLifecycleTransitionReceipts('task', state.task.id)).not.toEqual([])
|
||||
|
||||
db!.resetTasks()
|
||||
|
||||
@@ -94,9 +91,6 @@ describe('OrchestrationDb reset scopes', () => {
|
||||
expect(db!.getWorkerDispatch(state.started.dispatch.id)).toBeUndefined()
|
||||
expect(db!.getFederatedDispatch(state.started.dispatch.id)).toBeUndefined()
|
||||
expect(db!.getRemoteQuestion('question_1')).toBeUndefined()
|
||||
const sqlite = (db as unknown as { db: { prepare: (sql: string) => { all: () => unknown[] } } })
|
||||
.db
|
||||
expect(sqlite.prepare('SELECT * FROM lifecycle_transition_receipts').all()).toEqual([])
|
||||
expect(db!.getMessageById(state.localQuestion.message.id)).toBeDefined()
|
||||
expect(db!.getQuestion(state.localQuestion.message.id)).toMatchObject({
|
||||
status: 'closed',
|
||||
|
||||
@@ -402,7 +402,6 @@ describe('OrchestrationDb version-skew migration', () => {
|
||||
raw.close()
|
||||
|
||||
db = new OrchestrationDb(dbPath)
|
||||
expect(db.db.prepare('SELECT * FROM lifecycle_transition_receipts').all()).toEqual([])
|
||||
expect(db.db.prepare('SELECT * FROM attempt_observation_facts').all()).toEqual([])
|
||||
})
|
||||
it('repairs a v33 schema missing the pointer-enter column', () => {
|
||||
|
||||
@@ -347,39 +347,6 @@ describe('OrchestrationDb worker Dispatch state', () => {
|
||||
expect(d.getTask(task.id)?.status).toBe('blocked')
|
||||
})
|
||||
|
||||
it('rolls back stop-unknown projection when its receipt cannot be inserted', () => {
|
||||
const d = createDb()
|
||||
const task = d.createTask({ spec: 'atomic uncertain stop' })
|
||||
const started = d.createStartingWorkerDispatch({
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: Number.MAX_SAFE_INTEGER,
|
||||
taskId: task.id,
|
||||
startOptions: {}
|
||||
})
|
||||
d.markWorkerDispatchReady(started.dispatch.id)
|
||||
expect(d.beginWorkerStop(started.dispatch.id, 'runtime_test').disposition).toBe('stopping')
|
||||
d.db.exec(`
|
||||
CREATE TRIGGER reject_worker_stop_unknown_receipt
|
||||
BEFORE INSERT ON lifecycle_transition_receipts
|
||||
WHEN NEW.kind = 'worker_stop_unknown'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced stop-unknown receipt failure'); END;
|
||||
`)
|
||||
|
||||
expect(() => d.markWorkerStopUnknown(started.dispatch.id, 'stop response lost')).toThrow(
|
||||
'forced stop-unknown receipt failure'
|
||||
)
|
||||
expect(d.getWorkerDispatch(started.dispatch.id)).toMatchObject({
|
||||
state: 'stopping',
|
||||
stage: 'stop_requested',
|
||||
last_error: null
|
||||
})
|
||||
expect(
|
||||
d
|
||||
.getLifecycleTransitionReceipts('worker', started.dispatch.id)
|
||||
.some((receipt) => receipt.kind === 'worker_stop_unknown')
|
||||
).toBe(false)
|
||||
})
|
||||
|
||||
it('allows explicit stop recovery from uncertain local and remote starts', () => {
|
||||
const d = createDb()
|
||||
const task = d.createTask({ spec: 'uncertain local start' })
|
||||
|
||||
@@ -59,33 +59,6 @@ describe('worker start settled by an unobserved prompt', () => {
|
||||
expect(db.getTask(taskId)).toMatchObject({ status: 'completed', result: 'done the work' })
|
||||
expect(db.getDispatchContextById(dispatchId)?.status).toBe('completed')
|
||||
expect(db.getWorkerDispatch(dispatchId)).toMatchObject({ state: 'succeeded', stage: 'settled' })
|
||||
expect(db.getLifecycleTransitionReceipts('task', taskId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'failed',
|
||||
to_state: 'completed',
|
||||
kind: 'task_prompt_stall_corrected'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(db.getLifecycleTransitionReceipts('dispatch', dispatchId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'failed',
|
||||
to_state: 'completed',
|
||||
kind: 'dispatch_prompt_stall_corrected'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(db.getLifecycleTransitionReceipts('worker', dispatchId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'failed',
|
||||
to_state: 'succeeded',
|
||||
kind: 'worker_prompt_stall_corrected'
|
||||
})
|
||||
])
|
||||
)
|
||||
})
|
||||
|
||||
it('revokes and stays settled when the start failed for any other cause', () => {
|
||||
@@ -119,33 +92,6 @@ describe('worker start settled by an unobserved prompt', () => {
|
||||
last_failure: 'build broke on X'
|
||||
})
|
||||
expect(db.getWorkerDispatch(dispatchId)).toMatchObject({ state: 'failed', stage: 'settled' })
|
||||
expect(db.getLifecycleTransitionReceipts('task', taskId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'failed',
|
||||
to_state: 'failed',
|
||||
kind: 'task_prompt_stall_corrected'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(db.getLifecycleTransitionReceipts('dispatch', dispatchId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'failed',
|
||||
to_state: 'failed',
|
||||
kind: 'dispatch_prompt_stall_corrected'
|
||||
})
|
||||
])
|
||||
)
|
||||
expect(db.getLifecycleTransitionReceipts('worker', dispatchId)).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
from_state: 'failed',
|
||||
to_state: 'failed',
|
||||
kind: 'worker_prompt_stall_corrected'
|
||||
})
|
||||
])
|
||||
)
|
||||
|
||||
// The stalled cause is gone, so a repeat report has nothing left to correct.
|
||||
expect(
|
||||
@@ -154,17 +100,18 @@ describe('worker start settled by an unobserved prompt', () => {
|
||||
expect(db.getTask(taskId)?.result).toBe('build broke on X')
|
||||
})
|
||||
|
||||
it('rolls back every prompt-stall correction when a receipt cannot be inserted', () => {
|
||||
it('rolls back every prompt-stall correction when the worker transition fails', () => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
const { taskId, dispatchId } = startWorker('atomic correction')
|
||||
db.failWorkerStart(dispatchId, 'dispatch_input', 'agent_prompt_stalled', {
|
||||
retainCapability: true
|
||||
})
|
||||
// The worker correction is the last of the three, so aborting it must undo the other two.
|
||||
db.db.exec(`
|
||||
CREATE TRIGGER reject_worker_prompt_stall_correction
|
||||
BEFORE INSERT ON lifecycle_transition_receipts
|
||||
WHEN NEW.kind = 'worker_prompt_stall_corrected'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced prompt-stall receipt failure'); END;
|
||||
BEFORE UPDATE ON worker_dispatches
|
||||
WHEN NEW.state = 'succeeded'
|
||||
BEGIN SELECT RAISE(ABORT, 'forced prompt-stall correction failure'); END;
|
||||
`)
|
||||
|
||||
expect(() =>
|
||||
@@ -174,7 +121,7 @@ describe('worker start settled by an unobserved prompt', () => {
|
||||
outcome: 'succeeded',
|
||||
result: 'uncommitted result'
|
||||
})
|
||||
).toThrow('forced prompt-stall receipt failure')
|
||||
).toThrow('forced prompt-stall correction failure')
|
||||
expect(db.getTask(taskId)).toMatchObject({ status: 'failed', result: null })
|
||||
expect(db.getDispatchContextById(dispatchId)).toMatchObject({
|
||||
status: 'failed',
|
||||
@@ -185,15 +132,5 @@ describe('worker start settled by an unobserved prompt', () => {
|
||||
state: 'failed',
|
||||
stage: 'dispatch_input'
|
||||
})
|
||||
expect(
|
||||
db
|
||||
.getLifecycleTransitionReceipts('task', taskId)
|
||||
.some((receipt) => receipt.kind === 'task_prompt_stall_corrected')
|
||||
).toBe(false)
|
||||
expect(
|
||||
db
|
||||
.getLifecycleTransitionReceipts('dispatch', dispatchId)
|
||||
.some((receipt) => receipt.kind === 'dispatch_prompt_stall_corrected')
|
||||
).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -166,8 +166,7 @@ function applyConfirmedFederatedReleaseHomeProjection(
|
||||
stage: 'released',
|
||||
agent_terminal_handle: null,
|
||||
updated_at: new Date().toISOString()
|
||||
},
|
||||
receipt: { kind: 'federated_release_confirmed' }
|
||||
}
|
||||
})
|
||||
}
|
||||
// The remote handle is an execution-host fact; clear it after confirmation
|
||||
|
||||
@@ -69,16 +69,7 @@ export function completeWorkerTerminalRelease(
|
||||
const release = completeWorkerTerminalReleaseOnce(args)
|
||||
.then((receipt) => {
|
||||
if (activeRelease.recoveryRequested) {
|
||||
args.db.recordWorkerTerminalRecoveryAttempt(
|
||||
args.resource.id,
|
||||
receipt.state === 'released' || receipt.state === 'already_released'
|
||||
? 'released'
|
||||
: receipt.state === 'release_pending'
|
||||
? 'pending'
|
||||
: receipt.state === 'release_unknown'
|
||||
? 'unknown'
|
||||
: 'retained'
|
||||
)
|
||||
args.db.recordWorkerTerminalRecoveryAttempt(args.resource.id)
|
||||
}
|
||||
return receipt
|
||||
})
|
||||
|
||||
@@ -308,25 +308,11 @@ describe('orchestration worker release recovery', () => {
|
||||
recovery_attempt_count: 1,
|
||||
last_recovery_at: expect.any(String)
|
||||
})
|
||||
expect(db.getLifecycleTransitionReceipts('worker', resourceId!)).toEqual([
|
||||
expect.objectContaining({ kind: 'worker_terminal_recovery' })
|
||||
])
|
||||
expect(
|
||||
db
|
||||
.getLifecycleTransitionReceipts('worker', dispatchId)
|
||||
.filter((receipt) => receipt.kind === 'worker_terminal_recovery')
|
||||
).toEqual([
|
||||
expect.objectContaining({
|
||||
kind: 'worker_terminal_recovery',
|
||||
entity_id: dispatchId
|
||||
})
|
||||
])
|
||||
|
||||
await expect(reconcileRequestedWorkerTerminalReleases(runtime)).resolves.toMatchObject({
|
||||
attempted: 0
|
||||
})
|
||||
expect(db.getWorkerTerminalResource(resourceId!)?.recovery_attempt_count).toBe(1)
|
||||
expect(db.getLifecycleTransitionReceipts('worker', resourceId!)).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('keeps live terminals bounded across 50 settled workers while controls survive', async () => {
|
||||
|
||||
@@ -58,7 +58,6 @@ describe('orchestration worker release', () => {
|
||||
recovery_attempt_count: 0,
|
||||
last_recovery_at: null
|
||||
})
|
||||
expect(h.db.getLifecycleTransitionReceipts('worker', resourceId!)).toEqual([])
|
||||
})
|
||||
|
||||
it('releases a failed worker the same way', async () => {
|
||||
|
||||
@@ -343,8 +343,7 @@ describe('orchestration worker recovery', () => {
|
||||
id: started.dispatch.id,
|
||||
from: 'ready',
|
||||
to: 'ready',
|
||||
projection: { stage: 'released', agent_terminal_handle: null },
|
||||
receipt: { kind: 'test_release' }
|
||||
projection: { stage: 'released', agent_terminal_handle: null }
|
||||
})
|
||||
db.db
|
||||
.prepare(
|
||||
|
||||
Reference in New Issue
Block a user