mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
fix(orchestration): retire read deliveries and clarify mailbox recovery
This commit is contained in:
@@ -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.'
|
||||
]
|
||||
},
|
||||
{
|
||||
|
||||
@@ -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
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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'
|
||||
)
|
||||
})
|
||||
|
||||
|
||||
@@ -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`
|
||||
}
|
||||
|
||||
@@ -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')
|
||||
}
|
||||
)
|
||||
})
|
||||
@@ -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)
|
||||
})
|
||||
})
|
||||
@@ -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) {
|
||||
|
||||
@@ -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')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user