diff --git a/src/cli/specs/orchestration.ts b/src/cli/specs/orchestration.ts index e62911b2c76..3f63ff3a47d 100644 --- a/src/cli/specs/orchestration.ts +++ b/src/cli/specs/orchestration.ts @@ -110,8 +110,10 @@ export const ORCHESTRATION_COMMAND_SPECS: CommandSpec[] = [ notes: [ 'On Windows PowerShell, quote comma-separated type filters, e.g. --types "worker_done,escalation".', '--types is the wake condition for --wait; a returned Delivery is always the whole FIFO batch, so it is never filtered by type. Only --peek and --all filter their rows.', + 'Without --wait, --types does not filter a consuming check. An outstanding Delivery replays before the wake condition is tested.', + '--all shows the outstanding deliveryId for recovery; its history rows are not necessarily that batch. Check the full Delivery before --ack.', '--format renders the returned rows as local text only; it never writes to another terminal.', - 'A bound Run replays the same Delivery until --ack; process every message before acknowledging.' + 'A bound Run replays the same Delivery until --ack or all its messages are retired; process every message before acknowledging.' ] }, { diff --git a/src/main/runtime/orchestration/db/messages/message-inbox.ts b/src/main/runtime/orchestration/db/messages/message-inbox.ts index e9b800d431e..1ef57776b00 100644 --- a/src/main/runtime/orchestration/db/messages/message-inbox.ts +++ b/src/main/runtime/orchestration/db/messages/message-inbox.ts @@ -9,7 +9,8 @@ const MESSAGE_MUTATION_SAVEPOINT = 'message_id_mutation' function runBatchedMessageMutation( db: OrchestrationDb, ids: string[], - sqlForPlaceholders: (placeholders: string) => string + sqlForPlaceholders: (placeholders: string) => string, + retireReadDeliveries = false ): void { if (ids.length === 0) { return @@ -20,6 +21,14 @@ function runBatchedMessageMutation( const batch = ids.slice(offset, offset + MESSAGE_ID_UPDATE_BATCH_SIZE) const placeholders = batch.map(() => '?').join(',') db.db.prepare(sqlForPlaceholders(placeholders)).run(...batch) + if (retireReadDeliveries) { + const mailboxes = db.db + .prepare(`SELECT DISTINCT to_handle FROM messages WHERE id IN (${placeholders})`) + .all(...batch) as { to_handle: string }[] + for (const mailbox of mailboxes) { + db.retireReadMailboxDelivery(mailbox.to_handle) + } + } } db.db.exec(`RELEASE ${MESSAGE_MUTATION_SAVEPOINT}`) } catch (error) { @@ -160,7 +169,8 @@ export function markAsRead(this: OrchestrationDb, ids: string[]): void { `UPDATE messages SET read = 1, pointer_enter_pending = 0, pointer_pty_id = NULL, pointer_process_incarnation = NULL - WHERE id IN (${placeholders})` + WHERE id IN (${placeholders})`, + true ) } @@ -216,7 +226,8 @@ export function markAsReadAndDelivered(this: OrchestrationDb, ids: string[]): vo SET read = 1, delivered_at = COALESCE(delivered_at, datetime('now')), pointer_enter_pending = 0, pointer_pty_id = NULL, pointer_process_incarnation = NULL - WHERE id IN (${placeholders})` + WHERE id IN (${placeholders})`, + true ) } 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 c7553c089b7..8a5968c36f2 100644 --- a/src/main/runtime/orchestration/db/messages/role-mailbox-delivery.ts +++ b/src/main/runtime/orchestration/db/messages/role-mailbox-delivery.ts @@ -1,6 +1,6 @@ import type { DeliveryRow, MessageRow, MessageType } from '../../types' import { OrchestrationError } from '../../orchestration-error' -import { generateId } from '../generated-id' +import { generateId, isGeneratedId } from '../generated-id' import type { OrchestrationDb } from '../orchestration-db' import { exposeDeliveryTimestamps, exposeMessageListTimestamps } from '../utc-timestamp' import { ORCHESTRATION_DELIVERY_BATCH_LIMIT } from './mailbox-routing-page' @@ -43,9 +43,7 @@ export function getOrCreateMailboxDelivery( if (params.requireCurrentRunConsumer) { this.requireCurrentConsumer(params.runId, params.consumerGeneration) } - const existing = this.db - .prepare("SELECT * FROM deliveries WHERE mailbox_handle = ? AND status = 'outstanding'") - .get(params.mailboxHandle) as DeliveryRow | undefined + const existing = this.getOutstandingMailboxDelivery(params.mailboxHandle) if (existing) { if (existing.consumer_generation !== params.consumerGeneration) { throw new OrchestrationError( @@ -54,8 +52,11 @@ export function getOrCreateMailboxDelivery( ) } const messages = this.getDeliveryMessages(existing) - this.db.exec('COMMIT') - return { delivery: exposeDeliveryTimestamps(existing), messages, replayed: true } + if (messages.some((message) => message.read === 0)) { + this.db.exec('COMMIT') + return { delivery: exposeDeliveryTimestamps(existing), messages, replayed: true } + } + this.retireReadMailboxDelivery(params.mailboxHandle) } if (params.wakeTypes?.length) { const placeholders = params.wakeTypes.map(() => '?').join(',') @@ -124,6 +125,20 @@ export function acknowledgeMailboxDelivery( if (params.requireCurrentRunConsumer) { this.requireCurrentConsumer(params.runId, params.consumerGeneration) } + if (isGeneratedId(params.deliveryId, 'msg')) { + const outstanding = this.getOutstandingMailboxDelivery(params.mailboxHandle) + const containsMessage = + outstanding && + outstanding.run_id === params.runId && + (JSON.parse(outstanding.message_ids) as string[]).includes(params.deliveryId) + const guidance = containsMessage + ? `Process the entire batch, then use --ack ${outstanding.id}.` + : 'Run orchestration check to obtain the Delivery id.' + throw new OrchestrationError( + 'stale_delivery', + `--ack takes a Delivery id, not message id ${params.deliveryId}. ${guidance}` + ) + } const delivery = this.getDeliveryRaw(params.deliveryId) if ( !delivery || @@ -174,17 +189,35 @@ export function acknowledgeMailboxDelivery( } } +export function getOutstandingMailboxDelivery( + this: OrchestrationDb, + mailboxHandle: string +): DeliveryRow | undefined { + const row = this.db + .prepare("SELECT * FROM deliveries WHERE mailbox_handle = ? AND status = 'outstanding'") + .get(mailboxHandle) as DeliveryRow | undefined + return row ? exposeDeliveryTimestamps(row) : undefined +} + +// Acknowledgment keeps late consumer acks idempotent after lifecycle suppression. +export function retireReadMailboxDelivery(this: OrchestrationDb, mailboxHandle: string): void { + this.db + .prepare( + `UPDATE deliveries SET status = 'acknowledged', acknowledged_at = datetime('now') + WHERE mailbox_handle = ? AND status = 'outstanding' + AND NOT EXISTS ( + SELECT 1 FROM json_each(deliveries.message_ids) AS member + JOIN messages ON messages.id = member.value WHERE messages.read = 0 + )` + ) + .run(mailboxHandle) +} + export function hasOutstandingMailboxDelivery( this: OrchestrationDb, mailboxHandle: string ): boolean { - return Boolean( - this.db - .prepare( - "SELECT 1 FROM deliveries WHERE mailbox_handle = ? AND status = 'outstanding' LIMIT 1" - ) - .get(mailboxHandle) - ) + return this.getOutstandingMailboxDelivery(mailboxHandle) !== undefined } export function fenceOutstandingMailboxDelivery( @@ -199,6 +232,8 @@ export function fenceOutstandingMailboxDelivery( } export type RoleMailboxDeliveryMethods = { + getOutstandingMailboxDelivery: typeof getOutstandingMailboxDelivery + retireReadMailboxDelivery: typeof retireReadMailboxDelivery getDeliveryRaw: typeof getDeliveryRaw getDeliveryMessages: typeof getDeliveryMessages getOrCreateMailboxDelivery: typeof getOrCreateMailboxDelivery @@ -209,6 +244,8 @@ export type RoleMailboxDeliveryMethods = { export function attachRoleMailboxDelivery(ctor: { prototype: object }): void { Object.assign(ctor.prototype, { + getOutstandingMailboxDelivery, + retireReadMailboxDelivery, getDeliveryRaw, getDeliveryMessages, getOrCreateMailboxDelivery, diff --git a/src/main/runtime/orchestration/formatter.test.ts b/src/main/runtime/orchestration/formatter.test.ts index 148bfad9b73..e42b80f3f46 100644 --- a/src/main/runtime/orchestration/formatter.test.ts +++ b/src/main/runtime/orchestration/formatter.test.ts @@ -187,7 +187,7 @@ describe('formatMessagesForInjection', () => { describe('formatMessagePointer', () => { it('formats a singular pointer without message content', () => { expect(formatMessagePointer(1, 'run:run_1')).toBe( - '\nYou have 1 orchestration message. Run `orca orchestration check --run run_1`.\n' + '\n[Orca orchestration notification] You have 1 orchestration message. Run `orca orchestration check --run run_1`.\n' ) }) diff --git a/src/main/runtime/orchestration/formatter.ts b/src/main/runtime/orchestration/formatter.ts index 2dd4f86777b..540e1def1ba 100644 --- a/src/main/runtime/orchestration/formatter.ts +++ b/src/main/runtime/orchestration/formatter.ts @@ -118,5 +118,5 @@ export function formatMessagePointer( const runFlag = mailboxHandle?.startsWith('run:') ? ` --run ${mailboxHandle.slice('run:'.length)}` : '' - return `\nYou have ${count} orchestration ${noun}. Run \`${cliCommand} orchestration check${runFlag}\`.\n` + return `\n[Orca orchestration notification] You have ${count} orchestration ${noun}. Run \`${cliCommand} orchestration check${runFlag}\`.\n` } diff --git a/src/main/runtime/orchestration/orchestration-retired-delivery.test.ts b/src/main/runtime/orchestration/orchestration-retired-delivery.test.ts new file mode 100644 index 00000000000..820f69d147c --- /dev/null +++ b/src/main/runtime/orchestration/orchestration-retired-delivery.test.ts @@ -0,0 +1,142 @@ +import { afterEach, describe, expect, it } from 'vitest' +import { OrchestrationDb } from './db' +import { createRootDispatch } from './db/root-dispatch-test-fixture' +import { reconcileLifecycleMessage } from './lifecycle-reconciliation' + +describe('retired mailbox deliveries', () => { + let db: OrchestrationDb + afterEach(() => db?.close()) + + function setup() { + db = new OrchestrationDb(':memory:') + const run = db.createRun({ + objective: 'Retired delivery', + coordinatorHandle: 'term_coord', + coordinatorPaneKey: 'tab:11111111-1111-4111-8111-111111111111' + }) + const params = { runId: run.id, consumerGeneration: run.consumer_generation } + const insert = (subject: string) => + db.insertMessage({ runId: run.id, from: 'worker', to: `run:${run.id}`, subject }) + return { run, params, insert } + } + + it('retires the heartbeat delivery atomically when completion suppresses its contents', () => { + const { run, params } = setup() + const task = db.createTask({ runId: run.id, spec: 'work' }) + const dispatch = createRootDispatch(db, task.id, 'worker') + const insert = (type: 'heartbeat' | 'worker_done') => + db.insertMessage({ + runId: run.id, + from: 'worker', + to: `run:${run.id}`, + subject: type, + type, + payload: JSON.stringify({ taskId: task.id, dispatchId: dispatch.id, outcome: 'succeeded' }) + }) + insert('heartbeat') + const first = db.getOrCreateRunDelivery(params)! + const done = insert('worker_done') + expect(reconcileLifecycleMessage(db, done).action).toBe('completed') + expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('acknowledged') + expect(db.hasOutstandingRunDelivery(run.id)).toBe(false) + expect( + db + .getOrCreateRunDelivery({ ...params, wakeTypes: ['worker_done'] }) + ?.messages.map((m) => m.id) + ).toEqual([done.id]) + expect(db.acknowledgeRunDelivery({ ...params, deliveryId: first.delivery.id }).duplicate).toBe( + true + ) + }) + + it('repairs an already fully-read outstanding delivery before replay', () => { + const { params, insert } = setup() + const old = insert('old') + const first = db.getOrCreateRunDelivery(params)! + db.db.prepare('UPDATE messages SET read = 1 WHERE id = ?').run(old.id) + const next = insert('next') + const current = db.getOrCreateRunDelivery(params)! + expect(current.messages.map((m) => m.id)).toEqual([next.id]) + expect(current.replayed).toBe(false) + expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('acknowledged') + }) + + it('preserves the entire replay batch while any member is unread', () => { + const { params, insert } = setup() + const a = insert('a') + const b = insert('b') + const first = db.getOrCreateRunDelivery(params)! + db.markAsReadAndDelivered([a.id]) + insert('later') + const replay = db.getOrCreateRunDelivery(params)! + expect(replay.delivery.id).toBe(first.delivery.id) + expect(replay.messages.map((m) => m.id)).toEqual([a.id, b.id]) + expect(replay.replayed).toBe(true) + }) + + it('rolls retirement back with the enclosing lifecycle transaction', () => { + const { params, insert } = setup() + const message = insert('old') + const first = db.getOrCreateRunDelivery(params)! + db.db.exec('BEGIN') + db.markAsReadAndDelivered([message.id]) + expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('acknowledged') + db.db.exec('ROLLBACK') + expect(db.getMessageById(message.id)?.read).toBe(0) + expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('outstanding') + }) + + it('names the owning delivery when ack receives a message ID without consuming mail', () => { + const { params, insert } = setup() + const message = insert('pending') + const first = db.getOrCreateRunDelivery(params)! + expect(() => db.acknowledgeRunDelivery({ ...params, deliveryId: message.id })).toThrow( + `Process the entire batch, then use --ack ${first.delivery.id}.` + ) + expect(db.getMessageById(message.id)?.read).toBe(0) + expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('outstanding') + }) + + it('does not reveal a delivery from another mailbox in a wrong-ID error', () => { + const { params, insert } = setup() + const foreign = db.createRun({ + objective: 'other', + coordinatorHandle: 'other', + coordinatorPaneKey: 'other:22222222-2222-4222-9222-222222222222' + }) + const message = insert('private') + const first = db.getOrCreateRunDelivery(params)! + expect(() => + db.acknowledgeRunDelivery({ + runId: foreign.id, + consumerGeneration: foreign.consumer_generation, + deliveryId: message.id + }) + ).toThrow('Run orchestration check to obtain the Delivery id.') + expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('outstanding') + }) + + it.each(['markAsRead', 'markAsReadAndDelivered'] as const)( + '%s retires dispatch mail without retiring a different mailbox', + (method) => { + const { run, params, insert } = setup() + insert('coordinator mail') + const coordinator = db.getOrCreateRunDelivery(params)! + const task = db.createTask({ runId: run.id, spec: 'worker mail' }) + const dispatch = createRootDispatch(db, task.id, 'worker') + const mailboxHandle = `dispatch:${dispatch.id}` + const message = db.insertMessage({ + runId: run.id, + from: 'term_coord', + to: mailboxHandle, + subject: 'worker mail' + }) + const workerParams = { ...params, mailboxHandle } + const worker = db.getOrCreateMailboxDelivery(workerParams)! + db[method]([message.id]) + expect(db.getDeliveryRaw(worker.delivery.id)?.status).toBe('acknowledged') + expect(db.getOrCreateMailboxDelivery(workerParams)).toBeUndefined() + expect(db.getDeliveryRaw(coordinator.delivery.id)?.status).toBe('outstanding') + } + ) +}) diff --git a/src/main/runtime/rpc/methods/orchestration/messaging/check-delivery-history.test.ts b/src/main/runtime/rpc/methods/orchestration/messaging/check-delivery-history.test.ts new file mode 100644 index 00000000000..59a0718b66a --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration/messaging/check-delivery-history.test.ts @@ -0,0 +1,33 @@ +import { afterEach, describe, expect, it } from 'vitest' +import { createOrchestrationRpcHarness } from '../rpc-test-harness' + +describe('Run delivery history', () => { + const h = createOrchestrationRpcHarness() + afterEach(() => h.cleanup()) + + it('exposes the outstanding delivery without minting or acknowledging it', async () => { + const { db, ctx, activeRunId } = h.setup() + const params = { terminal: 'term_coord', run: activeRunId, all: true } + db.insertMessage({ + from: 'worker', + to: `run:${activeRunId}`, + runId: activeRunId, + subject: 'waiting' + }) + expect(await h.call('orchestration.check', params, ctx)).toMatchObject({ + deliveryId: null, + count: 1 + }) + expect(db.hasOutstandingRunDelivery(activeRunId!)).toBe(false) + const delivery = db.getOrCreateRunDelivery({ + runId: activeRunId!, + consumerGeneration: db.getRun(activeRunId!)!.consumer_generation + })! + expect(await h.call('orchestration.check', { ...params, format: true }, ctx)).toMatchObject({ + deliveryId: delivery.delivery.id, + count: 1 + }) + expect(db.getDeliveryRaw(delivery.delivery.id)?.status).toBe('outstanding') + expect(db.getMessageById(delivery.messages[0].id)?.read).toBe(0) + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration/messaging/check-run.ts b/src/main/runtime/rpc/methods/orchestration/messaging/check-run.ts index 6db89cf4a20..598c080623a 100644 --- a/src/main/runtime/rpc/methods/orchestration/messaging/check-run.ts +++ b/src/main/runtime/rpc/methods/orchestration/messaging/check-run.ts @@ -101,6 +101,7 @@ export async function checkRunMailbox(args: { const result = { messages: exposeMessages(messages), count: messages.length, + deliveryId: db.getOutstandingMailboxDelivery(address)?.id ?? null, acknowledged: acknowledged?.delivery.id ?? null } if (params.format || params.inject) { diff --git a/src/shared/orchestration-check-output.test.ts b/src/shared/orchestration-check-output.test.ts index 00ffbb79da3..df3248bd2fd 100644 --- a/src/shared/orchestration-check-output.test.ts +++ b/src/shared/orchestration-check-output.test.ts @@ -1,5 +1,8 @@ import { describe, expect, it } from 'vitest' -import { prepareOrchestrationCheckOutput } from './orchestration-check-output' +import { + formatOrchestrationCheckText, + prepareOrchestrationCheckOutput +} from './orchestration-check-output' describe('prepareOrchestrationCheckOutput', () => { it('keeps mixed read-only mail safe and current Run replies executable', () => { @@ -40,3 +43,19 @@ describe('prepareOrchestrationCheckOutput', () => { expect(prepared.formatted).not.toContain('unsafe stale formatter output') }) }) + +describe('formatted delivery acknowledgment', () => { + it('retains the delivery ID above formatted message bodies', () => { + expect( + formatOrchestrationCheckText( + { + messages: [{ id: 'msg_one', from_handle: 'worker' }], + count: 1, + deliveryId: 'delivery_one', + formatted: 'Message body' + }, + 'term_coord' + ) + ).toBe('Delivery delivery_one\nMessage body') + }) +}) diff --git a/src/shared/orchestration-check-output.ts b/src/shared/orchestration-check-output.ts index e736b92f4dd..5429c542341 100644 --- a/src/shared/orchestration-check-output.ts +++ b/src/shared/orchestration-check-output.ts @@ -81,7 +81,7 @@ export function formatOrchestrationCheckText( : '' const deliveryNotice = formatCurrentDeliveryNotice(prepared.legacyCompatibility?.currentDelivery) if (prepared.formatted) { - return `${legacyHeader}${prepared.formatted}${deliveryNotice}` + return `${legacyHeader}${prepared.deliveryId ? `Delivery ${prepared.deliveryId}\n` : ''}${prepared.formatted}${deliveryNotice}` } if (prepared.count === 0) { if (prepared.timedOut) {