mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 08:01:56 +00:00
fix(orchestration): point a structured chat's mail by the same rule as a terminal's
The chat lane pointed newer mail past a batch its reader had checked and not acknowledged, while the terminal lane skips a mailbox until that batch is acked. A chat now waits for the ack the same way a terminal does, and the lookup that let the chat lane filter the held batch out is removed.
This commit is contained in:
@@ -184,31 +184,6 @@ export function hasOutstandingMailboxDelivery(
|
||||
)
|
||||
}
|
||||
|
||||
/** The batch a consumer has read and not yet acknowledged on this mailbox, if any. */
|
||||
export function getOutstandingMailboxDelivery(
|
||||
this: OrchestrationDb,
|
||||
mailboxHandle: string
|
||||
): { id: string; messageIds: ReadonlySet<string> } | undefined {
|
||||
const row: unknown = this.db
|
||||
.prepare('SELECT id, message_ids FROM outstanding_deliveries WHERE mailbox_handle = ? LIMIT 1')
|
||||
.get(mailboxHandle)
|
||||
if (
|
||||
!row ||
|
||||
typeof row !== 'object' ||
|
||||
!('id' in row) ||
|
||||
typeof row.id !== 'string' ||
|
||||
!('message_ids' in row) ||
|
||||
typeof row.message_ids !== 'string'
|
||||
) {
|
||||
return undefined
|
||||
}
|
||||
const ids: unknown = JSON.parse(row.message_ids)
|
||||
return {
|
||||
id: row.id,
|
||||
messageIds: new Set(Array.isArray(ids) ? ids.filter((id) => typeof id === 'string') : [])
|
||||
}
|
||||
}
|
||||
|
||||
export function fenceUnacknowledgedMailboxDeliveries(
|
||||
this: OrchestrationDb,
|
||||
mailboxHandle: string
|
||||
@@ -226,7 +201,6 @@ export type RoleMailboxDeliveryMethods = {
|
||||
getOrCreateMailboxDelivery: typeof getOrCreateMailboxDelivery
|
||||
acknowledgeMailboxDelivery: typeof acknowledgeMailboxDelivery
|
||||
hasOutstandingMailboxDelivery: typeof hasOutstandingMailboxDelivery
|
||||
getOutstandingMailboxDelivery: typeof getOutstandingMailboxDelivery
|
||||
fenceUnacknowledgedMailboxDeliveries: typeof fenceUnacknowledgedMailboxDeliveries
|
||||
}
|
||||
|
||||
@@ -237,7 +211,6 @@ export function attachRoleMailboxDelivery(ctor: { prototype: object }): void {
|
||||
getOrCreateMailboxDelivery,
|
||||
acknowledgeMailboxDelivery,
|
||||
hasOutstandingMailboxDelivery,
|
||||
getOutstandingMailboxDelivery,
|
||||
fenceUnacknowledgedMailboxDeliveries
|
||||
})
|
||||
}
|
||||
|
||||
@@ -4,7 +4,6 @@ import {
|
||||
OrchestrationStructuredMailboxPointerDelivery,
|
||||
type StructuredMailboxPointerHost
|
||||
} from './structured-mailbox-pointer-delivery'
|
||||
import { formatMessagePointer } from './formatter'
|
||||
import type { StructuredPointerOperationRow } from './db/messages/structured-pointer-operation-store'
|
||||
import {
|
||||
structuredPointerBatchFingerprint,
|
||||
@@ -83,8 +82,6 @@ function harness(options: {
|
||||
/** The coordinator of this worker's Run is mid-batch: it checked and has not acked yet. */
|
||||
outstandingRunDelivery?: boolean
|
||||
outstandingOwnDelivery?: boolean
|
||||
/** Undelivered unread rows on the mailbox, oldest first. */
|
||||
unreadIds?: string[]
|
||||
/** The mailbox this worker owns; its own handle for direct peer mail outside a dispatch. */
|
||||
mailbox?: string
|
||||
dispatchId?: string | null
|
||||
@@ -103,18 +100,10 @@ function harness(options: {
|
||||
const stored = new Map<string, StructuredPointerOperationRow>()
|
||||
const db = {
|
||||
getDispatchContextById: () => ({ run_id: 'run_1' }),
|
||||
// The reader holds `m1` unacknowledged, on whichever mailbox the option names.
|
||||
getOutstandingMailboxDelivery: (handle: string) =>
|
||||
hasOutstandingMailboxDelivery: (handle: string) =>
|
||||
((options.outstandingRunDelivery ?? false) && handle.startsWith('run:')) ||
|
||||
((options.outstandingOwnDelivery ?? false) && !handle.startsWith('run:'))
|
||||
? { id: 'delivery_held', messageIds: new Set(['m1']) }
|
||||
: undefined,
|
||||
getUndeliveredUnreadMessages: () =>
|
||||
(options.unreadIds ?? ['m1']).map((id, index) => ({
|
||||
id,
|
||||
type: 'status',
|
||||
sequence: index + 3
|
||||
})),
|
||||
((options.outstandingOwnDelivery ?? false) && !handle.startsWith('run:')),
|
||||
getUndeliveredUnreadMessages: () => [{ id: 'm1', type: 'status', sequence: 3 }],
|
||||
markAsDelivered,
|
||||
getStructuredPointerOperation: (key: string) => stored.get(key),
|
||||
putStructuredPointerOperation: (row: StructuredPointerOperationRow) =>
|
||||
@@ -275,32 +264,15 @@ describe('structured mailbox pointer delivery', () => {
|
||||
expect(markAsDelivered).toHaveBeenCalledWith(['m1'])
|
||||
})
|
||||
|
||||
it('does not re-point mail the reader already holds unacknowledged', async () => {
|
||||
// It already has this batch, so a second pointer spends a whole provider turn telling it
|
||||
// something it was told.
|
||||
it('does not re-nudge a mailbox still holding its own unacked batch', async () => {
|
||||
// The other half of the same gate: the consumer already has this batch, so a second nudge
|
||||
// spends a whole provider turn telling it something it was told.
|
||||
const { delivery, send } = harness({ journal: idleJournal(), outstandingOwnDelivery: true })
|
||||
delivery.deliverForHandle('dispatch:d1')
|
||||
await flush()
|
||||
expect(send).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('points newer mail while the reader holds an unacknowledged batch, in the PTY pointer text', async () => {
|
||||
// The strand this pins: a chat reads a result, ends its turn without acking, and the gate on
|
||||
// "an unacknowledged batch exists" silenced every later result. The text stays the PTY lane's.
|
||||
const { delivery, send, markAsDelivered } = harness({
|
||||
journal: idleJournal(),
|
||||
outstandingOwnDelivery: true,
|
||||
unreadIds: ['m1', 'm2']
|
||||
})
|
||||
delivery.deliverForHandle('dispatch:d1')
|
||||
await flush()
|
||||
expect(send).toHaveBeenCalledTimes(1)
|
||||
expect(send.mock.calls[0]![0].body.blocks[0]).toMatchObject({
|
||||
text: formatMessagePointer(1, 'dispatch:d1', 'orca-dev').trim()
|
||||
})
|
||||
expect(markAsDelivered).toHaveBeenCalledWith(['m2'])
|
||||
})
|
||||
|
||||
it('retries a rejected nudge on the next journal edge, under the same id', async () => {
|
||||
// A rejection consumes no mail and nothing else redrives this mailbox, so leaving it unparked
|
||||
// stranded the worker until unrelated mail happened to arrive. The retry keeps the id: the host
|
||||
@@ -419,7 +391,7 @@ describe('forgetting one settled worker', () => {
|
||||
}))
|
||||
const db = {
|
||||
getDispatchContextById: () => ({ run_id: 'run_1' }),
|
||||
getOutstandingMailboxDelivery: () => undefined,
|
||||
hasOutstandingMailboxDelivery: () => false,
|
||||
getUndeliveredUnreadMessages: () => [{ id: 'm1', type: 'status', sequence: 3 }],
|
||||
markAsDelivered: vi.fn(),
|
||||
getStructuredPointerOperation: () => undefined,
|
||||
|
||||
@@ -163,16 +163,20 @@ export class OrchestrationStructuredMailboxPointerDelivery<
|
||||
if (!db || this.inFlight.has(mailboxHandle)) {
|
||||
return
|
||||
}
|
||||
// Eligibility is "not yet pointed" (`delivered_at`), never "has the consumer acked": a chat that
|
||||
// reads a batch and ends its turn without acking must still be pointed at the NEXT result. The
|
||||
// batch it holds is excluded; its own `check` replays that batch and names its ack.
|
||||
const outstanding = db.getOutstandingMailboxDelivery?.(mailboxHandle)
|
||||
// Don't re-nudge a mailbox whose consumer still holds an unacknowledged batch. The lookup is
|
||||
// keyed on the exact handle being nudged, so a coordinator's own `run:` delivery is invisible
|
||||
// to a worker's `dispatch:` gate and cannot suppress the nudges a coordinator sends its
|
||||
// workers. Worth more here than in the PTY lane: a structured nudge costs a whole provider
|
||||
// turn, not a line of text into a composer.
|
||||
if (db.hasOutstandingMailboxDelivery?.(mailboxHandle)) {
|
||||
return
|
||||
}
|
||||
const unread = selectOrchestrationPointerBatch({
|
||||
db,
|
||||
mailboxHandle,
|
||||
waiters: this.deps.getMessageWaiters(mailboxHandle),
|
||||
reservedTypes
|
||||
}).filter((message) => !outstanding?.messageIds.has(message.id))
|
||||
})
|
||||
if (unread.length === 0) {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -338,43 +338,6 @@ describe('a worker result reaches the structured chat that coordinates it', () =
|
||||
expect(await userTexts(COORDINATOR)).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('points the next result at a coordinator that read the last one without acking', async () => {
|
||||
// The strand this pins: a flagless `check` opens a delivery that `check` replays until acked,
|
||||
// and a lane gated on "an unacknowledged batch exists" never pointed the chat at a later result.
|
||||
const chat = await openChat(COORDINATOR)
|
||||
const { runId, taskId } = await coordinatorRunAndTask()
|
||||
await finishWorker(taskId)
|
||||
await vi.waitFor(() => expect(chat.turns).toHaveLength(1), WAIT)
|
||||
await settleTurn(COORDINATOR, 0)
|
||||
const first = await call('orchestration.check', {}, { sessionId: COORDINATOR })
|
||||
const heldDelivery = String(first.deliveryId)
|
||||
expect(first).toMatchObject({ count: 1, messages: [{ type: 'worker_done' }] })
|
||||
|
||||
const second = await call(
|
||||
'orchestration.taskCreate',
|
||||
{ spec: 'more' },
|
||||
{ sessionId: COORDINATOR }
|
||||
)
|
||||
await finishWorker(idOf(second.task), { handle: 'term_worker_2', paneKey: WORKER_2_PANE })
|
||||
await vi.waitFor(() => expect(chat.turns).toHaveLength(2), WAIT)
|
||||
// The PTY lane's text: `check` itself replays the held batch and names its ack.
|
||||
expect(turnText(chat.turns[1]!)).toBe(ptyPointer(`run:${runId}`))
|
||||
await settleTurn(COORDINATOR, 1)
|
||||
|
||||
// Exactly once per new message: a retry and the idle edge point nothing further.
|
||||
runtime.deliverPendingMessagesForHandle(`run:${runId}`)
|
||||
await new Promise((resolve) => setTimeout(resolve, 20))
|
||||
expect(chat.turns).toHaveLength(2)
|
||||
|
||||
const acked = await call(
|
||||
'orchestration.check',
|
||||
{ ack: heldDelivery },
|
||||
{ sessionId: COORDINATOR }
|
||||
)
|
||||
expect(acked).toMatchObject({ acknowledged: heldDelivery, count: 1 })
|
||||
expect(acked.messages).not.toEqual(first.messages)
|
||||
})
|
||||
|
||||
/** Fires both edges and waits until every gate read they started has answered. */
|
||||
async function edgesAnswered(): Promise<void> {
|
||||
const reads = vi.spyOn(host, 'journalSnapshot')
|
||||
|
||||
Reference in New Issue
Block a user