diff --git a/src/main/runtime/orchestration/context-only-dispatch-release.ts b/src/main/runtime/orchestration/context-only-dispatch-release.ts index ac2bd481c7a..ea99d78bf7d 100644 --- a/src/main/runtime/orchestration/context-only-dispatch-release.ts +++ b/src/main/runtime/orchestration/context-only-dispatch-release.ts @@ -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 } } diff --git a/src/main/runtime/orchestration/db-task-dispatch-invariant.test.ts b/src/main/runtime/orchestration/db-task-dispatch-invariant.test.ts index fc8876c0741..64401202bc0 100644 --- a/src/main/runtime/orchestration/db-task-dispatch-invariant.test.ts +++ b/src/main/runtime/orchestration/db-task-dispatch-invariant.test.ts @@ -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 }) - ]) - ) } ) diff --git a/src/main/runtime/orchestration/db-task-dispatch-lifecycle-guards.test.ts b/src/main/runtime/orchestration/db-task-dispatch-lifecycle-guards.test.ts index bc0d767e958..bb6d3472d4e 100644 --- a/src/main/runtime/orchestration/db-task-dispatch-lifecycle-guards.test.ts +++ b/src/main/runtime/orchestration/db-task-dispatch-lifecycle-guards.test.ts @@ -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)( diff --git a/src/main/runtime/orchestration/db-task-dispatch-races.test.ts b/src/main/runtime/orchestration/db-task-dispatch-races.test.ts index f59f30f6d2b..67ad9e06b11 100644 --- a/src/main/runtime/orchestration/db-task-dispatch-races.test.ts +++ b/src/main/runtime/orchestration/db-task-dispatch-races.test.ts @@ -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', () => { diff --git a/src/main/runtime/orchestration/db/decision-gate-lifecycle.test.ts b/src/main/runtime/orchestration/db/decision-gate-lifecycle.test.ts index c4b92b10652..219cf6fe212 100644 --- a/src/main/runtime/orchestration/db/decision-gate-lifecycle.test.ts +++ b/src/main/runtime/orchestration/db/decision-gate-lifecycle.test.ts @@ -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) }) }) diff --git a/src/main/runtime/orchestration/db/decision-gates/decision-gate-store.ts b/src/main/runtime/orchestration/db/decision-gates/decision-gate-store.ts index 2a98df34a0e..532fa52b22a 100644 --- a/src/main/runtime/orchestration/db/decision-gates/decision-gate-store.ts +++ b/src/main/runtime/orchestration/db/decision-gates/decision-gate-store.ts @@ -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 diff --git a/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts b/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts index aa5023f0478..2e281ca00f9 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts @@ -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 diff --git a/src/main/runtime/orchestration/db/dispatch-context/dispatch-context-store.ts b/src/main/runtime/orchestration/db/dispatch-context/dispatch-context-store.ts index 916f54935a0..3728278a670 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/dispatch-context-store.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/dispatch-context-store.ts @@ -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 = ?') diff --git a/src/main/runtime/orchestration/db/dispatch-context/task-dispatch-reconciliation.ts b/src/main/runtime/orchestration/db/dispatch-context/task-dispatch-reconciliation.ts index c3a2743a627..dd250552a2a 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/task-dispatch-reconciliation.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/task-dispatch-reconciliation.ts @@ -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 }) } diff --git a/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts b/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts index 6d9790c8965..3584a2e59a6 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts @@ -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( diff --git a/src/main/runtime/orchestration/db/federation/federated-dispatch-observation-fence.test.ts b/src/main/runtime/orchestration/db/federation/federated-dispatch-observation-fence.test.ts index 51a4393970b..e32bf27a00c 100644 --- a/src/main/runtime/orchestration/db/federation/federated-dispatch-observation-fence.test.ts +++ b/src/main/runtime/orchestration/db/federation/federated-dispatch-observation-fence.test.ts @@ -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( diff --git a/src/main/runtime/orchestration/db/lifecycle-transition.test.ts b/src/main/runtime/orchestration/db/lifecycle-transition.test.ts index 735d2a4448b..7b8d71284fe 100644 --- a/src/main/runtime/orchestration/db/lifecycle-transition.test.ts +++ b/src/main/runtime/orchestration/db/lifecycle-transition.test.ts @@ -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 }) }) }) diff --git a/src/main/runtime/orchestration/db/lifecycle-transition.ts b/src/main/runtime/orchestration/db/lifecycle-transition.ts index 956d56bc098..7c81d14c829 100644 --- a/src/main/runtime/orchestration/db/lifecycle-transition.ts +++ b/src/main/runtime/orchestration/db/lifecycle-transition.ts @@ -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 - /** 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 }) } diff --git a/src/main/runtime/orchestration/db/reset/orchestration-reset.ts b/src/main/runtime/orchestration/db/reset/orchestration-reset.ts index 355ea441053..f5532a2a1ae 100644 --- a/src/main/runtime/orchestration/db/reset/orchestration-reset.ts +++ b/src/main/runtime/orchestration/db/reset/orchestration-reset.ts @@ -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; diff --git a/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts b/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts index 183fcca68c0..92d56063763 100644 --- a/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts +++ b/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts @@ -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, diff --git a/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts b/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts index a77d8e9ecf8..f9190a350fd 100644 --- a/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts +++ b/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts @@ -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, diff --git a/src/main/runtime/orchestration/db/schema/migrate-mailbox-delivery-repair-v35.ts b/src/main/runtime/orchestration/db/schema/migrate-mailbox-delivery-repair-v35.ts index c04aa9ee5be..13ee89f35b0 100644 --- a/src/main/runtime/orchestration/db/schema/migrate-mailbox-delivery-repair-v35.ts +++ b/src/main/runtime/orchestration/db/schema/migrate-mailbox-delivery-repair-v35.ts @@ -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, diff --git a/src/main/runtime/orchestration/db/tasks/task-status-transition.ts b/src/main/runtime/orchestration/db/tasks/task-status-transition.ts index 9e77d1cf045..e1de8dc5b17 100644 --- a/src/main/runtime/orchestration/db/tasks/task-status-transition.ts +++ b/src/main/runtime/orchestration/db/tasks/task-status-transition.ts @@ -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') { diff --git a/src/main/runtime/orchestration/db/tasks/task-store.ts b/src/main/runtime/orchestration/db/tasks/task-store.ts index 5f5eb7c09d4..1b731f71628 100644 --- a/src/main/runtime/orchestration/db/tasks/task-store.ts +++ b/src/main/runtime/orchestration/db/tasks/task-store.ts @@ -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' }) } } diff --git a/src/main/runtime/orchestration/db/worker-dispatch/federated-worker-start-reconcile.ts b/src/main/runtime/orchestration/db/worker-dispatch/federated-worker-start-reconcile.ts index ea285d3ace4..9ce80695d5b 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/federated-worker-start-reconcile.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/federated-worker-start-reconcile.ts @@ -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) diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-abandon.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-abandon.ts index ced51a6270c..7549cb4eb26 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-abandon.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-abandon.ts @@ -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) diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts index 4b0932196a7..7da886c3d29 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts @@ -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') diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stage.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stage.ts index d9d52f0602f..0bee36d8f05 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stage.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stage.ts @@ -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) { diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-start.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-start.ts index 409606892af..bcb728f70cf 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-start.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-start.ts @@ -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 diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stop.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stop.ts index d254768db3a..8dbda2030c5 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stop.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-stop.ts @@ -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 diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts index 05cf8398bc3..97179eb0185 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts @@ -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') diff --git a/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-resource-store.ts b/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-resource-store.ts index 1d8bc9a4f6f..dc211f01a08 100644 --- a/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-resource-store.ts +++ b/src/main/runtime/orchestration/db/worker-terminal/worker-terminal-resource-store.ts @@ -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 diff --git a/src/main/runtime/orchestration/lifecycle-reconciliation.test.ts b/src/main/runtime/orchestration/lifecycle-reconciliation.test.ts index 9446049d677..3f0f0a7664f 100644 --- a/src/main/runtime/orchestration/lifecycle-reconciliation.test.ts +++ b/src/main/runtime/orchestration/lifecycle-reconciliation.test.ts @@ -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', () => { diff --git a/src/main/runtime/orchestration/orchestration-reset-db.test.ts b/src/main/runtime/orchestration/orchestration-reset-db.test.ts index 2df16d2b321..4b3e033a019 100644 --- a/src/main/runtime/orchestration/orchestration-reset-db.test.ts +++ b/src/main/runtime/orchestration/orchestration-reset-db.test.ts @@ -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', diff --git a/src/main/runtime/orchestration/orchestration-version-skew-migration.test.ts b/src/main/runtime/orchestration/orchestration-version-skew-migration.test.ts index 69732e391a5..923d45d701a 100644 --- a/src/main/runtime/orchestration/orchestration-version-skew-migration.test.ts +++ b/src/main/runtime/orchestration/orchestration-version-skew-migration.test.ts @@ -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', () => { diff --git a/src/main/runtime/orchestration/orchestration-worker-dispatch-db.test.ts b/src/main/runtime/orchestration/orchestration-worker-dispatch-db.test.ts index f34173054e7..dee969d70ea 100644 --- a/src/main/runtime/orchestration/orchestration-worker-dispatch-db.test.ts +++ b/src/main/runtime/orchestration/orchestration-worker-dispatch-db.test.ts @@ -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' }) diff --git a/src/main/runtime/orchestration/worker-start-unobserved-prompt-settlement.test.ts b/src/main/runtime/orchestration/worker-start-unobserved-prompt-settlement.test.ts index 0dd5660f6dc..382ec304bb6 100644 --- a/src/main/runtime/orchestration/worker-start-unobserved-prompt-settlement.test.ts +++ b/src/main/runtime/orchestration/worker-start-unobserved-prompt-settlement.test.ts @@ -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) }) }) diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-release.ts b/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-release.ts index 7adedea3559..26abad2cbd4 100644 --- a/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-release.ts +++ b/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-release.ts @@ -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 diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-release-completion.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-release-completion.ts index 9f1718903f2..5fa7ee4b1cd 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-release-completion.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-release-completion.ts @@ -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 }) diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-release-recovery.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-release-recovery.test.ts index 74c4e71d337..b9e80cf16e9 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-release-recovery.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-release-recovery.test.ts @@ -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 () => { diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-release.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-release.test.ts index 3982baceb34..f5fd76b682b 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-release.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-release.test.ts @@ -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 () => { diff --git a/src/main/runtime/rpc/methods/orchestration/worker/workers-recovery.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/workers-recovery.test.ts index 6fb7f7e8d88..a635a316b23 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/workers-recovery.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/workers-recovery.test.ts @@ -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(