fix(orchestration): validate live consumers and simplify batch revocation

This commit is contained in:
Jinwoo-H
2026-09-10 18:12:50 -04:00
parent d008782bce
commit 0e70dbcd40
10 changed files with 269 additions and 20 deletions
@@ -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
@@ -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')
@@ -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'
@@ -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])
}
)
})
@@ -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
@@ -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
})
}
@@ -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)
}
@@ -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
})
}
@@ -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' })
@@ -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