mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 08:02:43 +00:00
fix(orchestration): simplify delivery recovery and update nudge contracts
This commit is contained in:
@@ -109,11 +109,9 @@ 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.',
|
||||
'--types is the wake condition for --wait; a returned Delivery is always the whole FIFO batch, so it is never filtered by type. Without --wait it has no effect on consuming checks. Only --peek and --all filter their rows.',
|
||||
'--format renders the returned rows as local text only; it never writes to another terminal.',
|
||||
'A bound Run replays the same Delivery until --ack or all its messages are retired; process every message before acknowledging.'
|
||||
'A bound Run replays the same Delivery until --ack or all its messages are marked read, even with --types; process every message before acknowledging.'
|
||||
]
|
||||
},
|
||||
{
|
||||
|
||||
+2
-2
@@ -94,7 +94,7 @@ describe('OrcaRuntimeService', () => {
|
||||
runtime.onPtyData('pty-1', '\x1b]0;Codex done\x07', 101)
|
||||
expect(write).toHaveBeenCalledWith(
|
||||
'pty-1',
|
||||
'\nYou have 1 orchestration message. Run `orca-dev orchestration check --run run_mailbox`.\n'
|
||||
'\n[Orca orchestration notification] You have 1 orchestration message. Run `orca-dev orchestration check --run run_mailbox`.\n'
|
||||
)
|
||||
expect(write).not.toHaveBeenCalledWith(
|
||||
'pty-1',
|
||||
@@ -472,7 +472,7 @@ describe('OrcaRuntimeService', () => {
|
||||
await vi.waitFor(() => {
|
||||
expect(write).toHaveBeenCalledWith(
|
||||
'pty-1',
|
||||
'\nYou have 1 orchestration message. Run `orca-dev orchestration check --run run_codex_native_title`.\n'
|
||||
'\n[Orca orchestration notification] You have 1 orchestration message. Run `orca-dev orchestration check --run run_codex_native_title`.\n'
|
||||
)
|
||||
})
|
||||
db.close()
|
||||
|
||||
+1
-1
@@ -62,7 +62,7 @@ describe('OrcaRuntimeService', () => {
|
||||
.map(([, data]) => data)
|
||||
.filter((data): data is string => typeof data === 'string')
|
||||
expect(payloads).toContain(
|
||||
'\nYou have 1 orchestration message. Run `orca-dev orchestration check --run run_test`.\n'
|
||||
'\n[Orca orchestration notification] You have 1 orchestration message. Run `orca-dev orchestration check --run run_test`.\n'
|
||||
)
|
||||
expect(payloads.some((data) => data.includes('reserved completion'))).toBe(false)
|
||||
expect(status.delivered_at).toEqual(expect.any(String))
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { DeliveryRow, MessageRow, MessageType } from '../../types'
|
||||
import { OrchestrationError } from '../../orchestration-error'
|
||||
import { generateId, isGeneratedId } from '../generated-id'
|
||||
import { generateId } 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,7 +43,9 @@ export function getOrCreateMailboxDelivery(
|
||||
if (params.requireCurrentRunConsumer) {
|
||||
this.requireCurrentConsumer(params.runId, params.consumerGeneration)
|
||||
}
|
||||
const existing = this.getOutstandingMailboxDelivery(params.mailboxHandle)
|
||||
const existing = this.db
|
||||
.prepare("SELECT * FROM deliveries WHERE mailbox_handle = ? AND status = 'outstanding'")
|
||||
.get(params.mailboxHandle) as DeliveryRow | undefined
|
||||
if (existing) {
|
||||
if (existing.consumer_generation !== params.consumerGeneration) {
|
||||
throw new OrchestrationError(
|
||||
@@ -51,12 +53,11 @@ export function getOrCreateMailboxDelivery(
|
||||
'This mailbox Delivery belongs to a fenced consumer generation.'
|
||||
)
|
||||
}
|
||||
const messages = this.getDeliveryMessages(existing)
|
||||
if (messages.some((message) => message.read === 0)) {
|
||||
if (!this.retireReadMailboxDelivery(params.mailboxHandle)) {
|
||||
const messages = this.getDeliveryMessages(existing)
|
||||
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(',')
|
||||
@@ -125,20 +126,6 @@ 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 ||
|
||||
@@ -147,7 +134,7 @@ export function acknowledgeMailboxDelivery(
|
||||
) {
|
||||
throw new OrchestrationError(
|
||||
'stale_delivery',
|
||||
`Delivery ${params.deliveryId} does not belong to this mailbox.`
|
||||
`Delivery ${params.deliveryId} does not belong to this mailbox. --ack requires a delivery_* ID returned by orchestration check; process the entire batch before acknowledging.`
|
||||
)
|
||||
}
|
||||
if (
|
||||
@@ -189,19 +176,9 @@ 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
|
||||
export function retireReadMailboxDelivery(this: OrchestrationDb, mailboxHandle: string): boolean {
|
||||
const result = this.db
|
||||
.prepare(
|
||||
`UPDATE deliveries SET status = 'acknowledged', acknowledged_at = datetime('now')
|
||||
WHERE mailbox_handle = ? AND status = 'outstanding'
|
||||
@@ -211,13 +188,20 @@ export function retireReadMailboxDelivery(this: OrchestrationDb, mailboxHandle:
|
||||
)`
|
||||
)
|
||||
.run(mailboxHandle)
|
||||
return result.changes > 0
|
||||
}
|
||||
|
||||
export function hasOutstandingMailboxDelivery(
|
||||
this: OrchestrationDb,
|
||||
mailboxHandle: string
|
||||
): boolean {
|
||||
return this.getOutstandingMailboxDelivery(mailboxHandle) !== undefined
|
||||
return Boolean(
|
||||
this.db
|
||||
.prepare(
|
||||
"SELECT 1 FROM deliveries WHERE mailbox_handle = ? AND status = 'outstanding' LIMIT 1"
|
||||
)
|
||||
.get(mailboxHandle)
|
||||
)
|
||||
}
|
||||
|
||||
export function fenceOutstandingMailboxDelivery(
|
||||
@@ -232,7 +216,6 @@ export function fenceOutstandingMailboxDelivery(
|
||||
}
|
||||
|
||||
export type RoleMailboxDeliveryMethods = {
|
||||
getOutstandingMailboxDelivery: typeof getOutstandingMailboxDelivery
|
||||
retireReadMailboxDelivery: typeof retireReadMailboxDelivery
|
||||
getDeliveryRaw: typeof getDeliveryRaw
|
||||
getDeliveryMessages: typeof getDeliveryMessages
|
||||
@@ -244,7 +227,6 @@ export type RoleMailboxDeliveryMethods = {
|
||||
|
||||
export function attachRoleMailboxDelivery(ctor: { prototype: object }): void {
|
||||
Object.assign(ctor.prototype, {
|
||||
getOutstandingMailboxDelivery,
|
||||
retireReadMailboxDelivery,
|
||||
getDeliveryRaw,
|
||||
getDeliveryMessages,
|
||||
|
||||
@@ -86,33 +86,29 @@ describe('retired mailbox deliveries', () => {
|
||||
expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('outstanding')
|
||||
})
|
||||
|
||||
it('names the owning delivery when ack receives a message ID without consuming mail', () => {
|
||||
it('explains the delivery ID contract for invalid acknowledgements 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}.`
|
||||
'--ack requires a delivery_* ID returned by orchestration check; process the entire batch before acknowledging.'
|
||||
)
|
||||
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')
|
||||
it('checks the consumer generation before repairing an already read delivery', () => {
|
||||
const { run, params, insert } = setup()
|
||||
const message = insert('old')
|
||||
const first = db.getOrCreateRunDelivery(params)!
|
||||
db.db.prepare('UPDATE messages SET read = 1 WHERE id = ?').run(message.id)
|
||||
expect(() =>
|
||||
db.acknowledgeRunDelivery({
|
||||
runId: foreign.id,
|
||||
consumerGeneration: foreign.consumer_generation,
|
||||
deliveryId: message.id
|
||||
db.getOrCreateMailboxDelivery({
|
||||
...params,
|
||||
mailboxHandle: `run:${run.id}`,
|
||||
consumerGeneration: params.consumerGeneration + 1
|
||||
})
|
||||
).toThrow('Run orchestration check to obtain the Delivery id.')
|
||||
).toThrow(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
expect(db.getDeliveryRaw(first.delivery.id)?.status).toBe('outstanding')
|
||||
})
|
||||
|
||||
|
||||
+18
-5
@@ -5,7 +5,7 @@ describe('Run delivery history', () => {
|
||||
const h = createOrchestrationRpcHarness()
|
||||
afterEach(() => h.cleanup())
|
||||
|
||||
it('exposes the outstanding delivery without minting or acknowledging it', async () => {
|
||||
it('does not label filtered history as an acknowledgeable delivery', async () => {
|
||||
const { db, ctx, activeRunId } = h.setup()
|
||||
const params = { terminal: 'term_coord', run: activeRunId, all: true }
|
||||
db.insertMessage({
|
||||
@@ -15,7 +15,6 @@ describe('Run delivery history', () => {
|
||||
subject: 'waiting'
|
||||
})
|
||||
expect(await h.call('orchestration.check', params, ctx)).toMatchObject({
|
||||
deliveryId: null,
|
||||
count: 1
|
||||
})
|
||||
expect(db.hasOutstandingRunDelivery(activeRunId!)).toBe(false)
|
||||
@@ -23,10 +22,24 @@ describe('Run delivery history', () => {
|
||||
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
|
||||
db.insertMessage({
|
||||
from: 'worker',
|
||||
to: `run:${activeRunId}`,
|
||||
runId: activeRunId,
|
||||
subject: 'later completion',
|
||||
type: 'worker_done'
|
||||
})
|
||||
const history = await h.call(
|
||||
'orchestration.check',
|
||||
{
|
||||
...params,
|
||||
format: true,
|
||||
types: 'worker_done'
|
||||
},
|
||||
ctx
|
||||
)
|
||||
expect(history).toMatchObject({ count: 1, messages: [{ subject: 'later completion' }] })
|
||||
expect(history).not.toHaveProperty('deliveryId')
|
||||
expect(db.getDeliveryRaw(delivery.delivery.id)?.status).toBe('outstanding')
|
||||
expect(db.getMessageById(delivery.messages[0].id)?.read).toBe(0)
|
||||
})
|
||||
|
||||
@@ -101,7 +101,6 @@ 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) {
|
||||
|
||||
Reference in New Issue
Block a user