diff --git a/docs/reference/orchestration-delivery-storage.md b/docs/reference/orchestration-delivery-storage.md index ecb0ee2c402..bf4d4f927e4 100644 --- a/docs/reference/orchestration-delivery-storage.md +++ b/docs/reference/orchestration-delivery-storage.md @@ -1,6 +1,6 @@ # Orchestration delivery storage -A message records whether it still needs attention (`read`). A delivery records a stable batch ID, ordered message IDs, mailbox, consumer generation, and creation time. Acknowledgement time and fencing record actual events. There is no stored delivery status. +A message records whether it still needs attention (`read`). A delivery records a stable batch ID, ordered message IDs, mailbox, consumer generation, and creation time. Acknowledgement time records receipt. Fencing permanently revokes an unacknowledged batch when its consumer is replaced or legacy ownership is adopted, independently of message reads. Historical adoption can revoke a batch without changing its generation, so generation validation alone cannot replace this fact. There is no stored delivery status. `outstanding_deliveries` is a SQLite view, not a table or cached projection. It selects batches that have not been acknowledged or fenced and still contain unread messages. Both consuming checks and notification eligibility query this view. @@ -8,7 +8,7 @@ When completion makes an earlier heartbeat obsolete, only the message changes. I Acknowledgement marks the batch's messages read and records `acknowledged_at` in one transaction. A first explicit acknowledgement after lifecycle suppression records the actual event; subsequent acknowledgements are idempotent. Suppression never invents an acknowledgement timestamp. An acknowledgement cannot consume newer messages outside the batch. -Creation and acknowledgement use `BEGIN IMMEDIATE` and validate the current consumer inside that transaction. Run coordinators use the Run generation. Workers use the generation on their local Dispatch or federated attachment, according to the caller's existing routing. Loopback federation can contain both records, so the source is explicit. The insertion constraint consults the same view to prevent two outstanding batches in a mailbox, while allowing historical batches. +Creation and acknowledgement use `BEGIN IMMEDIATE` and validate the current consumer inside that transaction. Worker validation also checks lifecycle state under the same lock, so a settled worker cannot consume mail awaiting rerouting. Run coordinators use the Run generation. Workers use the generation on their local Dispatch or federated attachment, according to the caller's existing routing. Loopback federation can contain both records, so the source is explicit. The insertion constraint consults the same view to prevent two outstanding batches in a mailbox, while allowing historical batches. ## Migration and compatibility diff --git a/src/main/runtime/orchestration/db/dispatch-context/dispatch-capability.ts b/src/main/runtime/orchestration/db/dispatch-context/dispatch-capability.ts index 4f6861a17d0..007e81d1016 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/dispatch-capability.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/dispatch-capability.ts @@ -38,7 +38,7 @@ export function mintDispatchCapability( params.processIncarnation, params.dispatchId ) - this.fenceOutstandingMailboxDelivery(`dispatch:${params.dispatchId}`) + this.fenceUnacknowledgedMailboxDeliveries(`dispatch:${params.dispatchId}`) this.db.exec('COMMIT') } catch (error) { this.db.exec('ROLLBACK') diff --git a/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-authority.ts b/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-authority.ts index 5b2dc60615c..ac85fc7c3c9 100644 --- a/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-authority.ts +++ b/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-authority.ts @@ -71,7 +71,7 @@ export function prepareRemoteAttachmentAuthority( `Remote Dispatch ${params.dispatchId} is not starting.` ) } - this.fenceOutstandingMailboxDelivery(`dispatch:${params.dispatchId}`) + this.fenceUnacknowledgedMailboxDeliveries(`dispatch:${params.dispatchId}`) if (params.terminalOwnership && !this.getWorkerTerminalResourceByOwner(params.dispatchId)) { const resource = params.terminalOwnership === 'external' diff --git a/src/main/runtime/orchestration/db/messages/mailbox-consumer-lifecycle-fencing.test.ts b/src/main/runtime/orchestration/db/messages/mailbox-consumer-lifecycle-fencing.test.ts new file mode 100644 index 00000000000..3f81901b208 --- /dev/null +++ b/src/main/runtime/orchestration/db/messages/mailbox-consumer-lifecycle-fencing.test.ts @@ -0,0 +1,234 @@ +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it } from 'vitest' +import { ORCHESTRATION_CONTRACT_VERSION } from '../../../../../shared/protocol-version' +import { OrchestrationDb } from '../orchestration-db' +import { createRootDispatch } from '../root-dispatch-test-fixture' + +type Settlement = 'local completion' | 'local failure' | 'remote stop' | 'remote failure' +type DeliveryOperation = 'create' | 'acknowledge' + +describe('mailbox consumer lifecycle fencing', () => { + const connections: OrchestrationDb[] = [] + const directories: string[] = [] + + afterEach(() => { + for (const db of connections.splice(0)) { + db.close() + } + for (const directory of directories.splice(0)) { + rmSync(directory, { recursive: true, force: true }) + } + }) + + function open(path: string): OrchestrationDb { + const db = new OrchestrationDb(path) + connections.push(db) + return db + } + + function databasePath(): string { + const directory = mkdtempSync(join(tmpdir(), 'orca-mailbox-consumer-lifecycle-')) + directories.push(directory) + return join(directory, 'orchestration.db') + } + + function setup(settlement: Settlement): { + db: OrchestrationDb + peer: OrchestrationDb + messageId: string + params: { + runId: string + mailboxHandle: string + consumerGeneration: number + consumerSource: 'dispatch' | 'attachment' + } + settle: () => void + } { + const path = databasePath() + const db = open(path) + const run = db.createRun({ + objective: 'Fence settled mailbox consumers', + coordinatorHandle: 'coord', + coordinatorPaneKey: 'tab:11111111-1111-4111-8111-111111111111' + }) + const remote = settlement.startsWith('remote') + const dispatchId = remote + ? `ctx_${settlement.replace(' ', '_')}` + : createRootDispatch(db, db.createTask({ runId: run.id, spec: settlement }).id, 'worker').id + let consumerGeneration = 0 + + if (remote) { + db.createRemoteDispatchAttachment({ + runId: run.id, + dispatchId, + taskId: `task_${dispatchId}`, + homePeerFingerprint: 'home-peer', + protocolVersion: ORCHESTRATION_CONTRACT_VERSION, + runtimeEpoch: 'epoch-1', + mutationReceipt: { + callerFingerprint: 'home-peer', + requestId: `request_${dispatchId}`, + method: 'orchestration.federationAttachStart', + payloadHash: `hash_${dispatchId}` + } + }) + if (settlement === 'remote stop') { + db.prepareRemoteAttachmentAuthority({ + dispatchId, + paneKey: 'worker:22222222-2222-4222-9222-222222222222', + processIncarnation: 'runtime:worker:1', + worktreeId: 'folder', + terminalHandle: 'worker', + setupState: 'not_applicable', + effects: [] + }) + db.markRemoteAttachmentReady(dispatchId) + consumerGeneration = 1 + } + } + + const mailboxHandle = `dispatch:${dispatchId}` + const message = db.insertMessage({ + runId: run.id, + from: 'coord', + to: mailboxHandle, + subject: 'must remain unread' + }) + const peer = open(path) + const settle = (): void => { + if (settlement === 'local completion') { + peer.completeDispatch(dispatchId) + } else if (settlement === 'local failure') { + peer.failDispatch(dispatchId, 'settled by peer') + } else if (settlement === 'remote stop') { + peer.beginRemoteAttachmentStop(dispatchId) + peer.settleRemoteAttachmentStop(dispatchId) + } else { + peer.failRemoteAttachment(dispatchId, 'peer_failure', 'settled by peer', false) + } + } + + return { + db, + peer, + messageId: message.id, + params: { + runId: run.id, + mailboxHandle, + consumerGeneration, + consumerSource: remote ? 'attachment' : 'dispatch' + }, + settle + } + } + + function currentGeneration( + db: OrchestrationDb, + params: { + mailboxHandle: string + consumerSource: 'dispatch' | 'attachment' + } + ): number | undefined { + const dispatchId = params.mailboxHandle.slice('dispatch:'.length) + return params.consumerSource === 'dispatch' + ? db.getDispatchContextById(dispatchId)?.consumer_generation + : db.getRemoteDispatchAttachment(dispatchId)?.consumer_generation + } + + it.each<{ + operation: DeliveryOperation + settlement: Settlement + }>([ + { operation: 'create', settlement: 'local completion' }, + { operation: 'create', settlement: 'local failure' }, + { operation: 'create', settlement: 'remote stop' }, + { operation: 'create', settlement: 'remote failure' }, + { operation: 'acknowledge', settlement: 'local completion' }, + { operation: 'acknowledge', settlement: 'local failure' }, + { operation: 'acknowledge', settlement: 'remote stop' }, + { operation: 'acknowledge', settlement: 'remote failure' } + ])('rejects $operation after $settlement on another connection', ({ operation, settlement }) => { + const { db, peer, messageId, params, settle } = setup(settlement) + const delivery = + operation === 'acknowledge' ? db.getOrCreateMailboxDelivery(params)?.delivery : undefined + + settle() + + const operationCall = (): unknown => + operation === 'create' + ? db.getOrCreateMailboxDelivery(params) + : db.acknowledgeMailboxDelivery({ ...params, deliveryId: delivery!.id }) + expect(operationCall).toThrow(expect.objectContaining({ code: 'consumer_fenced' })) + expect(db.getMessageById(messageId)?.read).toBe(0) + expect(currentGeneration(peer, params)).toBe(params.consumerGeneration) + if (delivery) { + expect(db.getDeliveryRaw(delivery.id)?.acknowledged_at).toBeNull() + } + }) + + it.each(['start_unknown', 'stop_unknown'] as const)( + 'keeps a remote %s attachment eligible to consume mail', + (state) => { + const path = databasePath() + const db = open(path) + const run = db.createRun({ + objective: 'Preserve unverifiable remote consumers', + coordinatorHandle: 'coord', + coordinatorPaneKey: 'tab:11111111-1111-4111-8111-111111111111' + }) + const dispatchId = `ctx_${state}` + db.createRemoteDispatchAttachment({ + runId: run.id, + dispatchId, + taskId: `task_${state}`, + homePeerFingerprint: 'home-peer', + protocolVersion: ORCHESTRATION_CONTRACT_VERSION, + runtimeEpoch: 'epoch-1', + mutationReceipt: { + callerFingerprint: 'home-peer', + requestId: `request_${state}`, + method: 'orchestration.federationAttachStart', + payloadHash: `hash_${state}` + } + }) + let consumerGeneration = 0 + if (state === 'start_unknown') { + db.failRemoteAttachment(dispatchId, 'start_unknown', 'contact lost', true) + } else { + db.prepareRemoteAttachmentAuthority({ + dispatchId, + paneKey: 'worker:22222222-2222-4222-9222-222222222222', + processIncarnation: 'runtime:worker:1', + worktreeId: 'folder', + terminalHandle: 'worker', + setupState: 'not_applicable', + effects: [] + }) + db.markRemoteAttachmentReady(dispatchId) + db.beginRemoteAttachmentStop(dispatchId) + db.markRemoteAttachmentStopUnknown(dispatchId, 'contact lost') + consumerGeneration = 1 + } + const mailboxHandle = `dispatch:${dispatchId}` + const message = db.insertMessage({ + runId: run.id, + from: 'coord', + to: mailboxHandle, + subject: 'still deliverable' + }) + + expect( + db + .getOrCreateMailboxDelivery({ + runId: run.id, + mailboxHandle, + consumerGeneration, + consumerSource: 'attachment' + }) + ?.messages.map((row) => row.id) + ).toEqual([message.id]) + } + ) +}) diff --git a/src/main/runtime/orchestration/db/messages/mailbox-consumer.ts b/src/main/runtime/orchestration/db/messages/mailbox-consumer.ts index 9ed2b2d4cc9..66d72fb63d4 100644 --- a/src/main/runtime/orchestration/db/messages/mailbox-consumer.ts +++ b/src/main/runtime/orchestration/db/messages/mailbox-consumer.ts @@ -1,5 +1,15 @@ import type { OrchestrationDb } from '../orchestration-db' import { OrchestrationError } from '../../orchestration-error' +import { potentiallyLiveRemoteAttachmentSql } from '../federation/remote-attachment-liveness' + +const ACTIVE_DISPATCH_CONSUMER_SQL = ` + SELECT run_id, consumer_generation FROM dispatch_contexts + WHERE id = ? AND status IN ('pending', 'dispatched') +` +const ACTIVE_ATTACHMENT_CONSUMER_SQL = ` + SELECT home_run_id AS run_id, consumer_generation FROM remote_dispatch_attachments + WHERE dispatch_id = ? AND ${potentiallyLiveRemoteAttachmentSql()} +` // Validate inside the delivery transaction, so another connection cannot replace the consumer mid-check. export function requireMailboxConsumer( @@ -21,8 +31,8 @@ export function requireMailboxConsumer( // A loopback runtime has both records; use the counter belonging to the caller's attachment. const sql = params.consumerSource === 'attachment' - ? 'SELECT home_run_id AS run_id, consumer_generation FROM remote_dispatch_attachments WHERE dispatch_id = ?' - : 'SELECT run_id, consumer_generation FROM dispatch_contexts WHERE id = ?' + ? ACTIVE_ATTACHMENT_CONSUMER_SQL + : ACTIVE_DISPATCH_CONSUMER_SQL const consumer = db.db.prepare(sql).get(dispatchId) as | { run_id: string; consumer_generation: number } | undefined diff --git a/src/main/runtime/orchestration/db/messages/role-mailbox-delivery.ts b/src/main/runtime/orchestration/db/messages/role-mailbox-delivery.ts index d5ff9e04221..8c50b8ccf2f 100644 --- a/src/main/runtime/orchestration/db/messages/role-mailbox-delivery.ts +++ b/src/main/runtime/orchestration/db/messages/role-mailbox-delivery.ts @@ -177,14 +177,14 @@ export function hasOutstandingMailboxDelivery( ) } -export function fenceOutstandingMailboxDelivery( +export function fenceUnacknowledgedMailboxDeliveries( this: OrchestrationDb, mailboxHandle: string ): void { this.db .prepare( `UPDATE deliveries SET fenced = 1 - WHERE id IN (SELECT id FROM outstanding_deliveries WHERE mailbox_handle = ?)` + WHERE mailbox_handle = ? AND acknowledged_at IS NULL AND fenced = 0` ) .run(mailboxHandle) } @@ -195,7 +195,7 @@ export type RoleMailboxDeliveryMethods = { getOrCreateMailboxDelivery: typeof getOrCreateMailboxDelivery acknowledgeMailboxDelivery: typeof acknowledgeMailboxDelivery hasOutstandingMailboxDelivery: typeof hasOutstandingMailboxDelivery - fenceOutstandingMailboxDelivery: typeof fenceOutstandingMailboxDelivery + fenceUnacknowledgedMailboxDeliveries: typeof fenceUnacknowledgedMailboxDeliveries } export function attachRoleMailboxDelivery(ctor: { prototype: object }): void { @@ -205,6 +205,6 @@ export function attachRoleMailboxDelivery(ctor: { prototype: object }): void { getOrCreateMailboxDelivery, acknowledgeMailboxDelivery, hasOutstandingMailboxDelivery, - fenceOutstandingMailboxDelivery + fenceUnacknowledgedMailboxDeliveries }) } diff --git a/src/main/runtime/orchestration/db/runs/run-binding.ts b/src/main/runtime/orchestration/db/runs/run-binding.ts index a4110dd589f..e2ff77ec818 100644 --- a/src/main/runtime/orchestration/db/runs/run-binding.ts +++ b/src/main/runtime/orchestration/db/runs/run-binding.ts @@ -142,7 +142,7 @@ export function bindRun( WHERE id = ?` ) .run(params.coordinatorHandle, params.coordinatorPaneKey, params.runId) - this.fenceOutstandingDelivery(params.runId) + this.fenceUnacknowledgedMailboxDeliveries(`run:${params.runId}`) if (params.takeoverLegacy || replacesLegacyCoordinator) { this.promoteLegacyCoordinatorMailForTakeover(params.runId, retainedCoordinatorHandle) } diff --git a/src/main/runtime/orchestration/db/runs/run-lookup.ts b/src/main/runtime/orchestration/db/runs/run-lookup.ts index 7563effa8c3..194cceffcf0 100644 --- a/src/main/runtime/orchestration/db/runs/run-lookup.ts +++ b/src/main/runtime/orchestration/db/runs/run-lookup.ts @@ -141,7 +141,7 @@ export function unbindOtherRunsForPane( WHERE id = ?` ) .run(run.id) - this.fenceOutstandingDelivery(run.id) + this.fenceUnacknowledgedMailboxDeliveries(`run:${run.id}`) } } } @@ -152,10 +152,6 @@ export function requireRun(this: OrchestrationDb, runId: string): void { } } -export function fenceOutstandingDelivery(this: OrchestrationDb, runId: string): void { - this.fenceOutstandingMailboxDelivery(`run:${runId}`) -} - export type RunLookupMethods = { getRun: typeof getRun getLegacyAdoptedRunMailboxOwner: typeof getLegacyAdoptedRunMailboxOwner @@ -166,7 +162,6 @@ export type RunLookupMethods = { getRunRaw: typeof getRunRaw unbindOtherRunsForPane: typeof unbindOtherRunsForPane requireRun: typeof requireRun - fenceOutstandingDelivery: typeof fenceOutstandingDelivery } export function attachRunLookup(ctor: { prototype: object }): void { @@ -179,7 +174,6 @@ export function attachRunLookup(ctor: { prototype: object }): void { runsBoundToPane, getRunRaw, unbindOtherRunsForPane, - requireRun, - fenceOutstandingDelivery + requireRun }) } diff --git a/src/main/runtime/orchestration/db/schema/derived-delivery-migration.test.ts b/src/main/runtime/orchestration/db/schema/derived-delivery-migration.test.ts index fb83990d822..633541f8d35 100644 --- a/src/main/runtime/orchestration/db/schema/derived-delivery-migration.test.ts +++ b/src/main/runtime/orchestration/db/schema/derived-delivery-migration.test.ts @@ -72,6 +72,9 @@ describe('derived delivery migration', () => { fenced: 0 }) expect(db.getDeliveryRaw('history_fence')).toMatchObject({ acknowledged_at: null, fenced: 1 }) + expect(() => db.acknowledgeRunDelivery({ ...params, deliveryId: 'history_fence' })).toThrow( + expect.objectContaining({ code: 'consumer_fenced' }) + ) expect(db.getOrCreateRunDelivery(params)?.messages.map((message) => message.id)).toEqual([ next.id ]) @@ -139,6 +142,10 @@ describe('derived delivery migration', () => { coordinatorHandle: 'replacement', coordinatorPaneKey: 'other:22222222-2222-4222-9222-222222222222' })! + expect(first.getDeliveryRaw(batch.delivery.id)).toMatchObject({ + fenced: 1, + acknowledged_at: null + }) first.insertMessage({ runId: run.id, from: 'worker', to: `run:${run.id}`, subject: 'next' }) expect(() => first.getOrCreateRunDelivery(params)).toThrow( expect.objectContaining({ code: 'consumer_fenced' }) @@ -211,6 +218,10 @@ describe('derived delivery migration', () => { effects: [] }) } + expect(db.getDeliveryRaw(batch.delivery.id)).toMatchObject({ + fenced: 1, + acknowledged_at: null + }) peer.insertMessage({ runId: run.id, from: 'coord', to: mailboxHandle, subject: 'next' }) expect(() => db.getOrCreateMailboxDelivery(params)).toThrow( expect.objectContaining({ code: 'consumer_fenced' }) diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts index 57ec354a8cc..b33a8798dc3 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts @@ -75,7 +75,7 @@ export function prepareStartingWorkerAuthority( `Dispatch ${params.dispatchId} is not starting.` ) } - this.fenceOutstandingMailboxDelivery(`dispatch:${params.dispatchId}`) + this.fenceUnacknowledgedMailboxDeliveries(`dispatch:${params.dispatchId}`) const workerUpdate = this.db .prepare( `UPDATE worker_dispatches