From dd39dcba5f99710251bac1a33392bb41a1fd3945 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Mon, 5 Oct 2026 17:34:06 -0700 Subject: [PATCH] feat(orchestration): a native chat gets the orchestration pointer a CLI agent gets, through the same send as your messages (#25078) * feat(orchestration): agent mail to a busy chat waits in the chat's own queue A mail notice for a busy structured chat used to wait in the orchestration lane's own invisible "until the chat is free" gate. It now goes through the queue a person's message uses: sendAgentTurn(..., { delivery: 'queue', source }) makes the host hold it as a draft card the person sees and can Steer or delete, and the queue sends it when the turn ends. - The lane's busy gate (turnRunning / awaitingHuman) is gone; the queue's own hold decides. Its parking now only waits for its own send or card to settle. - A queued card's hand-off (sent under a fresh id) is matched by queuedMessageId, so it stamps the mail once and never sends a second notice. - A card from before an Orca restart is still the lane's: no second card. - A card the person deleted counts as handled for that batch; newer mail notifies again. No stored flag: read off the card and delivered_at. - `orchestration check` waits for the lane to withdraw a card whose mail it read, so a stale notice is never sent; more mail replaces the card. - Each card records who queued it in queued_messages.source_json (versioned, schema-validated): the person, or Orca for agents, with every distinct sender as an orchestration party plus host id, and the mailbox, dispatch, run and message ids. The restart pause holds only the person's cards. - sendAgentTurn answers with a snapshot, not the journal's live submission. * test(orchestration): type the mail fixture as a pointer batch message * fix(orchestration): the queue judges an agent's mail notice as it sends it A mail notice queued in a busy chat could go out stale or twice: it was kept true from outside the queue, by orchestration check withdrawing it, which missed other readers, raced an attempt in flight, lost the card when /clear carried it (a second notice), and moved it to the back when new mail replaced it. An agent's card behind a person's paused card never sent after a restart, so an unattended coordinator stalled. - The drain asks the card's sender, in its own serialized step, whether it still stands: send, restate (count and senders of the mail owed now, written onto the card in the hand-off's transaction), or withdraw as the host. A failing judge sends as written. Send-now restates but never withdraws. Orchestration registers the judge on the host. - check no longer touches any queue (back to main); onMailRead, notifyOrchestrationMailRead and the lane's withdrawals are gone. - One unsettled card per agent message, enforced at insert; new mail is counted when the card sends, never replaces or moves it. - The lane finds its card by what it is (a notice for the mailbox in the session the mailbox reaches now); a decline is a withdrawal a person's operation stamped. /clear moves agent cards as the host. Pointer rows last only while a direct send is in flight and end with the session. - Pauses (restart, Stop, /clear) hold only a person's cards, and a held card never traps an agent's card behind it, keyed on the message kind. - The source is named for the message (agent-session-message-source), one payload shape per message kind, no relative host id. * fix(orchestration): Orca's mail notice waits in the queue out of sight The notice card read as something the person typed and came and went on its own, the transient-state message the queue should not show. It is hidden until B can label who sent it. - The host leaves mail-notice cards out of the published queue, so the list, its count, Steer/Delete/Edit, steer-newest and the paused header never see them, on every client. Keyed on the message kind: D6's task card will be shown, labelled, in B. They still count toward the queue limit (temporary, until B). - A refusal no longer returns a hidden card for the person to act on: the host withdraws it, and the lane points again only after a turn ran since, the rule it already had for a refused notice. - Send-now no longer restates a card: no one can reach a hidden one. - The stored sender drops the pane key, a mailbox credential. - The gate covers an abandoned worker: a mailbox that reaches no session withdraws its notice (tested). * fix(orchestration): a hidden mail notice never waits on the person Hiding the notice left three paths where it waited on someone who cannot see it. - A failed hand-off write put a send_failed hold on it, which only a person's Send or Delete clears: that mailbox's notices stopped for good. The host now withdraws a hidden card instead and hands its mailbox back to the lane through the ordinary redrive, since an idle chat gets no other edge. A write that keeps failing is retried once per edge. - A person's Stop while the agent started on the notice put it back to waiting, where no pause holds it, so it was sent again at once. The host now withdraws it; the lane points again after a turn ran or when newer mail makes it a different notice, as on main. A restart still puts it back. - A notice queued before the person's message went first. One order rule now: the person's cards go in order and never past one waiting on them; a card they cannot see never delays one they can, and goes only when none of theirs may. The drain moves to its own module (structured-agent-session-queued- drain.ts). Comments that still described the notice as visible are corrected. * fix(orchestration): an accepted notice hand-off stamps its mail; a judge that cannot look decides nothing - When the agent opened its mail in the notice's own turn (the normal flow), an open batch left nothing owed, so the lane never marked the mail the accepted hand-off carried as delivered, unlike an accepted direct send. The lane now stamps exactly the ids a handed-off notice carried whenever undelivered mail remains, owed or not. A later conversation is no longer told again about mail already opened. - At quit the host registry is cleared before teardown, so the drain's judge could not resolve the mailbox and withdrew the notice. A judge that cannot read its inputs now defers: the card waits for a step that can look. A mailbox that resolves to another session or none is still withdrawn. - The judge, the owed-batch selection and the notice body move to structured-mail-notice.ts. A failing hand-back of a dropped card is logged. * refactor(orchestration): a chat receives the agent's mail itself, queued like a person's message A structured chat used to get a derived "You have N orchestration messages, run check" notice, which could go stale while it waited in the chat's queue, so the queue judged, restated or withdrew it as it sent. It now gets the mail itself as the turn, through sendAgentTurn with delivery 'queue', the path a person's message takes. - One turn per mailbox batch: the mailbox's unread mail at delivery, in mail order, each message led by "[message from ]", its type, subject, body, payload and reply hint (what check prints). Mail arriving while that card is still in the queue waits for the next card; a card is never edited. - An idle chat takes it at once; a busy one queues it as an ordinary card the person sees and can Steer or delete. Mail is marked read when the chat accepts the turn, so check does not return it again; a deleted card leaves its mail unread for check, and it is not pushed again. - The card stores who it is from (source_json): every distinct sender, and each message's id, run and sender. Kind 'mail-notice' becomes 'mail'. - Removed: the send-time judge (QueuedAgentCardVerdict), structured-mail- notice, structured-pointer-notice-cards, queued-message-restatement, the hidden-card rules, the agent-card exemptions from the restart, Stop and /clear pauses, the dropped-card hand-back, and the drain split. An agent's card now follows the same pauses as a person's. - Terminal agents are unchanged: they keep the typed pointer and check. * fix(orchestration): a chat's check skips mail already queued to it; restarts hold only the person's cards - When a structured chat runs `orchestration check` (consuming, --peek or --wait), mail an agent's card in that chat's own queue still carries (any card not deleted) is left out: it is on its way as a turn, so the agent does not read it twice. Derived from the queued rows at read time; a deleted card no longer carries it, so it comes back. Terminal callers and every other read are unchanged. New batches pass the exclusion to getOrCreateMailboxDelivery; peeks filter it. - A restart's queue pause now holds only the person's cards: an unattended coordinator's agent card sends after an app restart without waiting for a Resume, since the mailbox is the record and reading is marked. Stop and /clear still hold agent cards. One rule, in queuePauseHolding. An agent card queued behind the person's held card still waits behind it: the queue never reorders. - The duplicated direct-mailbox snapshot routing in check-run and check-worker becomes one helper, which keeps check-worker under its line limit. * fix(orchestration): the mailbox is the only record of read mail; a chat's check takes its waiting cards A card held an exclusive claim on its mail that nothing reconciled with the mailbox: queueing it stamped the mail delivered, and a chat's check hid every card that was not deleted. So a chat's own check --wait in one turn could not see a result its waiting card held, a hand-off that ended in doubt or came back "Not sent" stranded its mail, a check racing the card's build read the mail twice, and a card left behind by run-use still sent and marked it read. What a card or send holds is now derived each time from the chat's queue and its sends: an accepted send marks its mail read; a waiting card or an unanswered send holds its mail; everything else unread is pushed again or returned by check. A card the person deleted is recorded as not to be pushed, at the delete. The lane withdraws returned, partly read and moved cards. A chat's consuming check withdraws the waiting cards holding what it reads, one at a time with the lane. The sender is stored on the sent submission too (host-only), so read state survives the card row's prune. A refusal before anything started ends its operation; a restart pause is raised only by the person's cards; an unreadable agent source stays an agent's. * refactor(orchestration): a chat gets the pointer a terminal agent gets, sent through the chat's own queue A structured chat now receives exactly what a terminal agent is typed: "You have N orchestration messages ... run `orca orchestration check`" (formatMessagePointer, same CLI name). It goes through the shared sendAgentTurn with delivery 'queue', the composer's queue-if-active: an idle chat gets it at once, a busy chat's queue holds it as a normal card and sends it when the turn ends. The lane's idle gate is gone; the send decides, as for the person's message. The lane queues no second pointer while one of its cards still waits, read from the queued rows; mail arriving meanwhile, or after the card is sent, is pointed again as the terminal lane points new unread mail. A queued pointer counts as delivered, as an accepted one does. `check` and mail read state are main's: only `check` reads mail. The card records who it is from (queued_messages.source_json, kind 'agent', message 'mail-notice', with its senders) for the next PR to render. A restart's pause is raised by and holds only the person's cards; Stop and /clear still hold an agent's. Removed from the earlier designs: the message-as-turn batching, the check exclusion and taking, the derived claims, lock and reconciliation, the sender on submissions, and their tests. * fix(orchestration): a chat's pointer follows a card that leaves its queue unsent The lane holds new mail while a pointer card waits, and retried it only on the chat's next status change. A card that leaves the queue without a turn after it (the person deletes it while Stop holds it, or its hand-off comes back) made none, so that mail sat unpointed until something unrelated happened. The shared queue wiring now tells its host when an agent's card stops waiting, read only when the queue changed, and the runtime redrives that chat's parked mail. The runtime forwards the new host dep like the others. Tests: that case end to end, and /clear carrying an agent card keeps who it is from. Agent card bodies in tests are the pointer text; the restart pause's header says why an agent card may send. * refactor(orchestration): no special handling for agent notices in the chat queue A structured chat's orchestration notice is now sent through the chat's own send like any message: an idle chat takes it as a turn, a busy chat's queue holds it and sends it at turn end, under the same Stop, restart and /clear pauses as the person's cards. Native chat no longer branches on who a card is from; the card only records it. - Restore main's queue pause logic (no restart exemption for agent cards). - Remove the queue watch that redrove mail when an agent card left the queue, and the host's queued-row read the lane used for it. - The lane reads nothing of the chat's queue. It sends no second notice while mail it already pointed at is unread, read off the mailbox alone. * fix(orchestration): point new mail like a terminal, even while earlier pointed mail is unread Drops the native-chat-only rule that held back a notice while the agent had not yet read mail it was already pointed at. Also fixes main's stop-note test, which still built a provider handle in the shape #24991 replaced. * test(native-chat): keep the queued-message rig fixture under the line limit * fix(test): drop main's duplicate codexProviderHandle import (same line as #25713) --- .../journal-queued-messages.ts | 2 + .../journal-submission-positions.test.ts | 1 - .../journal-submission-queued-link.test.ts | 6 +- ...queued-message-bookkeeping-failure.test.ts | 6 +- .../queued-message-delivered-echo.test.ts | 3 +- .../queued-message-pause.test.ts | 4 +- .../queued-message-schema.ts | 4 +- .../queued-message-store.test.ts | 51 ++- .../queued-message-table.ts | 32 +- ...structured-agent-session-host-mutations.ts | 4 + ...session-queued-message-rig.test-fixture.ts | 12 +- ...ured-agent-session-queued-messages.test.ts | 26 ++ ...tructured-agent-session-queued-messages.ts | 11 +- ...ructured-agent-session-queued-mutations.ts | 3 +- .../structured-agent-session-queued-send.ts | 2 + ...ent-session-restore-without-import.test.ts | 3 +- .../orchestration/agent-facing-parity.test.ts | 6 +- .../orchestration-caller-identity.ts | 19 +- .../send-agent-turn-host.test.ts | 36 +- .../orchestration/send-agent-turn.test.ts | 20 +- .../runtime/orchestration/send-agent-turn.ts | 24 +- .../structured-mail-source.test.ts | 43 +++ .../orchestration/structured-mail-source.ts | 53 +++ ...tructured-mailbox-pointer-delivery.test.ts | 210 ++++------ .../structured-mailbox-pointer-delivery.ts | 48 ++- .../structured-mailbox-pointer-host.test.ts | 57 ++- .../structured-mailbox-pointer-host.ts | 30 +- ...tructured-session-pointer-delivery.test.ts | 32 -- .../structured-session-pointer-delivery.ts | 56 +-- ...ctured-chat-coordinator-mail-queue.test.ts | 139 +++++++ ...-chat-coordinator-mail-rig.test-fixture.ts | 324 ++++++++++++++++ .../structured-chat-coordinator-mail.test.ts | 359 +++--------------- src/shared/agent-session-message-source.ts | 88 +++++ src/shared/orchestration-party-identity.ts | 18 + 34 files changed, 1103 insertions(+), 629 deletions(-) create mode 100644 src/main/runtime/orchestration/structured-mail-source.test.ts create mode 100644 src/main/runtime/orchestration/structured-mail-source.ts create mode 100644 src/main/runtime/structured-chat-coordinator-mail-queue.test.ts create mode 100644 src/main/runtime/structured-chat-coordinator-mail-rig.test-fixture.ts create mode 100644 src/shared/agent-session-message-source.ts create mode 100644 src/shared/orchestration-party-identity.ts diff --git a/src/main/native-chat/agent-session-journal/journal-queued-messages.ts b/src/main/native-chat/agent-session-journal/journal-queued-messages.ts index 0828740ca6f..8c190196bdf 100644 --- a/src/main/native-chat/agent-session-journal/journal-queued-messages.ts +++ b/src/main/native-chat/agent-session-journal/journal-queued-messages.ts @@ -33,6 +33,7 @@ import { type QueuedMessageHoldReason, type QueuedMessageRow } from './queued-message-table' +import type { AgentSessionMessageSource } from '../../../shared/agent-session-message-source' import { draftsDeliveredByAppliedEcho } from './queued-message-delivered-echo' import { pruneQueuedMessages, retainedSubmissionVerdict } from './queued-message-retention' import { @@ -110,6 +111,7 @@ export class JournalQueuedMessages { fingerprint: string hostInstance: string carriedFrom?: string + source: AgentSessionMessageSource }, receipt?: JournalOperationReceipt ): Promise { diff --git a/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts b/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts index 37416378605..3e05bbdbbcc 100644 --- a/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-submission-positions.test.ts @@ -18,7 +18,6 @@ import { type AgentJournalMessageItem, type AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types' -import { codexProviderHandle } from '../../../shared/agent-session-provider-handle-encoding' import { readAgentSessionHydrationPage } from '../agent-session-wire/agent-session-history-page' import { createTrackedJournalOpener } from './journal-host-database-test-support' diff --git a/src/main/native-chat/agent-session-journal/journal-submission-queued-link.test.ts b/src/main/native-chat/agent-session-journal/journal-submission-queued-link.test.ts index f9588b71ba3..d0306fa9467 100644 --- a/src/main/native-chat/agent-session-journal/journal-submission-queued-link.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-submission-queued-link.test.ts @@ -51,7 +51,8 @@ async function handOff(journal: AgentSessionJournal, draftId: string, submission messageId: draftId, body: BODY, fingerprint: 'fp', - hostInstance: 'p' + hostInstance: 'p', + source: { kind: 'user' } }) await journal.appendSubmission( { clientMessageId: submissionId, payloadFingerprint: 'fp', body: BODY, fence: 0 }, @@ -100,7 +101,8 @@ describe('the submission names the queued draft it hands off', () => { messageId: 'draft-1', body: BODY, fingerprint: 'fp', - hostInstance: 'p' + hostInstance: 'p', + source: { kind: 'user' } }) await expect( journal.appendSubmission( diff --git a/src/main/native-chat/agent-session-journal/queued-message-bookkeeping-failure.test.ts b/src/main/native-chat/agent-session-journal/queued-message-bookkeeping-failure.test.ts index 5f9a95a53d2..60b613ea307 100644 --- a/src/main/native-chat/agent-session-journal/queued-message-bookkeeping-failure.test.ts +++ b/src/main/native-chat/agent-session-journal/queued-message-bookkeeping-failure.test.ts @@ -65,7 +65,8 @@ async function queueAndConsume(journal: AgentSessionJournal, messageId: string): messageId, body, fingerprint: `fp-${messageId}`, - hostInstance: 'proc-1' + hostInstance: 'proc-1', + source: { kind: 'user' } }) await journal.appendSubmission( { @@ -172,7 +173,8 @@ describe('draft bookkeeping inside a journal append', () => { messageId: 'draft-1', body: BODY, fingerprint: 'fp-draft-1', - hostInstance: 'proc-1' + hostInstance: 'proc-1', + source: { kind: 'user' } }) expect(journal.queuedMessages.list()).toMatchObject([{ state: 'waiting' }]) const commit = failNextCommit() diff --git a/src/main/native-chat/agent-session-journal/queued-message-delivered-echo.test.ts b/src/main/native-chat/agent-session-journal/queued-message-delivered-echo.test.ts index a12630d6566..88caf16949a 100644 --- a/src/main/native-chat/agent-session-journal/queued-message-delivered-echo.test.ts +++ b/src/main/native-chat/agent-session-journal/queued-message-delivered-echo.test.ts @@ -77,7 +77,8 @@ async function handOffAndReject( messageId: 'draft-1', body, fingerprint, - hostInstance: 'p' + hostInstance: 'p', + source: { kind: 'user' } }) await journal.appendSubmission( { diff --git a/src/main/native-chat/agent-session-journal/queued-message-pause.test.ts b/src/main/native-chat/agent-session-journal/queued-message-pause.test.ts index a2fa19a6a5c..5989a10e967 100644 --- a/src/main/native-chat/agent-session-journal/queued-message-pause.test.ts +++ b/src/main/native-chat/agent-session-journal/queued-message-pause.test.ts @@ -63,6 +63,7 @@ function queueDraft(journal: AgentSessionJournal, messageId: string, carriedFrom body: message(messageId), fingerprint: `fp-${messageId}`, hostInstance: HOST, + source: { kind: 'user' }, ...(carriedFrom ? { carriedFrom } : {}) }) } @@ -495,7 +496,8 @@ describe("a restart's pause", () => { messageId: 'draft-restart', body: message('written before the restart'), fingerprint: 'fp-draft-restart', - hostInstance: 'proc-0' + hostInstance: 'proc-0', + source: { kind: 'user' } }) expect(reason(journal)).toBe('restarted') await queueDraft(journal, 'draft-legacy') diff --git a/src/main/native-chat/agent-session-journal/queued-message-schema.ts b/src/main/native-chat/agent-session-journal/queued-message-schema.ts index 3f9d0122fda..d19c2e09a2d 100644 --- a/src/main/native-chat/agent-session-journal/queued-message-schema.ts +++ b/src/main/native-chat/agent-session-journal/queued-message-schema.ts @@ -12,7 +12,8 @@ const NULLABLE_COLUMNS: readonly (readonly [name: string, type: string])[] = [ ['consumed_as', 'TEXT'], ['carried_from', 'TEXT'], ['queued_epoch', 'TEXT'], - ['queued_sequence', 'INTEGER'] + ['queued_sequence', 'INTEGER'], + ['source_json', 'TEXT'] ] /** @@ -43,6 +44,7 @@ CREATE TABLE IF NOT EXISTS queued_messages ( carried_from TEXT, queued_epoch TEXT, queued_sequence INTEGER, + source_json TEXT, PRIMARY KEY (session_id, message_id) ); `) diff --git a/src/main/native-chat/agent-session-journal/queued-message-store.test.ts b/src/main/native-chat/agent-session-journal/queued-message-store.test.ts index abad56fa6dd..a9c0ef7c70c 100644 --- a/src/main/native-chat/agent-session-journal/queued-message-store.test.ts +++ b/src/main/native-chat/agent-session-journal/queued-message-store.test.ts @@ -22,6 +22,7 @@ import { QueuedMessageNotConsumableError } from './journal-queued-messages' import type { AgentSessionJournal } from './journal-store' +import type { AgentSessionMessageSource } from '../../../shared/agent-session-message-source' import { closeTestJournalHostDatabases, createTrackedJournalOpener @@ -81,7 +82,8 @@ async function queueDraft(journal: AgentSessionJournal, messageId: string, text messageId, body: message(text), fingerprint: `fp-${messageId}`, - hostInstance: 'proc-1' + hostInstance: 'proc-1', + source: { kind: 'user' } }) } @@ -174,6 +176,53 @@ describe('draft rows', () => { expect(journal.queuedMessages.list()).toHaveLength(1) }) + it("keeps who queued a card across reopen; a card with no readable sender is the person's", async () => { + const agent: AgentSessionMessageSource = { + kind: 'agent', + senders: [ + { + party: { + address: 'structworker_1', + terminalHandle: 'structworker_1', + orcaSessionId: null + } + } + ], + orchestration: { + message: 'mail-notice', + mailbox: 'run:r1', + dispatchId: 'd1', + messages: [{ messageId: 'm1', runId: 'r1', from: 'structworker_1' }] + } + } + const first = await open() + await first.queuedMessages.insert({ + messageId: 'agent-card', + body: message('You have 1 orchestration message. Run `orca orchestration check --run r1`.'), + fingerprint: 'fp-agent-card', + hostInstance: 'proc-1', + source: agent + }) + await queueDraft(first, 'before-the-column') + await queueDraft(first, 'unreadable') + await first.close() + closeTestJournalHostDatabases() + const db = new Database(journalDatabasePath(root)) + db.prepare('UPDATE queued_messages SET source_json = NULL WHERE message_id = ?').run( + 'before-the-column' + ) + db.prepare('UPDATE queued_messages SET source_json = \'{"v":9}\' WHERE message_id = ?').run( + 'unreadable' + ) + db.close() + const reopened = await open() + expect(reopened.queuedMessages.list().map((row) => [row.messageId, row.source])).toEqual([ + ['agent-card', agent], + ['before-the-column', { kind: 'user' }], + ['unreadable', { kind: 'user' }] + ]) + }) + it('drafts survive epoch replacement, which deletes only journal rows', async () => { const journal = await open() await queueDraft(journal, 'draft-1') diff --git a/src/main/native-chat/agent-session-journal/queued-message-table.ts b/src/main/native-chat/agent-session-journal/queued-message-table.ts index a2bd9c3e05a..1a82925f1f9 100644 --- a/src/main/native-chat/agent-session-journal/queued-message-table.ts +++ b/src/main/native-chat/agent-session-journal/queued-message-table.ts @@ -13,6 +13,11 @@ import type { AgentJournalCursor, AgentJournalMessageItem } from '../../../shared/agent-session-journal-types' +import { + readAgentSessionMessageSource, + serializeAgentSessionMessageSource, + type AgentSessionMessageSource +} from '../../../shared/agent-session-message-source' import { rejectedDraftSettlement } from './journal-dispatch-settlement' import { readStoredRejectionFact } from './journal-dispatch-reducer' @@ -59,10 +64,12 @@ export type QueuedMessageRow = { /** Where the journal stood when it was queued: a Stop's pause holds only cards queued before * it. Null on rows from builds before it was recorded, which read as queued before any Stop. */ queuedAt: AgentJournalCursor | null + /** Who it is from: the person, or another agent through Orca. */ + source: AgentSessionMessageSource } const COLUMNS = - 'session_id, message_id, position, body_json, fingerprint, created_at, host_instance, state, hold_reason, returned_reason, returned_rejection, settled_at, settled_by_op, consumed_as, carried_from, queued_epoch, queued_sequence' + 'session_id, message_id, position, body_json, fingerprint, created_at, host_instance, state, hold_reason, returned_reason, returned_rejection, settled_at, settled_by_op, consumed_as, carried_from, queued_epoch, queued_sequence, source_json' export function insertQueuedMessage( db: Database.Database, @@ -74,6 +81,7 @@ export function insertQueuedMessage( hostInstance: string carriedFrom?: string queuedAt: AgentJournalCursor + source: AgentSessionMessageSource now: number } ): QueuedMessageRow { @@ -84,7 +92,7 @@ export function insertQueuedMessage( const position = Number(highest?.p ?? 0) + 1 db.prepare( `INSERT INTO queued_messages (${COLUMNS}) - VALUES (?, ?, ?, ?, ?, ?, ?, 'waiting', NULL, NULL, NULL, NULL, NULL, NULL, ?, ?, ?)` + VALUES (?, ?, ?, ?, ?, ?, ?, 'waiting', NULL, NULL, NULL, NULL, NULL, NULL, ?, ?, ?, ?)` ).run( input.sessionId, input.messageId, @@ -95,7 +103,8 @@ export function insertQueuedMessage( input.hostInstance, input.carriedFrom ?? null, input.queuedAt.epoch, - input.queuedAt.sequence + input.queuedAt.sequence, + serializeAgentSessionMessageSource(input.source) ) return { sessionId: input.sessionId, @@ -113,7 +122,8 @@ export function insertQueuedMessage( settledByOp: null, consumedAs: null, carriedFrom: input.carriedFrom ?? null, - queuedAt: input.queuedAt + queuedAt: input.queuedAt, + source: input.source } } @@ -295,6 +305,7 @@ function toStoredRow(row: unknown): QueuedMessageRow | null { carried_from: string | null queued_epoch: string | null queued_sequence: number | null + source_json: string | null } let body: AgentJournalMessageItem try { @@ -333,10 +344,21 @@ function toStoredRow(row: unknown): QueuedMessageRow | null { queuedAt: record.queued_epoch !== null && typeof record.queued_sequence === 'number' ? { epoch: record.queued_epoch, sequence: record.queued_sequence } - : null + : null, + source: storedSource(record.source_json) } } +function storedSource(json: string | null): AgentSessionMessageSource { + let stored: unknown = null + try { + stored = json === null ? null : JSON.parse(json) + } catch { + // An unreadable value is read as no value; the source reader decides what that means. + } + return readAgentSessionMessageSource(stored) +} + function storedRejection(json: string | null): UnreadAgentSessionFailureFact | null { if (json === null) { return null diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts index 1133342270d..bceb95e4dbd 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts @@ -37,6 +37,7 @@ import { setOptionPlan } from './structured-agent-session-mutation-plans' import { runQueueableStructuredAgentSessionSend } from './structured-agent-session-queued-send' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' import { cancelStructuredAgentSessionPrompt } from './structured-agent-session-prompt-cancel' import { mutateWithChatStop } from './structured-agent-session-chat-stop' export type { StructuredAgentSessionMutationContext } from './structured-agent-session-mutation-context' @@ -60,6 +61,9 @@ export function sendStructuredAgentSessionTurn( * Orchestration mail, a restart continuation and `agent.launch`'s host-sent * prompt never set it. */ userSend?: true + /** Host-local, never on the wire: who a host-side `queue-if-active` send queues for, recorded + * on its card. A client's send is always its person's (`userSend`). */ + source?: AgentMessageSource beforeRun?: () => void }, arrival?: Parameters[2] diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-message-rig.test-fixture.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-message-rig.test-fixture.ts index f1271f4e339..00319720942 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-message-rig.test-fixture.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-message-rig.test-fixture.ts @@ -10,6 +10,7 @@ import { agentSessionFailureWords } from '../../../shared/agent-session-failure- import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' import type { AgentJournalSubmission } from '../../../shared/agent-session-journal-types' import type { AgentSessionQueuePause } from '../../../shared/agent-session-wire' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' import { openTestAgentSessionRecordStore } from '../../runtime/agent-session-record-store-test-harness' import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter' import type { StructuredAgentSessionEventSink } from './structured-agent-session-event-sink' @@ -30,6 +31,7 @@ import { createStructuredAgentSessionLogger } from './structured-agent-session-l import { codexProviderHandle } from '../../../shared/agent-session-provider-handle-encoding' export const QUEUED_RIG_CALLER = { callerKey: 'client-1' } +type RigSendOptions = { internal?: true; source?: AgentMessageSource } export function eventually(assertion: () => void | Promise): Promise { return vi.waitFor(assertion, { timeout: 10_000 }) @@ -140,16 +142,16 @@ export async function createQueuedMessageTestRig( } /** A client's send, as the `agentSession.send` RPC hands it to the host; - * `internal` is a host-side sender (orchestration mail, a restart continuation). */ - function send(text: string, delivery?: 'queue-if-active', options?: { internal?: true }) { + * `internal` is a host-side sender (orchestration mail, a restart continuation), and `source` + * who it is from. */ + function send(text: string, delivery?: 'queue-if-active', options?: RigSendOptions) { const body = hostTestMessage(text) const clientOperationId = hostTestOperationId() const fields = { body, ...(delivery ? { delivery } : {}) } const result = host.send(QUEUED_RIG_CALLER, { envelope: envelope(fields, 'agentSession.send', clientOperationId), - body, - ...(delivery ? { delivery } : {}), - ...(options?.internal ? {} : { userSend: true as const }) + ...fields, + ...(options?.internal ? { source: options.source } : { userSend: true as const }) }) return { id: clientOperationId, result } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.test.ts index 8a7c5c4c3d3..eefd106b3ea 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.test.ts @@ -12,6 +12,7 @@ import { type AgentSessionSubscribeEvent } from '../../../shared/agent-session-wire' import { ConversationCommandParams } from '../../../shared/rpc-contract/structured-agent-session-params' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' import { AgentSessionJournal } from '../agent-session-journal/journal-store' import { JournalQueuedMessages } from '../agent-session-journal/journal-queued-messages' import { @@ -636,6 +637,31 @@ describe('/clear', () => { }) }) + it('carries who each card is from', async () => { + const notice = { + message: 'mail-notice', + mailbox: 'run:r1', + dispatchId: null, + messages: [] + } as const + const source: AgentMessageSource = { kind: 'agent', senders: [], orchestration: notice } + const working = await workingSend() + await send('pointer', 'queue-if-active', { internal: true, source }).result + await send('typed', 'queue-if-active').result + await stop() + await settleAccepted(working, 'a') + const cleared = await clear(hostTestOperationId()) + const replacementId = cleared.ok ? cleared.value.replacementSessionId : undefined + if (!replacementId) { + throw new Error('expected a replacement session') + } + const journal = host.collaboratorsForTests().sessions.get(replacementId)?.journal + expect(journal?.queuedMessages.list().map((row) => row.source)).toEqual([ + source, + { kind: 'user' } + ]) + }) + it("the replacement's 'cleared' pause lifts through Resume exactly like a Stop's", async () => { const [firstId] = await pausedDrafts() const cleared = await clear(hostTestOperationId()) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.ts index 248534097b8..d4c7c470f12 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-messages.ts @@ -8,6 +8,10 @@ import { randomUUID } from 'node:crypto' import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types' +import { + USER_MESSAGE_SOURCE, + type AgentMessageSource +} from '../../../shared/agent-session-message-source' import { QUEUED_MESSAGE_PAUSED_SEND_FAILED, type AgentSessionSendResult, @@ -191,6 +195,10 @@ export async function maybeQueueStructuredAgentSessionSend( envelope: { clientOperationId: string } body: AgentJournalMessageItem delivery?: 'queue-if-active' + /** A person's send at a chat surface; it outranks any `source`. */ + userSend?: true + /** Who a host-side send is from. */ + source?: AgentMessageSource } ): Promise< | { ok: true; value: AgentSessionSendResult } @@ -231,7 +239,8 @@ export async function maybeQueueStructuredAgentSessionSend( messageId: clientMessageId, body: params.body, fingerprint: queuedMessageFingerprint(ctx.sessionId, params.body), - hostInstance: structuredAgentSessionHostInstance() + hostInstance: structuredAgentSessionHostInstance(), + source: params.userSend ? USER_MESSAGE_SOURCE : (params.source ?? USER_MESSAGE_SOURCE) }, ctx.operationReceipt ) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts index ee32b3648a9..8554bb08154 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-mutations.ts @@ -117,7 +117,8 @@ export async function carryQueuedMessagesToClearReplacement( body: row.body, fingerprint: queuedMessageFingerprint(input.replacementSessionId, row.body), hostInstance: structuredAgentSessionHostInstance(), - carriedFrom: ctx.sessionId + carriedFrom: ctx.sessionId, + source: row.source }) } await withdrawQueuedMessagesForOperation(ctx.journal, { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-send.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-send.ts index d411efc4239..9abdc713f0e 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-queued-send.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-queued-send.ts @@ -4,6 +4,7 @@ import type { AgentSessionSendResult } from '../../../shared/agent-session-wire' import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' import { maybeQueueStructuredAgentSessionSend } from './structured-agent-session-queued-messages' import type { AgentSessionTurnContext, TurnOutcome } from './structured-agent-session-turns' @@ -16,6 +17,7 @@ export async function runQueueableStructuredAgentSessionSend( body: AgentJournalMessageItem delivery?: 'queue-if-active' userSend?: true + source?: AgentMessageSource }, immediate: () => Promise> ): Promise> { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts index 20cd3b298d5..a1bb644d84e 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts @@ -519,7 +519,8 @@ describe('startup restore of chats still in their per-chat files', () => { messageId: 'draft-1', body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'later' }] }, fingerprint: 'fp-draft-1', - hostInstance: 'proc-1' + hostInstance: 'proc-1', + source: { kind: 'user' } }) expect(journal.importPending).toBe(false) diff --git a/src/main/runtime/orchestration/agent-facing-parity.test.ts b/src/main/runtime/orchestration/agent-facing-parity.test.ts index f682dfc7ce2..ac765825e69 100644 --- a/src/main/runtime/orchestration/agent-facing-parity.test.ts +++ b/src/main/runtime/orchestration/agent-facing-parity.test.ts @@ -113,7 +113,7 @@ async function renderPreamble(worker: 'chat' | 'terminal'): Promise { /** The turn text the structured lane sends a chat for one message on `mailbox`. */ async function renderChatPointer(mailbox: string): Promise { const texts: string[] = [] - db.insertMessage({ from: 'term_peer', to: mailbox, subject: 'hi' }) + const message = db.insertMessage({ from: 'term_peer', to: mailbox, subject: 'hi' }) const delivery = new OrchestrationStructuredMailboxPointerDelivery({ getDb: () => db, getMessageWaiters: () => undefined, @@ -121,7 +121,7 @@ async function renderChatPointer(mailbox: string): Promise { // The runtime's wiring of the structured lane. getCliCommand: localOrchestrationCliCommand, host: { - readGateFacts: async () => ({ turnRunning: false, awaitingHuman: false, submissions: [] }), + readSessionFacts: async () => ({ submissions: [] }), currentFence: () => 1, send: async (input) => { for (const block of input.body.blocks) { @@ -133,6 +133,8 @@ async function renderChatPointer(mailbox: string): Promise { }) delivery.deliverForHandle(mailbox) await vi.waitFor(() => expect(texts).toHaveLength(1)) + // Read, as the agent's `check` reads it, so this mailbox's next mail is pointed too. + db.markAsRead([message.id]) return texts[0]! } diff --git a/src/main/runtime/orchestration/orchestration-caller-identity.ts b/src/main/runtime/orchestration/orchestration-caller-identity.ts index 9e34634ff3e..6909ac7cf65 100644 --- a/src/main/runtime/orchestration/orchestration-caller-identity.ts +++ b/src/main/runtime/orchestration/orchestration-caller-identity.ts @@ -2,23 +2,10 @@ import type { RunRow } from './types' import { isEquivalentPaneKey } from './db/pane-key-match' import { currentRunCoordinatorOrcaSessionId } from './db/runs/run-coordinator-orca-session' import { formatOrcaSessionAddress, type OrcaSessionId } from '../../../shared/orca-session-address' +import type { OrchestrationPartyIdentity } from '../../../shared/orchestration-party-identity' -/** - * Who an orchestration caller is, as Run binding and mail routing match it. - * - * A PTY agent is its terminal: a handle and a pane key, no Orca session id. An agent that is a - * structured session is its Orca session id, addressed as `orca_session_id:`; a structured worker also - * has the handle and pane key it was minted, and an ordinary chat has neither. Methods pass this - * through whole and never branch on which fields are set; the lookups below own that. - */ -export type OrchestrationCallerIdentity = Readonly<{ - /** Mailbox address the caller sends from and reads: its terminal handle, else its session address. */ - address: string - terminalHandle: string | null - paneKey: string | null - /** The bare Orca session id the caller is addressed by; mail spells it `orca_session_id:`. */ - orcaSessionId: OrcaSessionId | null -}> +/** Who an orchestration caller is; shared so a queued message can name its sender the same way. */ +export type OrchestrationCallerIdentity = OrchestrationPartyIdentity /** The part of a caller a Run binding stores and matches. */ export type OrchestrationCoordinatorKey = Pick< diff --git a/src/main/runtime/orchestration/send-agent-turn-host.test.ts b/src/main/runtime/orchestration/send-agent-turn-host.test.ts index 5eb76ee1a55..0f92d95cee8 100644 --- a/src/main/runtime/orchestration/send-agent-turn-host.test.ts +++ b/src/main/runtime/orchestration/send-agent-turn-host.test.ts @@ -2,6 +2,7 @@ // host's own admission: a fingerprint over other fields than the send carries is refused there. import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' import { createQueuedMessageTestRig, eventually, @@ -20,6 +21,25 @@ import { let rig: QueuedMessageTestRig +const MAIL_SOURCE: AgentMessageSource = { + kind: 'agent', + senders: [ + { + party: { + address: 'term_peer', + terminalHandle: 'term_peer', + orcaSessionId: null + } + } + ], + orchestration: { + message: 'mail-notice', + mailbox: 'dispatch:d1', + dispatchId: 'd1', + messages: [{ messageId: 'm1', runId: 'r1', from: 'term_peer' }] + } +} + beforeEach(async () => { rig = await createQueuedMessageTestRig() }) @@ -36,7 +56,12 @@ function sendTurn( host, sessionId: SESSION, callerKey: 'trusted-local:orchestration:d1', - turn: { body: hostTestMessage('mail'), delivery, operationId, expectedRuntimeFence: 1 } + turn: { + body: hostTestMessage('mail'), + operationId, + expectedRuntimeFence: 1, + ...(delivery === 'queue' ? { delivery, source: MAIL_SOURCE } : { delivery }) + } }) } @@ -47,6 +72,15 @@ describe('sendAgentTurn through the real host', () => { kind: 'queued', queued: { position: 1, state: 'waiting' } }) + // Stored with the card, read back whole: who it is from survives the round trip. + expect( + rig.host + .collaboratorsForTests() + .sessions.get(SESSION) + ?.journal.queuedMessages.list() + .map(({ state, source }) => ({ state, source })) + ).toEqual([{ state: 'waiting', source: MAIL_SOURCE }]) + // Shown in the chat's queue like the person's own card. expect(await rig.drafts()).toMatchObject([{ state: 'waiting' }]) }) diff --git a/src/main/runtime/orchestration/send-agent-turn.test.ts b/src/main/runtime/orchestration/send-agent-turn.test.ts index c018f27d28b..599ef009ae5 100644 --- a/src/main/runtime/orchestration/send-agent-turn.test.ts +++ b/src/main/runtime/orchestration/send-agent-turn.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it, vi } from 'vitest' import type { AgentJournalSubmission } from '../../../shared/agent-session-journal-types' import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' import type { AgentSessionSendResult } from '../../../shared/agent-session-wire' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' import { ORCHESTRATION_READINESS_TIMEOUT_MS } from '../../../shared/orchestration-timing-budgets' import { dispatchPreambleSendOptions } from './preamble' import { @@ -52,6 +53,17 @@ function structuredHost(answer: HostSendAnswer, settled?: AgentJournalSubmission return { host, send, waitForSendSettlement } } +const MAIL_SOURCE: AgentMessageSource = { + kind: 'agent', + senders: [], + orchestration: { + message: 'mail-notice', + mailbox: 'dispatch:d1', + dispatchId: 'd1', + messages: [{ messageId: 'm1', runId: 'r1', from: 'term_peer' }] + } +} + const turn: StructuredSessionTurn = { body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'hello' }] }, delivery: 'now', @@ -142,7 +154,7 @@ describe('sendAgentTurn to a structured session', () => { }) ) await expect( - sendAgentTurn(structured(fake.host, { ...turn, delivery: 'queue' })) + sendAgentTurn(structured(fake.host, { ...turn, delivery: 'queue', source: MAIL_SOURCE })) ).resolves.toEqual({ kind: 'queued', clientMessageId: 'op-1', @@ -158,7 +170,9 @@ describe('sendAgentTurn to a structured session', () => { payloadFingerprint: hostFingerprint({ body: turn.body, delivery: 'queue-if-active' }) }, body: turn.body, - delivery: 'queue-if-active' + delivery: 'queue-if-active', + // Host-local: who the card is from rides beside the envelope, outside its fingerprint. + source: MAIL_SOURCE } ) expect(fake.waitForSendSettlement).not.toHaveBeenCalled() @@ -172,7 +186,7 @@ describe('sendAgentTurn to a structured session', () => { }) ) await expect( - sendAgentTurn(structured(fake.host, { ...turn, delivery: 'queue' })) + sendAgentTurn(structured(fake.host, { ...turn, delivery: 'queue', source: MAIL_SOURCE })) ).resolves.toMatchObject({ kind: 'queued', queued: { state: 'returned' } }) }) }) diff --git a/src/main/runtime/orchestration/send-agent-turn.ts b/src/main/runtime/orchestration/send-agent-turn.ts index 0a1921b96ba..340bcc82e87 100644 --- a/src/main/runtime/orchestration/send-agent-turn.ts +++ b/src/main/runtime/orchestration/send-agent-turn.ts @@ -17,6 +17,7 @@ import { type AgentSessionQueuedSendReceipt } from '../../../shared/agent-session-wire' import type { AgentSessionWireRefusal } from '../../../shared/agent-session-wire-refusals' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' import { ORCHESTRATION_READINESS_TIMEOUT_MS } from '../../../shared/orchestration-timing-budgets' import { structuredAgentSessionMessageSendMutation } from '../../../shared/structured-agent-session-send-mutation' import type { StructuredAgentSessionHost } from '../../native-chat/agent-session-wire/structured-agent-session-host' @@ -36,11 +37,14 @@ export type StructuredAgentTurnHost = Pick< export type StructuredSessionTurn = { body: AgentJournalMessageItem - delivery: AgentTurnDelivery /** Reused on a retry, so the host replays its recorded answer instead of sending twice. */ operationId: string expectedRuntimeFence: number -} +} & ( + | { delivery: 'now' } + /** A queued card records who it is from. */ + | { delivery: 'queue'; source: AgentMessageSource } +) export type StructuredSessionTurnSend = { kind: 'structured-session' @@ -123,15 +127,17 @@ async function sendStructuredSessionTurn( send: StructuredSessionTurnSend ): Promise { const { turn } = send + const message = structuredAgentSessionMessageSendMutation({ + sessionId: send.sessionId, + clientOperationId: turn.operationId, + expectedRuntimeFence: turn.expectedRuntimeFence, + body: turn.body, + delivery: turn.delivery === 'queue' ? 'queue-if-active' : undefined + }) const result = await send.host.send( { callerKey: send.callerKey }, - structuredAgentSessionMessageSendMutation({ - sessionId: send.sessionId, - clientOperationId: turn.operationId, - expectedRuntimeFence: turn.expectedRuntimeFence, - body: turn.body, - delivery: turn.delivery === 'queue' ? 'queue-if-active' : undefined - }) + // The source is host-local and outside the fingerprint: a retry under the same id replays. + turn.delivery === 'queue' ? { ...message, source: turn.source } : message ) if (!result.ok) { return { kind: 'refused', refusal: result.refusal } diff --git a/src/main/runtime/orchestration/structured-mail-source.test.ts b/src/main/runtime/orchestration/structured-mail-source.test.ts new file mode 100644 index 00000000000..f6f394d722e --- /dev/null +++ b/src/main/runtime/orchestration/structured-mail-source.test.ts @@ -0,0 +1,43 @@ +import { describe, expect, it } from 'vitest' +import { structuredMailSource } from './structured-mail-source' + +const SESSION = '4a1f6c2e-8b3d-4e7a-9c15-0d2b6e8f1a37' + +describe('who delivered mail is from', () => { + it("names each sender once, without the pane key that would open its mailbox, and each message's own sender", () => { + // No database: a terminal handle and a session address name their party by themselves. + const source = structuredMailSource({ + db: null, + mailboxHandle: 'run:r1', + dispatchId: null, + batch: [ + { id: 'm1', from_handle: 'term_a', run_id: 'r1' }, + { id: 'm2', from_handle: `orca_session_id:${SESSION}`, run_id: 'r2' }, + { id: 'm3', from_handle: 'term_a', run_id: 'r1' } + ] + }) + expect(source).toEqual({ + kind: 'agent', + senders: [ + { party: { address: 'term_a', terminalHandle: 'term_a', orcaSessionId: null } }, + { + party: { + address: `orca_session_id:${SESSION}`, + terminalHandle: null, + orcaSessionId: SESSION + } + } + ], + orchestration: { + message: 'mail-notice', + mailbox: 'run:r1', + dispatchId: null, + messages: [ + { messageId: 'm1', runId: 'r1', from: 'term_a' }, + { messageId: 'm2', runId: 'r2', from: `orca_session_id:${SESSION}` }, + { messageId: 'm3', runId: 'r1', from: 'term_a' } + ] + } + }) + }) +}) diff --git a/src/main/runtime/orchestration/structured-mail-source.ts b/src/main/runtime/orchestration/structured-mail-source.ts new file mode 100644 index 00000000000..ad26c3b8be0 --- /dev/null +++ b/src/main/runtime/orchestration/structured-mail-source.ts @@ -0,0 +1,53 @@ +/** + * Who the mail a chat is pointed at is from: every distinct sender, named the + * way orchestration names a party, and each message's own sender and records. The run, dispatch + * and message ids join back to orchestration's own rows while those exist. + */ + +import { parseOrcaSessionAddress } from '../../../shared/orca-session-address' +import type { + AgentMessageSource, + AgentMessageSender +} from '../../../shared/agent-session-message-source' +import type { MessageRow, OrchestrationDb } from './db' +import { resolveOrchestrationParty } from './orchestration-party' + +export type MailSourceMessage = Pick + +export function structuredMailSource(input: { + db: OrchestrationDb | null + mailboxHandle: string + dispatchId: string | null + batch: readonly MailSourceMessage[] +}): AgentMessageSource { + const senders = new Map() + for (const { from_handle: address } of input.batch) { + if (!senders.has(address)) { + senders.set(address, { party: senderParty(address, input.db) }) + } + } + return { + kind: 'agent', + senders: [...senders.values()], + orchestration: { + message: 'mail-notice', + mailbox: input.mailboxHandle, + dispatchId: input.dispatchId, + messages: input.batch.map((message) => ({ + messageId: message.id, + runId: message.run_id, + from: message.from_handle + })) + } + } +} + +function senderParty(address: string, db: OrchestrationDb | null): AgentMessageSender['party'] { + try { + const { paneKey: _credential, ...party } = resolveOrchestrationParty(address, db) + return party + } catch { + // A worker this host lost the identity of: what the address itself says. + return { address, terminalHandle: null, orcaSessionId: parseOrcaSessionAddress(address) } + } +} diff --git a/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.test.ts b/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.test.ts index 9b40bc08b0d..46090094c0d 100644 --- a/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.test.ts +++ b/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.test.ts @@ -1,5 +1,4 @@ import { describe, expect, it, vi } from 'vitest' -import type { AgentJournalRenderItem } from '../../../shared/agent-session-journal-types' import { OrchestrationStructuredMailboxPointerDelivery, type StructuredMailboxPointerHost @@ -9,7 +8,6 @@ import { structuredPointerBatchFingerprint, type StructuredPointerSubmission } from './structured-pointer-operation-id' -import { structuredSessionGateFacts } from './structured-session-pointer-delivery' import type { StructuredWorkerIdentity } from '../structured-worker-identity' const IDENTITY: StructuredWorkerIdentity = { @@ -22,63 +20,12 @@ const IDENTITY: StructuredWorkerIdentity = { hostScope: { kind: 'local', hostId: 'local' } } -function idleJournal(): AgentJournalRenderItem[] { - return [ - { - itemId: 'i1', - observedAt: 1, - body: { kind: 'status', text: 'done', turnLifecycle: { state: 'completed', turnId: 't1' } } - } as unknown as AgentJournalRenderItem - ] -} - -function runningJournal(): AgentJournalRenderItem[] { - return [ - { - itemId: 'i1', - observedAt: 1, - body: { kind: 'status', text: 'working', turnLifecycle: { state: 'running', turnId: 't1' } } - } as unknown as AgentJournalRenderItem - ] -} - -/** What a worker's journal looks like once it has finished a substantial turn: history, and no - * turnLifecycle row anywhere, because settlement tombstones it. */ -function settledLongJournal(): AgentJournalRenderItem[] { - return Array.from( - { length: 120 }, - (_unused, index) => - ({ - itemId: `tool-${index}`, - observedAt: index, - body: { kind: 'tool-call', name: 'Bash', input: {}, state: 'completed' } - }) as unknown as AgentJournalRenderItem - ) -} - -/** A prompt raised at the very start of a long turn, far outside any bounded tail window. */ -function staleAttentionJournal(): AgentJournalRenderItem[] { - return [...attentionJournal(), ...settledLongJournal()] -} - -function attentionJournal(): AgentJournalRenderItem[] { - return [ - { - itemId: 'i1', - observedAt: 1, - body: { - kind: 'question', - question: 'which?', - options: [], - resolution: { state: 'pending' } - } - } as unknown as AgentJournalRenderItem - ] -} - function harness(options: { - journal: AgentJournalRenderItem[] | null + /** False: the session cannot be read (not attached). */ + attached?: boolean dispatchState?: 'accepted' | 'rejected' | 'unknown' + /** The chat was busy: its queue holds the pointer as a card. */ + queued?: true /** The coordinator of this worker's Run is mid-batch: it checked and has not acked yet. */ outstandingRunDelivery?: boolean outstandingOwnDelivery?: boolean @@ -88,14 +35,24 @@ function harness(options: { }) { const mailbox = options.mailbox ?? 'dispatch:d1' const dispatchId = options.dispatchId === undefined ? 'd1' : options.dispatchId - let journal = options.journal + let attached = options.attached ?? true // The session's recorded sends, as its journal reports them. let submissions: StructuredPointerSubmission[] = [] - const markAsDelivered = vi.fn() - const send: StructuredMailboxPointerHost['send'] = vi.fn(async () => ({ - kind: 'sent' as const, - state: options.dispatchState ?? ('accepted' as const) - })) + // The mailbox's unread mail; a pointed message is no longer selected for a pointer. + const mail = [ + { id: 'm1', type: 'status', sequence: 3, from_handle: 'term_coord', run_id: 'run_1' } + ] + const pointed = new Set() + const markAsDelivered = vi.fn((ids: string[]) => { + for (const id of ids) { + pointed.add(id) + } + }) + const send: StructuredMailboxPointerHost['send'] = vi.fn(async () => + options.queued + ? { kind: 'queued' as const } + : { kind: 'sent' as const, state: options.dispatchState ?? ('accepted' as const) } + ) const sendMock = vi.mocked(send) const stored = new Map() const db = { @@ -103,7 +60,7 @@ function harness(options: { hasOutstandingMailboxDelivery: (handle: string) => ((options.outstandingRunDelivery ?? false) && handle.startsWith('run:')) || ((options.outstandingOwnDelivery ?? false) && !handle.startsWith('run:')), - getUndeliveredUnreadMessages: () => [{ id: 'm1', type: 'status', sequence: 3 }], + getUndeliveredUnreadMessages: () => mail.filter((message) => !pointed.has(message.id)), markAsDelivered, getStructuredPointerOperation: (key: string) => stored.get(key), putStructuredPointerOperation: (row: StructuredPointerOperationRow) => @@ -117,8 +74,7 @@ function harness(options: { mailboxHandle === mailbox ? { sessionId: IDENTITY.sessionId, dispatchId } : null, getCliCommand: () => 'orca-dev', host: { - readGateFacts: async () => - journal === null ? null : { ...structuredSessionGateFacts(journal), submissions }, + readSessionFacts: async () => (attached ? { submissions } : null), currentFence: () => 4, send } @@ -128,8 +84,11 @@ function harness(options: { markAsDelivered, send: sendMock, stored, - setJournal: (next: AgentJournalRenderItem[] | null) => { - journal = next + setAttached: (next: boolean) => { + attached = next + }, + receive: (id: string, sequence: number) => { + mail.push({ id, type: 'status', sequence, from_handle: 'term_coord', run_id: 'run_1' }) }, setSubmissions: (next: StructuredPointerSubmission[]) => { submissions = next @@ -141,13 +100,13 @@ const flush = () => new Promise((resolve) => setTimeout(resolve, 0)) describe('structured mailbox pointer delivery', () => { it('claims only mailboxes whose assignee is a structured worker', () => { - const { delivery } = harness({ journal: idleJournal() }) + const { delivery } = harness({}) expect(delivery.deliverForHandle('dispatch:d1')).toBe(true) expect(delivery.deliverForHandle('run:run_1')).toBe(false) }) it('sends the pointer as a turn and consumes mail on an accepted dispatch', async () => { - const { delivery, markAsDelivered, send } = harness({ journal: idleJournal() }) + const { delivery, markAsDelivered, send } = harness({}) delivery.deliverForHandle('dispatch:d1') await flush() expect(send).toHaveBeenCalledTimes(1) @@ -157,7 +116,6 @@ describe('structured mailbox pointer delivery', () => { it('nudges through the worker`s own handle for direct peer mail outside a dispatch', async () => { const { delivery, send, markAsDelivered } = harness({ - journal: idleJournal(), mailbox: IDENTITY.handle, dispatchId: null }) @@ -176,7 +134,6 @@ describe('structured mailbox pointer delivery', () => { it('retains mail when the dispatch settles unknown', async () => { const { delivery, markAsDelivered } = harness({ - journal: idleJournal(), dispatchState: 'unknown' }) delivery.deliverForHandle('dispatch:d1') @@ -184,40 +141,8 @@ describe('structured mailbox pointer delivery', () => { expect(markAsDelivered).not.toHaveBeenCalled() }) - it('retains mail while a turn is running', async () => { - const { delivery, send, markAsDelivered } = harness({ journal: runningJournal() }) - delivery.deliverForHandle('dispatch:d1') - await flush() - expect(send).not.toHaveBeenCalled() - expect(markAsDelivered).not.toHaveBeenCalled() - }) - - it('retains mail while a prompt is waiting for a human', async () => { - const { delivery, send } = harness({ journal: attentionJournal() }) - delivery.deliverForHandle('dispatch:d1') - await flush() - expect(send).not.toHaveBeenCalled() - }) - - it('delivers to a worker whose finished turn left a long history and no lifecycle row', async () => { - // The steady state after a worker's first substantial turn. Gating on a bounded tail page read - // this as permanently busy, so every later nudge parked forever and the worker went unnudged. - const { delivery, send, markAsDelivered } = harness({ journal: settledLongJournal() }) - delivery.deliverForHandle('dispatch:d1') - await flush() - expect(send).toHaveBeenCalledTimes(1) - expect(markAsDelivered).toHaveBeenCalledWith(['m1']) - }) - - it('retains mail for a prompt that scrolled out of the tail window', async () => { - const { delivery, send } = harness({ journal: staleAttentionJournal() }) - delivery.deliverForHandle('dispatch:d1') - await flush() - expect(send).not.toHaveBeenCalled() - }) - it('retains mail when the session is not attached', async () => { - const { delivery, send } = harness({ journal: null }) + const { delivery, send } = harness({ attached: false }) delivery.deliverForHandle('dispatch:d1') await flush() expect(send).not.toHaveBeenCalled() @@ -226,27 +151,54 @@ describe('structured mailbox pointer delivery', () => { it('redrives a detached session when the journal replays on re-attach', async () => { // A transient detach parks nothing to be woken unless `session-not-attached` waits for the // journal edge, and the dispatch preamble tells the worker not to poll. - const { delivery, send, setJournal, markAsDelivered } = harness({ journal: null }) + const { delivery, send, setAttached, markAsDelivered } = harness({ attached: false }) delivery.deliverForHandle('dispatch:d1') await flush() expect(send).not.toHaveBeenCalled() - setJournal(idleJournal()) + setAttached(true) delivery.onJournalActivity('session-1') await flush() expect(send).toHaveBeenCalledTimes(1) expect(markAsDelivered).toHaveBeenCalledWith(['m1']) }) - it('retries a parked pointer when the journal moves', async () => { - const { delivery, send, setJournal, markAsDelivered } = harness({ journal: runningJournal() }) + it('sends the pointer through the chat, with who it is from, and counts it pointed once queued', async () => { + const { delivery, send, markAsDelivered, stored } = harness({ queued: true }) delivery.deliverForHandle('dispatch:d1') await flush() - expect(send).not.toHaveBeenCalled() - setJournal(idleJournal()) - delivery.onJournalActivity('session-1') + expect(send).toHaveBeenCalledTimes(1) + expect(send.mock.calls[0]![0].body.blocks[0]).toMatchObject({ + text: expect.stringContaining('orca-dev orchestration check') + }) + expect(send.mock.calls[0]![0].source).toMatchObject({ + kind: 'agent', + senders: [{ party: { address: 'term_coord' } }], + orchestration: { + message: 'mail-notice', + mailbox: 'dispatch:d1', + messages: [{ messageId: 'm1', runId: 'run_1', from: 'term_coord' }] + } + }) + // The chat's queue holds it now, as it holds the person's: the same mail is not pointed again. + expect(markAsDelivered).toHaveBeenCalledWith(['m1']) + expect(stored.has('dispatch:d1')).toBe(false) + }) + + it('points mail that arrives while earlier pointed mail is still unread, counting only the new mail', async () => { + const { delivery, send, receive } = harness({}) + delivery.deliverForHandle('dispatch:d1') await flush() expect(send).toHaveBeenCalledTimes(1) - expect(markAsDelivered).toHaveBeenCalledWith(['m1']) + delivery.deliverForHandle('dispatch:d1') + await flush() + expect(send).toHaveBeenCalledTimes(1) + receive('m2', 4) + delivery.deliverForHandle('dispatch:d1') + await flush() + expect(send).toHaveBeenCalledTimes(2) + expect(send.mock.calls[1]![0].body.blocks[0]).toMatchObject({ + text: expect.stringContaining('You have 1 orchestration message.') + }) }) it('nudges the worker while its coordinator holds an unacked Run delivery', async () => { @@ -255,7 +207,6 @@ describe('structured mailbox pointer delivery', () => { // coordinator's `run:` delivery is invisible here — gating the WORKER's dispatch mailbox on it // dropped the nudge with nothing parked, and the worker sat idle on mail it was never told of. const { delivery, send, markAsDelivered } = harness({ - journal: idleJournal(), outstandingRunDelivery: true }) delivery.deliverForHandle('dispatch:d1') @@ -267,7 +218,7 @@ describe('structured mailbox pointer delivery', () => { 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 }) + const { delivery, send } = harness({ outstandingOwnDelivery: true }) delivery.deliverForHandle('dispatch:d1') await flush() expect(send).not.toHaveBeenCalled() @@ -278,7 +229,6 @@ describe('structured mailbox pointer delivery', () => { // stranded the worker until unrelated mail happened to arrive. The retry keeps the id: the host // replays a recorded refusal rather than starting the agent again. const { delivery, send, markAsDelivered } = harness({ - journal: idleJournal(), dispatchState: 'rejected' }) delivery.deliverForHandle('dispatch:d1') @@ -294,7 +244,6 @@ describe('structured mailbox pointer delivery', () => { it('points again under a new id once a later send ran', async () => { const { delivery, send, setSubmissions } = harness({ - journal: idleJournal(), dispatchState: 'unknown' }) delivery.deliverForHandle('dispatch:d1') @@ -312,7 +261,6 @@ describe('structured mailbox pointer delivery', () => { it('points once more under a new id for a send an earlier process left in doubt', async () => { const { delivery, send, stored, setSubmissions } = harness({ - journal: idleJournal(), dispatchState: 'unknown' }) stored.set('dispatch:d1', { @@ -344,7 +292,6 @@ describe('structured mailbox pointer delivery', () => { vi.useFakeTimers({ toFake: ['Date'] }) try { const { delivery, send, setSubmissions } = harness({ - journal: idleJournal(), dispatchState: 'unknown' }) // The wall clock steps back an hour after the lane started: its own row is still its own. @@ -374,7 +321,6 @@ describe('structured mailbox pointer delivery', () => { vi.useFakeTimers({ toFake: ['Date'] }) try { const { delivery, send, setSubmissions } = harness({ - journal: idleJournal(), dispatchState: 'unknown' }) const personTurn = { @@ -408,7 +354,6 @@ describe('structured mailbox pointer delivery', () => { it('stamps a pointer whose echo arrived after the lane stopped waiting, sending nothing more', async () => { const { delivery, send, markAsDelivered, stored, setSubmissions } = harness({ - journal: idleJournal(), dispatchState: 'unknown' }) delivery.deliverForHandle('dispatch:d1') @@ -428,7 +373,6 @@ describe('structured mailbox pointer delivery', () => { it('reuses one operation id for the same batch and re-mints when it grows', async () => { const { delivery, send, stored } = harness({ - journal: idleJournal(), dispatchState: 'unknown' }) delivery.deliverForHandle('dispatch:d1') @@ -445,10 +389,10 @@ describe('structured mailbox pointer delivery', () => { }) describe('forgetting one settled worker', () => { - /** Two workers, each mid-turn and so each parked on its OWN session's journal edge. */ + /** Two workers, each detached and so each parked on its OWN session's journal edge. */ function twoWorkerHarness() { let resolves = true - let journal = runningJournal() + let attached = false const sessionByMailbox: Record = { 'dispatch:d1': 'session-1', 'dispatch:d2': 'session-2' @@ -460,7 +404,9 @@ describe('forgetting one settled worker', () => { const db = { getDispatchContextById: () => ({ run_id: 'run_1' }), hasOutstandingMailboxDelivery: () => false, - getUndeliveredUnreadMessages: () => [{ id: 'm1', type: 'status', sequence: 3 }], + getUndeliveredUnreadMessages: () => [ + { id: 'm1', type: 'status', sequence: 3, from_handle: 'term_coord', run_id: 'run_1' } + ], markAsDelivered: vi.fn(), getStructuredPointerOperation: () => undefined, putStructuredPointerOperation: () => {}, @@ -477,7 +423,7 @@ describe('forgetting one settled worker', () => { }, getCliCommand: () => 'orca', host: { - readGateFacts: async () => ({ ...structuredSessionGateFacts(journal), submissions: [] }), + readSessionFacts: async () => (attached ? { submissions: [] } : null), currentFence: () => 4, send } @@ -485,8 +431,8 @@ describe('forgetting one settled worker', () => { return { delivery, send: vi.mocked(send), - goIdle: () => { - journal = idleJournal() + attach: () => { + attached = true }, stopResolving: () => { resolves = false @@ -501,7 +447,7 @@ describe('forgetting one settled worker', () => { // The bug: `forgetSession` re-resolved every parked mailbox and pruned the ones that answered // null. A momentarily null DB reference or a session mid-teardown made that EVERY worker, so // the sibling's mail stayed durable but lost the edge that would have woken it. - const { delivery, send, goIdle, stopResolving, resumeResolving } = twoWorkerHarness() + const { delivery, send, attach, stopResolving, resumeResolving } = twoWorkerHarness() delivery.deliverForHandle('dispatch:d1') delivery.deliverForHandle('dispatch:d2') await flush() @@ -511,7 +457,7 @@ describe('forgetting one settled worker', () => { delivery.forgetSession('session-1') resumeResolving() - goIdle() + attach() delivery.onJournalActivity('session-2') await flush() expect(send).toHaveBeenCalledTimes(1) @@ -519,7 +465,7 @@ describe('forgetting one settled worker', () => { }) it('still drops what the settled worker itself had parked', async () => { - const { delivery, send, goIdle, stopResolving } = twoWorkerHarness() + const { delivery, send, attach, stopResolving } = twoWorkerHarness() delivery.deliverForHandle('dispatch:d1') await flush() expect(send).not.toHaveBeenCalled() @@ -529,7 +475,7 @@ describe('forgetting one settled worker', () => { stopResolving() delivery.forgetSession('session-1') - goIdle() + attach() delivery.onJournalActivity('session-1') await flush() expect(send).not.toHaveBeenCalled() diff --git a/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts b/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts index e84aa8e30e9..f442c520d4b 100644 --- a/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts +++ b/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts @@ -3,9 +3,9 @@ * * The PTY lane types the nudge into a live pane and reads the idle edge off the terminal title. * Neither exists here, so this is a sibling of `OrchestrationMailboxPointerDelivery` rather than a - * branch inside it: batch selection is literally shared (`selectOrchestrationPointerBatch`), and - * everything below it is different — the nudge is a session turn, the idle edge is the journal, - * and only an `accepted` dispatch may consume mail. + * branch inside it: batch selection and the pointer text are literally shared, and everything + * below it is different — the nudge goes through the chat's own send, as a person's message does, + * and the retry edge is the journal. * * Coordinators are in scope here, unlike the PTY lane's reasoning: a PTY coordinator blocks in * `check --wait`, where a waiter preempts pointer delivery, but a structured coordinator is a chat @@ -13,7 +13,8 @@ */ import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types' -import type { OrchestrationDb } from './db' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' +import type { MessageRow, OrchestrationDb } from './db' import { formatMessagePointer } from './formatter' import type { OrchestrationCliCommand } from './cli-command' import { @@ -24,13 +25,12 @@ import { resolveStructuredPointerOperation, type StructuredPointerSubmission } from './structured-pointer-operation-id' +import { structuredMailSource } from './structured-mail-source' import { - decideStructuredSessionPointerDelivery, retainReasonForDispatch, structuredDispatchDelivered, type StructuredDispatchState, - type StructuredPointerRetainReason, - type StructuredSessionGateFacts + type StructuredPointerRetainReason } from './structured-session-pointer-delivery' export type StructuredPointerTarget = { @@ -50,22 +50,25 @@ type ParkedPointerDelivery = { export type StructuredPointerSendOutcome = | { kind: 'sent'; state: StructuredDispatchState } + /** The chat's queue took it, as it takes a person's message. */ + | { kind: 'queued' } | { kind: 'unattached' } -export type StructuredPointerGateFacts = StructuredSessionGateFacts & { +export type StructuredPointerSessionFacts = { /** Every send the session recorded, oldest first: what the lane's own sends settled as. */ submissions: readonly StructuredPointerSubmission[] } export type StructuredMailboxPointerHost = { - /** The idle gate, read off the session's full reduced timeline; `null` when it cannot be read. */ - readGateFacts: (sessionId: string) => Promise + /** `null` when the session cannot be read. */ + readSessionFacts: (sessionId: string) => Promise send: (input: { sessionId: string dispatchId: string | null operationId: string expectedRuntimeFence: number body: AgentJournalMessageItem + source: AgentMessageSource }) => Promise /** Current lease fence; `null` when no record backs the session any more. */ currentFence: (sessionId: string) => number | null @@ -194,14 +197,13 @@ export class OrchestrationStructuredMailboxPointerDelivery< db: OrchestrationDb, mailboxHandle: string, target: StructuredPointerTarget, - unread: readonly { id: string; type: string; sequence: number }[], + unread: readonly MessageRow[], reservedTypes: ReadonlySet | undefined ): Promise { const sessionId = target.sessionId - const session = await this.deps.host.readGateFacts(sessionId) - const decision = decideStructuredSessionPointerDelivery({ session }) - if (!decision.deliver) { - this.retain(mailboxHandle, sessionId, decision.retain, reservedTypes) + const session = await this.deps.host.readSessionFacts(sessionId) + if (!session) { + this.retain(mailboxHandle, sessionId, 'session-not-attached', reservedTypes) return } const fence = this.deps.host.currentFence(sessionId) @@ -225,7 +227,7 @@ export class OrchestrationStructuredMailboxPointerDelivery< mailboxHandle, sessionId, messageIds: staged, - submissions: session?.submissions ?? [], + submissions: session.submissions, sentByThisProcess: this.sentOperationIds.get(mailboxHandle) }) if (operation.kind === 'stamp') { @@ -245,13 +247,20 @@ export class OrchestrationStructuredMailboxPointerDelivery< dispatchId: target.dispatchId, operationId: operation.operationId, expectedRuntimeFence: fence, - body + body, + source: structuredMailSource({ + db, + mailboxHandle, + dispatchId: target.dispatchId, + batch: unread + }) }) if (outcome.kind === 'unattached') { this.retain(mailboxHandle, sessionId, 'session-not-attached', reservedTypes) return } - if (!structuredDispatchDelivered(outcome.state)) { + // A queued pointer is the chat's queue's to send, as a person's queued message is. + if (outcome.kind === 'sent' && !structuredDispatchDelivered(outcome.state)) { // The row stays: resending under its id replays this verdict and starts nothing. this.retain(mailboxHandle, sessionId, retainReasonForDispatch(outcome.state), reservedTypes) return @@ -264,7 +273,8 @@ export class OrchestrationStructuredMailboxPointerDelivery< } /** - * No `markAsUndelivered` is owed: rows are marked delivered only after an accepted dispatch. + * No `markAsUndelivered` is owed: rows are marked delivered only after an accepted dispatch, or + * once the chat's queue holds the pointer. * * Every reason parks for the session's next journal edge. `unknown` may mean the nudge already * sits in the provider's input queue, so an immediate retry can stack duplicate nudges; diff --git a/src/main/runtime/orchestration/structured-mailbox-pointer-host.test.ts b/src/main/runtime/orchestration/structured-mailbox-pointer-host.test.ts index 72953f1e904..1d1700e76bc 100644 --- a/src/main/runtime/orchestration/structured-mailbox-pointer-host.test.ts +++ b/src/main/runtime/orchestration/structured-mailbox-pointer-host.test.ts @@ -1,5 +1,6 @@ import { beforeEach, describe, expect, it, vi } from 'vitest' import type { AgentJournalRenderItem } from '../../../shared/agent-session-journal-types' +import type { AgentMessageSource } from '../../../shared/agent-session-message-source' const hostRef: { current: unknown } = { current: null } @@ -9,6 +10,7 @@ vi.mock('../../native-chat/agent-session-wire/structured-agent-session-registry' const { createStructuredMailboxPointerHost, + readStructuredSessionGateFacts, structuredPointerCallerKey, structuredSessionPointerCallerKey } = await import('./structured-mailbox-pointer-host') @@ -33,6 +35,12 @@ function transcript(count: number): AgentJournalRenderItem[] { ) } +const NOTICE_SOURCE: AgentMessageSource = { + kind: 'agent', + senders: [], + orchestration: { message: 'mail-notice', mailbox: 'dispatch:d1', dispatchId: 'd1', messages: [] } +} + describe('structured mailbox pointer host', () => { beforeEach(() => { hostRef.current = null @@ -41,30 +49,33 @@ describe('structured mailbox pointer host', () => { it('reads the gate facts from the FULL timeline, never a bounded tail', async () => { // The defect this pins: a running turn is announced by ONE lifecycle item, and settlement // tombstones it rather than rewriting it. A long tool-calling turn pushes that item arbitrarily - // far from the tail, so any page-sized read reports a busy worker as idle — and the pointer is - // then delivered mid-turn, which Codex coalesces into the running turn and Claude folds into - // it -- either way folded into work already in flight rather than read as a new instruction. + // far from the tail, so any page-sized read reports a busy worker as idle — and `@idle` then + // wakes it mid-turn. const items = [runningTurn(), ...transcript(500)] - const submissions = [{ clientMessageId: 'op1', dispatchState: 'unknown' }] - hostRef.current = { journalSnapshot: () => ({ items, submissions }) } - // The recorded sends ride along: the lane reads what its own operation id settled as. - expect(await createStructuredMailboxPointerHost().readGateFacts('s1')).toEqual({ + hostRef.current = { journalSnapshot: () => ({ items, submissions: [] }) } + expect(await readStructuredSessionGateFacts('s1')).toEqual({ turnRunning: true, - awaitingHuman: false, + awaitingHuman: false + }) + }) + + it("reads what the session's sends settled as", async () => { + const submissions = [{ clientMessageId: 'op1', dispatchState: 'unknown' }] + hostRef.current = { journalSnapshot: () => ({ items: [], submissions }) } + expect(await createStructuredMailboxPointerHost().readSessionFacts('s1')).toEqual({ submissions }) }) - it('answers null rather than idle when the session cannot be read', async () => { - // Null retains the pointer; `{turnRunning:false}` would deliver a nudge into a session this - // runtime cannot see at all. - expect(await createStructuredMailboxPointerHost().readGateFacts('s1')).toBeNull() + it('answers null rather than nothing recorded when the session cannot be read', async () => { + // Null retains the pointer; an empty answer would send into a session this runtime cannot see. + expect(await createStructuredMailboxPointerHost().readSessionFacts('s1')).toBeNull() hostRef.current = { journalSnapshot: () => { throw new Error('agent_session_ownership_unknown') } } - expect(await createStructuredMailboxPointerHost().readGateFacts('s1')).toBeNull() + expect(await createStructuredMailboxPointerHost().readSessionFacts('s1')).toBeNull() }) it('reports an unattached host rather than a rejection when nothing can be sent', async () => { @@ -108,25 +119,31 @@ describe('structured mailbox pointer host', () => { expect(send.mock.calls[0]![1]!.retryUnknown).toBeUndefined() }) - it('reads a queued answer as unknown, so the pointer is retained', async () => { - hostRef.current = { - send: async () => ({ + it('asks a busy chat to queue the pointer as a card, with who it is from', async () => { + const send = vi.fn( + async (_caller: unknown, _payload: { delivery?: string; source?: unknown }) => ({ ok: true, value: { clientMessageId: 'op1', queued: { messageId: 'op1', position: 0, state: 'waiting' } } }) - } + ) + hostRef.current = { send } await expect( createStructuredMailboxPointerHost().send({ sessionId: 's1', dispatchId: 'd1', operationId: 'op1', expectedRuntimeFence: 1, - body: { kind: 'message', role: 'user', blocks: [] } - } as never) - ).resolves.toEqual({ kind: 'sent', state: 'unknown' }) + body: { kind: 'message', role: 'user', blocks: [] }, + source: NOTICE_SOURCE + }) + ).resolves.toEqual({ kind: 'queued' }) + expect(send.mock.calls[0]![1]).toMatchObject({ + delivery: 'queue-if-active', + source: NOTICE_SOURCE + }) }) it('consumes mail once an accepted nudge is delivered while the worker starts (W10)', async () => { diff --git a/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts b/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts index c948a169a89..7bb79d01e3e 100644 --- a/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts +++ b/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts @@ -10,7 +10,7 @@ import { AGENT_SESSION_NOT_ATTACHED } from '../../native-chat/agent-session-wire import { getStructuredAgentSessionHost } from '../../native-chat/agent-session-wire/structured-agent-session-registry' import type { StructuredMailboxPointerHost, - StructuredPointerGateFacts + StructuredPointerSessionFacts } from './structured-mailbox-pointer-delivery' import type { AgentJournalSnapshot } from '../../../shared/agent-session-journal-types' import { @@ -36,12 +36,12 @@ export function structuredSessionPointerCallerKey(sessionId: string): string { } /** - * The idle gate for a structured session, read off its FULL reduced timeline. + * Whether a structured session is idle, for group addressing (`@idle`), read off its FULL reduced + * timeline. * * Never a bounded page. A settled turn's lifecycle item is revised in place, so on any tail window * an idle session and a busy one whose lifecycle item scrolled off look identical — and - * idle-with-history is the normal steady state of a working agent. Shared so the pointer lane and - * group addressing cannot disagree about it. + * idle-with-history is the normal steady state of a working agent. */ export async function readStructuredSessionGateFacts( sessionId: string @@ -50,12 +50,12 @@ export async function readStructuredSessionGateFacts( return snapshot ? structuredSessionGateFacts(snapshot.items) : null } -/** The pointer lane's gate: the shared idle facts, plus what each recorded send settled as. */ -async function readPointerGateFacts(sessionId: string): Promise { +/** What each recorded send settled as. */ +async function readPointerSessionFacts( + sessionId: string +): Promise { const snapshot = await readSessionJournal(sessionId) - return snapshot - ? { ...structuredSessionGateFacts(snapshot.items), submissions: snapshot.submissions } - : null + return snapshot ? { submissions: snapshot.submissions } : null } async function readSessionJournal(sessionId: string): Promise { @@ -77,8 +77,8 @@ async function readSessionJournal(sessionId: string): Promise { }) }) -describe('decideStructuredSessionPointerDelivery', () => { - it('delivers to an attached, idle session', () => { - expect(decideStructuredSessionPointerDelivery({ session: IDLE })).toEqual({ - deliver: true - }) - }) - - it('retains when the session is not attached on this host', () => { - expect(decideStructuredSessionPointerDelivery({ session: null })).toEqual({ - deliver: false, - retain: 'session-not-attached' - }) - }) - - it('retains mid-turn rather than delegating the race to the provider', () => { - expect( - decideStructuredSessionPointerDelivery({ - session: { turnRunning: true, awaitingHuman: false } - }) - ).toEqual({ deliver: false, retain: 'turn-unsettled' }) - }) - - it('names the human prompt ahead of the turn, so the retain reason is the actionable one', () => { - expect( - decideStructuredSessionPointerDelivery({ - session: { turnRunning: true, awaitingHuman: true } - }) - ).toEqual({ deliver: false, retain: 'awaiting-human' }) - }) -}) - describe('dispatch outcome classification', () => { it('marks mail delivered only on an accepted dispatch', () => { expect(structuredDispatchDelivered('accepted')).toBe(true) diff --git a/src/main/runtime/orchestration/structured-session-pointer-delivery.ts b/src/main/runtime/orchestration/structured-session-pointer-delivery.ts index 4e838e313f8..fc7758e07a1 100644 --- a/src/main/runtime/orchestration/structured-session-pointer-delivery.ts +++ b/src/main/runtime/orchestration/structured-session-pointer-delivery.ts @@ -1,13 +1,10 @@ /** - * Delivery decisions for an orchestration mail pointer aimed at a host-owned - * structured ("native") agent session. + * What orchestration mail delivery reads of a host-owned structured ("native") agent session. * - * A structured session has no PTY the pointer can be typed into, so the nudge - * travels as a session turn instead of as bytes. Everything here is pure: the - * caller supplies the session's gate facts, and gets back a decision it can - * act on. Orchestration's database stays the source of truth — - * no decision here ever consumes mail, it only says whether the nudge may be - * attempted now. + * A structured session has no PTY the pointer can be typed into, so the nudge travels as a session + * turn instead of as bytes, and a busy session's own queue holds it until the turn ends. + * Everything here is pure. Orchestration's database stays the source of truth: nothing here + * consumes mail. */ import type { AgentJournalRenderItem } from '../../../shared/agent-session-journal-types' @@ -20,22 +17,18 @@ import { export type StructuredPointerRetainReason = | 'session-not-attached' | 'turn-unsettled' - | 'awaiting-human' | 'dispatch-rejected' | 'dispatch-unknown' -export type StructuredPointerDecision = - | { deliver: true } - | { deliver: false; retain: StructuredPointerRetainReason } - /** The dispatch states both provider adapters converge on. */ export type StructuredDispatchState = 'accepted' | 'rejected' | 'unknown' /** - * What the delivery gate needs to know about a session, read once per attempt. + * Whether a session is busy, in the vocabulary group addressing (`@idle`) matches on. * * Deliberately two booleans rather than the journal: the caller reads the FULL reduced timeline - * (see `readGateFacts`), so nothing downstream can be tempted to re-derive them from a page. + * (see `readStructuredSessionGateFacts`), so nothing downstream can be tempted to re-derive them + * from a page. */ export type StructuredSessionGateFacts = { turnRunning: boolean @@ -46,7 +39,7 @@ export type StructuredSessionGateFacts = { /** * Projects the gate facts off a session's live items. * - * Reuses the projection the chat view already reads, so the delivery gate and the visible + * Reuses the projection the chat view already reads, so `@idle` and the visible * "working" state can never disagree. Both must be answered from the fully reduced timeline: a * settled turn is TOMBSTONED rather than rewritten to `completed`, so on a bounded tail page an * idle session and a running turn whose lifecycle item was pushed off the end look identical — @@ -61,37 +54,6 @@ export function structuredSessionGateFacts( } } -/** - * Decide whether the nudge may be sent right now. - * - * Mid-turn delivery is refused for both providers rather than delegated to - * them. Neither refuses the frame: Codex COALESCES a mid-turn `turn/start` into - * the running turn -- measured on codex-cli 0.147.0, 0.150.1 and 0.153.4, none - * of which refuse it and none of which fire a second `turn/started` -- and - * Claude folds it into the running turn (or runs it as the next turn when the - * turn ends first). Both therefore - * fold the nudge into work already in flight, where it reads as part of the - * running turn rather than a new instruction. Waiting for the turn to settle is - * the one contract that holds for both, and it preserves orchestration's - * existing idle-edge-only delivery policy. - */ -export function decideStructuredSessionPointerDelivery(input: { - session: StructuredSessionGateFacts | null -}): StructuredPointerDecision { - if (!input.session) { - return { deliver: false, retain: 'session-not-attached' } - } - // Checked before the turn gate: a pending prompt has no running turn, so the turn test alone - // reads it as idle, and sending there queues a nudge behind something only a human can clear. - if (input.session.awaitingHuman) { - return { deliver: false, retain: 'awaiting-human' } - } - if (input.session.turnRunning) { - return { deliver: false, retain: 'turn-unsettled' } - } - return { deliver: true } -} - /** * Only an accepted dispatch may mark mail delivered. * diff --git a/src/main/runtime/structured-chat-coordinator-mail-queue.test.ts b/src/main/runtime/structured-chat-coordinator-mail-queue.test.ts new file mode 100644 index 00000000000..104f767e333 --- /dev/null +++ b/src/main/runtime/structured-chat-coordinator-mail-queue.test.ts @@ -0,0 +1,139 @@ +import './rpc/unused-default-rpc-methods.test-fixture' +// A busy structured chat holds the orchestration pointer as a card in its own queue, sent when the +// turn ends, as it holds a message the person sends then; the queue does nothing else with it. End +// to end on the coordinator-mail rig. + +import { describe, expect, it, vi } from 'vitest' +import type { FakeConnection } from './structured-chat-coordinator-fake-codex-fixture' +import { idOf } from './rpc/orchestration-session-caller-test-fixture' +import { + COORDINATOR, + WORKER_2_PANE, + WAIT, + call, + coordinatorRunAndTask, + db, + finishWorker, + host, + openChat, + ptyPointer, + queuedCardTexts, + runtime, + sendUserMessage, + settleTurn, + turnText +} from './structured-chat-coordinator-mail-rig.test-fixture' + +/** The person's turn, started and still running; resolves to its end. */ +async function runningUserTurn(chat: FakeConnection): Promise<() => Promise> { + expect(await sendUserMessage(COORDINATOR, 'go')).toMatchObject({ ok: true }) + await vi.waitFor(() => expect(chat.turns).toHaveLength(1), WAIT) + const notify = (method: string, params: unknown) => chat.handlers.onNotification?.(method, params) + notify('turn/started', { turn: { id: 'turn-1' } }) + notify('item/completed', { + item: { + type: 'userMessage', + id: 'echo-go', + clientId: chat.turns[0]!.clientUserMessageId, + content: [{ type: 'text', text: 'go' }] + } + }) + await host.flushStreamedEvents(COORDINATOR) + return async () => { + notify('turn/completed', { turn: { id: 'turn-1' } }) + await host.flushStreamedEvents(COORDINATOR) + } +} + +/** Idle edges with nothing owed: whatever they would send gets the time to show. */ +async function idleEdgesSettled(): Promise { + for (let edge = 0; edge < 3; edge += 1) { + runtime.onStructuredSessionStatusForMail({ sessionId: COORDINATOR, status: 'idle' }) + await new Promise((resolve) => setTimeout(resolve, 100)) + } +} + +/** The chat's queue as its journal stores it. */ +function queuedRows() { + return host.collaboratorsForTests().sessions.get(COORDINATOR)?.journal.queuedMessages.list() ?? [] +} + +/** A second task, for a second worker result. */ +async function secondTask(): Promise { + return idOf( + (await call('orchestration.taskCreate', { spec: 'more' }, { sessionId: COORDINATOR })).task + ) +} + +describe("a busy chat's orchestration pointer waits in its queue", () => { + it('queues the pointer as a card, with who it is from, and sends it once when the turn ends', async () => { + const chat = await openChat(COORDINATOR) + const { runId, taskId } = await coordinatorRunAndTask() + const endTurn = await runningUserTurn(chat) + await finishWorker(taskId) + await vi.waitFor( + async () => expect(await queuedCardTexts()).toEqual([ptyPointer(`run:${runId}`)]), + WAIT + ) + expect(chat.turns).toHaveLength(1) + const [card] = queuedRows() + const [mail] = db.getAllMessages(`run:${runId}`) + expect(card?.source).toEqual({ + kind: 'agent', + senders: [ + { party: { address: 'term_worker', terminalHandle: 'term_worker', orcaSessionId: null } } + ], + orchestration: { + message: 'mail-notice', + mailbox: `run:${runId}`, + dispatchId: null, + messages: [{ messageId: mail!.id, runId, from: 'term_worker' }] + } + }) + + await endTurn() + await vi.waitFor(() => expect(chat.turns).toHaveLength(2), WAIT) + expect(turnText(chat.turns[1]!)).toBe(ptyPointer(`run:${runId}`)) + expect(await queuedCardTexts()).toEqual([]) + await settleTurn(COORDINATOR, 1) + await idleEdgesSettled() + expect(chat.turns).toHaveLength(2) + }) + + it('queues a second card for mail that arrives while the first waits, each counting its own mail', async () => { + const chat = await openChat(COORDINATOR) + const { runId, taskId } = await coordinatorRunAndTask() + const second = await secondTask() + const endTurn = await runningUserTurn(chat) + await finishWorker(taskId) + await vi.waitFor(async () => expect(await queuedCardTexts()).toHaveLength(1), WAIT) + await finishWorker(second, { handle: 'term_worker_2', paneKey: WORKER_2_PANE }) + const pointer = ptyPointer(`run:${runId}`) + await vi.waitFor(async () => expect(await queuedCardTexts()).toEqual([pointer, pointer]), WAIT) + + await endTurn() + await vi.waitFor(() => expect(chat.turns).toHaveLength(2), WAIT) + await settleTurn(COORDINATOR, 1) + await vi.waitFor(() => expect(chat.turns).toHaveLength(3), WAIT) + expect(turnText(chat.turns[2]!)).toBe(pointer) + await settleTurn(COORDINATOR, 2) + await idleEdgesSettled() + expect(chat.turns).toHaveLength(3) + expect(await queuedCardTexts()).toEqual([]) + }) + + it("leaves the chat's own `check` as it is: the mail stays readable, and the card stays", async () => { + const chat = await openChat(COORDINATOR) + const { runId, taskId } = await coordinatorRunAndTask() + const endTurn = await runningUserTurn(chat) + await finishWorker(taskId) + await vi.waitFor(async () => expect(await queuedCardTexts()).toHaveLength(1), WAIT) + const [mail] = db.getAllMessages(`run:${runId}`) + expect(await call('orchestration.check', {}, { sessionId: COORDINATOR })).toMatchObject({ + count: 1, + messages: [{ id: mail!.id }] + }) + expect(await queuedCardTexts()).toEqual([ptyPointer(`run:${runId}`)]) + await endTurn() + }) +}) diff --git a/src/main/runtime/structured-chat-coordinator-mail-rig.test-fixture.ts b/src/main/runtime/structured-chat-coordinator-mail-rig.test-fixture.ts new file mode 100644 index 00000000000..6fadd54a13b --- /dev/null +++ b/src/main/runtime/structured-chat-coordinator-mail-rig.test-fixture.ts @@ -0,0 +1,324 @@ +// The coordinator-mail rig, shared by every suite that drives a worker's result into a structured +// chat end to end in one process. +// +// Real: the structured agent-session host, its record store, journal, lease and Codex adapter; the +// orchestration database, RPC dispatcher and methods; the runtime's pointer lanes. Fake: only the +// Codex app-server child, which answers the JSON-RPC calls the real one does. Importing it +// registers the rig's own beforeEach/afterEach for the importing file. + +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, expect, vi } from 'vitest' +import type { AgentJournalRenderItem } from '../../shared/agent-session-journal-types' +import { computeAgentSessionPayloadFingerprint } from '../../shared/agent-session-mutation-envelope' +import { ORCHESTRATION_CONTRACT_VERSION } from '../../shared/protocol-version' +import type { StructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-host' +import { agentSessionProviderHandleChainHead } from '../../shared/agent-session-provider-handle' +import { OrcaRuntimeService } from './orca-runtime' +import { OrchestrationDb } from './orchestration/db' +import { localOrchestrationCliCommand } from './orchestration/cli-command' +import { formatMessagePointer } from './orchestration/formatter' +import { RpcDispatcher } from './rpc/dispatcher' +import { ORCHESTRATION_METHODS } from './rpc/methods/orchestration' +import { idOf, isRecord, resultOf } from './rpc/orchestration-session-caller-test-fixture' +import { + ensureStructuredAgentSessionHost, + stopStructuredAgentSessionRuntime +} from './structured-agent-session-runtime' +import { createCoordinatorMailObservationClock } from './structured-chat-coordinator-observation-clock.test-fixture' +import { + attachParams, + fakeCodex, + operationId, + resetProviderFaults, + type FakeConnection +} from './structured-chat-coordinator-fake-codex-fixture' +import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger' + +export const COORDINATOR = '4a1f6c2e-8b3d-4e7a-9c15-0d2b6e8f1a37' +export const PEER_CHAT = '7e3b9d15-2c4a-4f86-a0b1-5c9e2d7f3b64' +export const WORKER_PANE = 'tab_worker:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb' +export const WORKER_2_PANE = 'tab_worker2:cccccccc-cccc-4ccc-8ccc-cccccccccccc' + +export let codex: ReturnType +export let root: string +export let runtime: OrcaRuntimeService +export let db: OrchestrationDb +export let host: StructuredAgentSessionHost +export let dispatcher: RpcDispatcher +export let requests = 0 +export const observationClock = createCoordinatorMailObservationClock(() => host, COORDINATOR) + +export function request( + method: string, + params: Record, + options: { sessionId?: string } = {} +): Parameters[0] { + requests += 1 + return { + id: `rpc-${requests}`, + authToken: 'test', + method, + params, + orchestrationContractVersion: ORCHESTRATION_CONTRACT_VERSION, + orchestrationRequestId: `req-${requests}`, + ...(options.sessionId + ? { orchestrationCompatibilityEvidence: { agentSessionId: options.sessionId } } + : {}) + } +} + +export async function call( + method: string, + params: Record, + options?: { sessionId?: string } +): Promise> { + const response = await dispatcher.dispatch(request(method, params, options)) + if (!response.ok) { + throw new Error(`${method} failed: ${JSON.stringify(response)}`) + } + return resultOf(response) +} + +export async function openChat(sessionId: string): Promise { + const attached = await host.attach({ callerKey: 'test-surface' }, attachParams(sessionId)) + expect(attached, JSON.stringify(attached)).toMatchObject({ ok: true }) + await host.setSessionTabVisibility(sessionId, true) + threadBySession.set(sessionId, codex.connections.at(-1)!.threadId!) + return connectionFor(sessionId) +} + +export const threadBySession = new Map() + +export function connectionFor(sessionId: string): FakeConnection { + // A cleared chat's successor starts on its first message; its record then names its thread. + const head = agentSessionProviderHandleChainHead( + host.deps.store.getRecord(sessionId)?.providerHandleChain ?? [] + ) + const thread = threadBySession.get(sessionId) ?? head?.handle.nativeId + const connection = codex.connections.findLast((candidate) => candidate.threadId === thread) + if (!connection) { + throw new Error(`no app-server for ${sessionId}`) + } + return connection +} + +/** Codex's own sequence for a turn: it starts, echoes the user message, and completes. */ +export async function settleTurn(sessionId: string, turnIndex: number): Promise { + const connection = connectionFor(sessionId) + const turn = connection.turns[turnIndex]! + const turnId = `turn-${turnIndex + 1}` + const notify = (method: string, params: unknown) => + connection.handlers.onNotification?.(method, params) + notify('turn/started', { turn: { id: turnId } }) + notify('item/completed', { + item: { + type: 'userMessage', + id: `echo-${turn.clientUserMessageId}`, + clientId: turn.clientUserMessageId, + content: [{ type: 'text', text: 'pointer' }] + } + }) + notify('turn/completed', { turn: { id: turnId } }) + await host.flushStreamedEvents(sessionId) +} + +/** A user message typed into the chat, as the chat surface sends it. */ +export function sendUserMessage(sessionId: string, text: string) { + const body = { + kind: 'message' as const, + role: 'user' as const, + blocks: [{ type: 'text' as const, text }] + } + return host.send( + { callerKey: 'test-surface' }, + { + envelope: { + sessionId, + clientOperationId: operationId(), + expectedRuntimeFence: host.deps.store.getRecord(sessionId)!.lease.runtimeFence, + payloadFingerprint: computeAgentSessionPayloadFingerprint({ + method: 'agentSession.send', + sessionId, + fields: { body } + }) + }, + body + } + ) +} + +export async function userTexts(sessionId: string): Promise { + return (await host.journalSnapshot(sessionId)).items.flatMap((item: AgentJournalRenderItem) => + item.body?.kind === 'message' && item.body.role === 'user' + ? item.body.blocks.map((block) => (block.type === 'text' ? block.text : '')) + : [] + ) +} + +/** A supervised terminal worker under the coordinator's Run, and its worker_done. */ +export async function finishWorker( + taskId: string, + worker: { handle: string; paneKey: string } = { handle: 'term_worker', paneKey: WORKER_PANE } +): Promise { + const started = db.createStartingWorkerDispatch({ + creator: { kind: 'system' }, + maxDepth: Number.MAX_SAFE_INTEGER, + taskId, + startOptions: {} + }) + db.prepareStartingWorkerAuthority({ + dispatchId: started.dispatch.id, + handle: worker.handle, + paneKey: worker.paneKey, + processIncarnation: `runtime_test:${worker.handle}:1`, + worktreeId: 'repo::worker', + effects: [], + setupState: 'not_applicable' + }) + db.markWorkerDispatchReady(started.dispatch.id) + await call('orchestration.send', { + from: worker.handle, + subject: 'Done', + type: 'worker_done', + payload: JSON.stringify({ taskId, dispatchId: started.dispatch.id, outcome: 'succeeded' }) + }) +} + +export async function coordinatorRunAndTask(): Promise<{ runId: string; taskId: string }> { + const created = await call( + 'orchestration.runCreate', + { objective: 'ship' }, + { + sessionId: COORDINATOR + } + ) + const runId = idOf(created.run) + const task = await call( + 'orchestration.taskCreate', + { spec: 'build it' }, + { + sessionId: COORDINATOR + } + ) + return { runId, taskId: idOf(task.task) } +} + +/** `/clear` as the chat surface runs it: the conversation continues in a new session. */ +export async function clearChat(sessionId: string): Promise { + const command = 'clear' as const + const cleared = await host.conversationCommand( + { callerKey: 'test-surface' }, + { + command, + envelope: { + sessionId, + clientOperationId: operationId(), + expectedRuntimeFence: host.deps.store.getRecord(sessionId)!.lease.runtimeFence, + payloadFingerprint: computeAgentSessionPayloadFingerprint({ + method: 'agentSession.conversationCommand', + sessionId, + fields: { command } + }) + } + } + ) + const successor = cleared.ok ? cleared.value.replacementSessionId : undefined + if (!successor) { + throw new Error(`clear failed: ${JSON.stringify(cleared)}`) + } + // The surface swaps the tab over to the session that continues the chat. + await host.setSessionTabVisibility(sessionId, false) + await host.setSessionTabVisibility(successor, true) + return successor +} + +/** A cleared chat's successor runs once the user writes to it; only then can its agent act. */ +export async function startSuccessor(successor: string): Promise { + expect(await sendUserMessage(successor, 'hello')).toMatchObject({ ok: true }) + await vi.waitFor(() => expect(connectionFor(successor).turns).toHaveLength(1), WAIT) + await settleTurn(successor, 0) +} + +beforeEach(async () => { + resetProviderFaults() + root = await mkdtemp(join(tmpdir(), 'orca-structured-coordinator-mail-')) + codex = fakeCodex() + db = new OrchestrationDb(':memory:') + runtime = startRuntime() + host = await ensureStructuredAgentSessionHost({ + logger: createStructuredAgentSessionLogger(), + stateDirectory: root, + hostId: 'local', + claimKeyId: 'key-1', + resolveWorkspacePath: async (workspaceId) => `/repos/${workspaceId}`, + resolveCodexCommand: () => '/usr/local/bin/codex', + resolveClaudeAuthPolicy: () => ({ stripAuthEnv: true }), + resolveEnvironment: async () => ({ PATH: '/usr/bin' }), + openCodexConnection: codex.openConnection, + readProcessStartTime: async () => 1_700_000_000_000, + // The same calls the runtime's own host install makes. + onSessionStatusChanged: (summary) => runtime.onStructuredSessionStatusForMail(summary) + }) + dispatcher = new RpcDispatcher({ runtime, methods: ORCHESTRATION_METHODS }) +}) + +/** The runtime over the shared database; a second call is what an Orca restart leaves behind. */ +export function startRuntime(): OrcaRuntimeService { + const started = new OrcaRuntimeService() + started.setOrchestrationDb(db) + vi.spyOn(started, 'ensureStructuredAgentSessionHost').mockResolvedValue() + vi.spyOn(started, 'getTerminalPaneKey').mockImplementation((handle) => + handle === 'term_worker' ? WORKER_PANE : handle === 'term_worker_2' ? WORKER_2_PANE : null + ) + return started +} + +afterEach(async () => { + try { + await stopStructuredAgentSessionRuntime() + db.close() + await observationClock.drainClosedDatabaseRepair() + vi.restoreAllMocks() + await rm(root, { recursive: true, force: true }) + } finally { + observationClock.restore() + } +}) + +// Pointers are sent on asynchronous edges; the default 1s wait is too tight under a loaded parallel run. +export const WAIT = { timeout: 10_000 } + +export const POINTER = + /You have 1 orchestration message\. Run `orca(-dev)? orchestration check --run run_\w+`\./ + +/** The text the PTY lane types into a local terminal for this mailbox, byte for byte. */ +export function ptyPointer(mailboxHandle: string): string { + return formatMessagePointer(1, mailboxHandle, localOrchestrationCliCommand()).trim() +} + +/** The text of a turn the fake provider received. */ +export function turnText(turn: { text: string }): string { + const input: unknown = JSON.parse(turn.text) + return Array.isArray(input) + ? input.map((item: unknown) => (isRecord(item) ? String(item.text) : '')).join('') + : '' +} + +/** What an Orca restart leaves behind: a new runtime over the same database and host. */ +export function restartRuntime(): void { + runtime = startRuntime() + dispatcher = new RpcDispatcher({ runtime, methods: ORCHESTRATION_METHODS }) +} + +/** The text of every card the chat lists in its queue, as the person sees it. */ +export async function queuedCardTexts(sessionId = COORDINATOR): Promise { + const page = await host.history({ sessionId, direction: 'tail' }) + if (!page.ok) { + throw new Error('history refused') + } + return (page.page.queuedMessages ?? []).flatMap((card) => + card.body.blocks.map((block) => (block.type === 'text' ? block.text : '')) + ) +} diff --git a/src/main/runtime/structured-chat-coordinator-mail.test.ts b/src/main/runtime/structured-chat-coordinator-mail.test.ts index 29a2d5bb79f..5eb0f57209c 100644 --- a/src/main/runtime/structured-chat-coordinator-mail.test.ts +++ b/src/main/runtime/structured-chat-coordinator-mail.test.ts @@ -1,319 +1,51 @@ import './rpc/unused-default-rpc-methods.test-fixture' -// A worker's result reaching the structured chat that coordinates it, end to end in one process. -// -// Real: the structured agent-session host, its record store, journal, lease and Codex adapter; the -// orchestration database, RPC dispatcher and methods; the runtime's pointer lanes. Fake: only the -// Codex app-server child, which answers the JSON-RPC calls the real one does. +// A worker's result reaching the structured chat that coordinates it, end to end in one process, +// on the coordinator-mail rig. -import { mkdtemp, rm } from 'node:fs/promises' -import { tmpdir } from 'node:os' -import { join } from 'node:path' -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' -import type { AgentJournalRenderItem } from '../../shared/agent-session-journal-types' +import { describe, expect, it, vi } from 'vitest' import { agentJournalSubmissionKey } from '../../shared/agent-session-journal-item-key' import { computeAgentSessionPayloadFingerprint } from '../../shared/agent-session-mutation-envelope' -import { ORCHESTRATION_CONTRACT_VERSION } from '../../shared/protocol-version' import { AgentSessionAcquisitionRefusal, AgentSessionPreSpawnError } from '../native-chat/agent-session-wire/structured-agent-session-adapter' -import type { StructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-host' import type { AgentSessionJournal } from '../native-chat/agent-session-journal/journal-store' import { AGENT_SESSION_MAX_NEW_OPERATION_AGE_MS } from '../../shared/agent-session-host-authority' import { refuse } from '../../shared/agent-session-wire-refusals' -import { agentSessionProviderHandleChainHead } from '../../shared/agent-session-provider-handle' -import { OrcaRuntimeService } from './orca-runtime' -import { OrchestrationDb } from './orchestration/db' import { localOrchestrationCliCommand } from './orchestration/cli-command' import { formatMessagePointer } from './orchestration/formatter' import { currentRunCoordinatorOrcaSessionId } from './orchestration/db/runs/run-coordinator-orca-session' -import { RpcDispatcher } from './rpc/dispatcher' -import { ORCHESTRATION_METHODS } from './rpc/methods/orchestration' -import { idOf, isRecord, resultOf } from './rpc/orchestration-session-caller-test-fixture' +import { idOf } from './rpc/orchestration-session-caller-test-fixture' +import { operationId, providerFaults } from './structured-chat-coordinator-fake-codex-fixture' + import { - ensureStructuredAgentSessionHost, - stopStructuredAgentSessionRuntime -} from './structured-agent-session-runtime' -import { createCoordinatorMailObservationClock } from './structured-chat-coordinator-observation-clock.test-fixture' -import { - attachParams, - fakeCodex, - operationId, - providerFaults, - resetProviderFaults, - type FakeConnection -} from './structured-chat-coordinator-fake-codex-fixture' -import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger' - -const COORDINATOR = '4a1f6c2e-8b3d-4e7a-9c15-0d2b6e8f1a37' -const PEER_CHAT = '7e3b9d15-2c4a-4f86-a0b1-5c9e2d7f3b64' -const WORKER_PANE = 'tab_worker:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb' -const WORKER_2_PANE = 'tab_worker2:cccccccc-cccc-4ccc-8ccc-cccccccccccc' - -let codex: ReturnType -let root: string -let runtime: OrcaRuntimeService -let db: OrchestrationDb -let host: StructuredAgentSessionHost -let dispatcher: RpcDispatcher -let requests = 0 -const observationClock = createCoordinatorMailObservationClock(() => host, COORDINATOR) - -function request( - method: string, - params: Record, - options: { sessionId?: string } = {} -): Parameters[0] { - requests += 1 - return { - id: `rpc-${requests}`, - authToken: 'test', - method, - params, - orchestrationContractVersion: ORCHESTRATION_CONTRACT_VERSION, - orchestrationRequestId: `req-${requests}`, - ...(options.sessionId - ? { orchestrationCompatibilityEvidence: { agentSessionId: options.sessionId } } - : {}) - } -} - -async function call( - method: string, - params: Record, - options?: { sessionId?: string } -): Promise> { - const response = await dispatcher.dispatch(request(method, params, options)) - if (!response.ok) { - throw new Error(`${method} failed: ${JSON.stringify(response)}`) - } - return resultOf(response) -} - -async function openChat(sessionId: string): Promise { - const attached = await host.attach({ callerKey: 'test-surface' }, attachParams(sessionId)) - expect(attached, JSON.stringify(attached)).toMatchObject({ ok: true }) - await host.setSessionTabVisibility(sessionId, true) - threadBySession.set(sessionId, codex.connections.at(-1)!.threadId!) - return connectionFor(sessionId) -} - -const threadBySession = new Map() - -function connectionFor(sessionId: string): FakeConnection { - // A cleared chat's successor starts on its first message; its record then names its thread. - const head = agentSessionProviderHandleChainHead( - host.deps.store.getRecord(sessionId)?.providerHandleChain ?? [] - ) - const thread = threadBySession.get(sessionId) ?? head?.handle.nativeId - const connection = codex.connections.findLast((candidate) => candidate.threadId === thread) - if (!connection) { - throw new Error(`no app-server for ${sessionId}`) - } - return connection -} - -/** Codex's own sequence for a turn: it starts, echoes the user message, and completes. */ -async function settleTurn(sessionId: string, turnIndex: number): Promise { - const connection = connectionFor(sessionId) - const turn = connection.turns[turnIndex]! - const turnId = `turn-${turnIndex + 1}` - const notify = (method: string, params: unknown) => - connection.handlers.onNotification?.(method, params) - notify('turn/started', { turn: { id: turnId } }) - notify('item/completed', { - item: { - type: 'userMessage', - id: `echo-${turn.clientUserMessageId}`, - clientId: turn.clientUserMessageId, - content: [{ type: 'text', text: 'pointer' }] - } - }) - notify('turn/completed', { turn: { id: turnId } }) - await host.flushStreamedEvents(sessionId) -} - -/** A user message typed into the chat, as the chat surface sends it. */ -function sendUserMessage(sessionId: string, text: string) { - const body = { - kind: 'message' as const, - role: 'user' as const, - blocks: [{ type: 'text' as const, text }] - } - return host.send( - { callerKey: 'test-surface' }, - { - envelope: { - sessionId, - clientOperationId: operationId(), - expectedRuntimeFence: host.deps.store.getRecord(sessionId)!.lease.runtimeFence, - payloadFingerprint: computeAgentSessionPayloadFingerprint({ - method: 'agentSession.send', - sessionId, - fields: { body } - }) - }, - body - } - ) -} - -async function userTexts(sessionId: string): Promise { - return (await host.journalSnapshot(sessionId)).items.flatMap((item: AgentJournalRenderItem) => - item.body?.kind === 'message' && item.body.role === 'user' - ? item.body.blocks.map((block) => (block.type === 'text' ? block.text : '')) - : [] - ) -} - -/** A supervised terminal worker under the coordinator's Run, and its worker_done. */ -async function finishWorker( - taskId: string, - worker: { handle: string; paneKey: string } = { handle: 'term_worker', paneKey: WORKER_PANE } -): Promise { - const started = db.createStartingWorkerDispatch({ - creator: { kind: 'system' }, - maxDepth: Number.MAX_SAFE_INTEGER, - taskId, - startOptions: {} - }) - db.prepareStartingWorkerAuthority({ - dispatchId: started.dispatch.id, - handle: worker.handle, - paneKey: worker.paneKey, - processIncarnation: `runtime_test:${worker.handle}:1`, - worktreeId: 'repo::worker', - effects: [], - setupState: 'not_applicable' - }) - db.markWorkerDispatchReady(started.dispatch.id) - await call('orchestration.send', { - from: worker.handle, - subject: 'Done', - type: 'worker_done', - payload: JSON.stringify({ taskId, dispatchId: started.dispatch.id, outcome: 'succeeded' }) - }) -} - -async function coordinatorRunAndTask(): Promise<{ runId: string; taskId: string }> { - const created = await call( - 'orchestration.runCreate', - { objective: 'ship' }, - { - sessionId: COORDINATOR - } - ) - const runId = idOf(created.run) - const task = await call( - 'orchestration.taskCreate', - { spec: 'build it' }, - { - sessionId: COORDINATOR - } - ) - return { runId, taskId: idOf(task.task) } -} - -/** `/clear` as the chat surface runs it: the conversation continues in a new session. */ -async function clearChat(sessionId: string): Promise { - const command = 'clear' as const - const cleared = await host.conversationCommand( - { callerKey: 'test-surface' }, - { - command, - envelope: { - sessionId, - clientOperationId: operationId(), - expectedRuntimeFence: host.deps.store.getRecord(sessionId)!.lease.runtimeFence, - payloadFingerprint: computeAgentSessionPayloadFingerprint({ - method: 'agentSession.conversationCommand', - sessionId, - fields: { command } - }) - } - } - ) - const successor = cleared.ok ? cleared.value.replacementSessionId : undefined - if (!successor) { - throw new Error(`clear failed: ${JSON.stringify(cleared)}`) - } - // The surface swaps the tab over to the session that continues the chat. - await host.setSessionTabVisibility(sessionId, false) - await host.setSessionTabVisibility(successor, true) - return successor -} - -/** A cleared chat's successor runs once the user writes to it; only then can its agent act. */ -async function startSuccessor(successor: string): Promise { - expect(await sendUserMessage(successor, 'hello')).toMatchObject({ ok: true }) - await vi.waitFor(() => expect(connectionFor(successor).turns).toHaveLength(1), WAIT) - await settleTurn(successor, 0) -} - -beforeEach(async () => { - resetProviderFaults() - root = await mkdtemp(join(tmpdir(), 'orca-structured-coordinator-mail-')) - codex = fakeCodex() - db = new OrchestrationDb(':memory:') - runtime = startRuntime() - host = await ensureStructuredAgentSessionHost({ - logger: createStructuredAgentSessionLogger(), - stateDirectory: root, - hostId: 'local', - claimKeyId: 'key-1', - resolveWorkspacePath: async (workspaceId) => `/repos/${workspaceId}`, - resolveCodexCommand: () => '/usr/local/bin/codex', - resolveClaudeAuthPolicy: () => ({ stripAuthEnv: true }), - resolveEnvironment: async () => ({ PATH: '/usr/bin' }), - openCodexConnection: codex.openConnection, - readProcessStartTime: async () => 1_700_000_000_000, - // The same call the runtime's own host install makes on every status change. - onSessionStatusChanged: (summary) => runtime.onStructuredSessionStatusForMail(summary) - }) - dispatcher = new RpcDispatcher({ runtime, methods: ORCHESTRATION_METHODS }) -}) - -/** The runtime over the shared database; a second call is what an Orca restart leaves behind. */ -function startRuntime(): OrcaRuntimeService { - const started = new OrcaRuntimeService() - started.setOrchestrationDb(db) - vi.spyOn(started, 'ensureStructuredAgentSessionHost').mockResolvedValue() - vi.spyOn(started, 'getTerminalPaneKey').mockImplementation((handle) => - handle === 'term_worker' ? WORKER_PANE : handle === 'term_worker_2' ? WORKER_2_PANE : null - ) - return started -} - -afterEach(async () => { - try { - await stopStructuredAgentSessionRuntime() - db.close() - await observationClock.drainClosedDatabaseRepair() - vi.restoreAllMocks() - await rm(root, { recursive: true, force: true }) - } finally { - observationClock.restore() - } -}) - -// Pointers are sent on asynchronous edges; the default 1s wait is too tight under a loaded parallel run. -const WAIT = { timeout: 10_000 } - -const POINTER = - /You have 1 orchestration message\. Run `orca(-dev)? orchestration check --run run_\w+`\./ - -/** The text the PTY lane types into a local terminal for this mailbox, byte for byte. */ -function ptyPointer(mailboxHandle: string): string { - return formatMessagePointer(1, mailboxHandle, localOrchestrationCliCommand()).trim() -} - -/** The text of a turn the fake provider received. */ -function turnText(turn: { text: string }): string { - const input: unknown = JSON.parse(turn.text) - return Array.isArray(input) - ? input.map((item: unknown) => (isRecord(item) ? String(item.text) : '')).join('') - : '' -} + COORDINATOR, + PEER_CHAT, + WORKER_2_PANE, + codex, + runtime, + db, + host, + dispatcher, + observationClock, + request, + call, + openChat, + connectionFor, + settleTurn, + sendUserMessage, + userTexts, + finishWorker, + coordinatorRunAndTask, + clearChat, + startSuccessor, + WAIT, + POINTER, + ptyPointer, + turnText, + queuedCardTexts, + restartRuntime +} from './structured-chat-coordinator-mail-rig.test-fixture' describe('a worker result reaches the structured chat that coordinates it', () => { it('lands as a turn in the coordinator journal, and a flagless check returns the worker_done', async () => { @@ -522,8 +254,7 @@ describe('a worker result reaches the structured chat that coordinates it', () = // The next process: a fresh runtime over the same database redrives restored mail. The // provider still dies, so exactly one start proves it is pointed once, not in a loop. - runtime = startRuntime() - dispatcher = new RpcDispatcher({ runtime, methods: ORCHESTRATION_METHODS }) + restartRuntime() const before = providerFaults.starts await vi.waitFor(() => expect(providerFaults.turnStarts).toBe(2), WAIT) await observationClock.observe(1_500) @@ -658,9 +389,10 @@ describe('a worker result reaches the structured chat that coordinates it', () = } }) - it('holds mail a refused turn left in doubt until the next result, then points it once', async () => { + it("holds mail a refused turn left in doubt, then queues the next pointer behind it as the person's message would wait", async () => { // A failed turn/start cannot prove the turn never started, so the host records it `unknown` - // and a resend under its id replays that; new mail is a new send. + // and a resend under its id replays that. A live doubt counts as work still owed, so the next + // result's pointer waits in the chat's queue, as a message the person sent then would. const chat = await openChat(COORDINATOR) const { runId, taskId } = await coordinatorRunAndTask() const second = await call( @@ -676,10 +408,14 @@ describe('a worker result reaches the structured chat that coordinates it', () = expect(chat.turns).toHaveLength(0) await finishWorker(idOf(second.task), { handle: 'term_worker_2', paneKey: WORKER_2_PANE }) - await vi.waitFor(() => expect(chat.turns).toHaveLength(1), WAIT) - expect(turnText(chat.turns[0]!)).toBe( - formatMessagePointer(2, `run:${runId}`, localOrchestrationCliCommand()).trim() + await vi.waitFor( + async () => + expect(await queuedCardTexts()).toEqual([ + formatMessagePointer(2, `run:${runId}`, localOrchestrationCliCommand()).trim() + ]), + WAIT ) + expect(chat.turns).toHaveLength(0) expect(codex.connections.length).toBe(before) }) @@ -879,10 +615,11 @@ describe('a /clear keeps the chat its orchestration address', () => { expect(sent).toMatchObject({ message: { to_handle: `orca_session_id:${PEER_CHAT}` } }) await vi.waitFor(() => expect(connectionFor(successor).turns).toHaveLength(index + 1), WAIT) await settleTurn(successor, index) + // Read each ping before the next is sent, so each check holds exactly one. + const checked = await call('orchestration.check', {}, { sessionId: successor }) + expect(checked).toMatchObject({ count: 1, messages: [{ subject: `ping ${index}` }] }) + await call('orchestration.check', { ack: checked.deliveryId }, { sessionId: successor }) } - await expect(call('orchestration.check', {}, { sessionId: successor })).resolves.toMatchObject({ - count: 3 - }) }) }) diff --git a/src/shared/agent-session-message-source.ts b/src/shared/agent-session-message-source.ts new file mode 100644 index 00000000000..f0ab9a71f13 --- /dev/null +++ b/src/shared/agent-session-message-source.ts @@ -0,0 +1,88 @@ +// Who a chat message is from: the person at the composer, or another agent through Orca. +// Persisted with a queued card (`queued_messages.source_json`), so the chat can name each sender. + +import { z } from 'zod' +import { isOrcaSessionId, type OrcaSessionId } from './orca-session-address' +import type { OrchestrationPartyIdentity } from './orchestration-party-identity' + +/** + * An agent a message is from, named by the orchestration database of the host that stores the + * message: the only host whose agents can send today. A relayed sender adds its host here. No pane + * key: it reads and consumes that agent's mailbox, so the host resolves it from the handle. + */ +export type AgentMessageSender = Readonly<{ party: Omit }> + +/** One orchestration message a notice points at: its record, and its sender's `senders` address. */ +export type OrchestrationMailMessage = Readonly<{ messageId: string; runId: string; from: string }> + +/** "You have N orchestration messages": the pointer a terminal agent is typed, for a mailbox's + * unread mail, which the agent reads with `check`. */ +export type OrchestrationMailNotice = Readonly<{ + message: 'mail-notice' + mailbox: string + dispatchId: string | null + messages: readonly OrchestrationMailMessage[] +}> + +/** What Orca delivers for other agents, one shape per message kind. */ +export type OrchestrationAgentMessage = OrchestrationMailNotice + +export type AgentMessageSource = Readonly<{ + kind: 'agent' + /** Every distinct sender of the messages it carries, in mail order. */ + senders: readonly AgentMessageSender[] + orchestration: OrchestrationAgentMessage +}> + +export type AgentSessionMessageSource = Readonly<{ kind: 'user' }> | AgentMessageSource + +export const USER_MESSAGE_SOURCE: AgentSessionMessageSource = { kind: 'user' } + +const MESSAGE_SOURCE_VERSION = 1 + +const orcaSessionIdSchema = z.custom( + (value) => typeof value === 'string' && isOrcaSessionId(value) +) + +const mailNoticeSchema = z.object({ + message: z.literal('mail-notice'), + mailbox: z.string(), + dispatchId: z.string().nullable(), + messages: z.array(z.object({ messageId: z.string(), runId: z.string(), from: z.string() })) +}) + +// Not strict: a newer build may add a field, which this one keeps no use for and must not reject. +const storedSourceSchema = z.discriminatedUnion('kind', [ + z.object({ v: z.literal(MESSAGE_SOURCE_VERSION), kind: z.literal('user') }), + z.object({ + v: z.literal(MESSAGE_SOURCE_VERSION), + kind: z.literal('agent'), + senders: z.array( + z.object({ + party: z.object({ + address: z.string(), + terminalHandle: z.string().nullable(), + orcaSessionId: orcaSessionIdSchema.nullable() + }) + }) + ), + orchestration: z.discriminatedUnion('message', [mailNoticeSchema]) + }) +]) + +export function serializeAgentSessionMessageSource(source: AgentSessionMessageSource): string { + return JSON.stringify({ v: MESSAGE_SOURCE_VERSION, ...source }) +} + +/** + * The stored value read back. Absent (a card from before the column) is the person's: only the + * composer queued then. So is a value this build cannot read; either way it is sent as written. + */ +export function readAgentSessionMessageSource(stored: unknown): AgentSessionMessageSource { + const parsed = storedSourceSchema.safeParse(stored) + if (!parsed.success) { + return USER_MESSAGE_SOURCE + } + const { v: _version, ...source } = parsed.data + return source +} diff --git a/src/shared/orchestration-party-identity.ts b/src/shared/orchestration-party-identity.ts new file mode 100644 index 00000000000..8bbce97e9f7 --- /dev/null +++ b/src/shared/orchestration-party-identity.ts @@ -0,0 +1,18 @@ +import type { OrcaSessionId } from './orca-session-address' + +/** + * Who an orchestration party is, as Run binding and mail routing match it. + * + * A PTY agent is its terminal: a handle and a pane key, no Orca session id. An agent that is a + * structured session is its Orca session id, addressed as `orca_session_id:`; a structured worker also + * has the handle and pane key it was minted, and an ordinary chat has neither. Methods pass this + * through whole and never branch on which fields are set; the lookups that build it own that. + */ +export type OrchestrationPartyIdentity = Readonly<{ + /** Mailbox address the party sends from and reads: its terminal handle, else its session address. */ + address: string + terminalHandle: string | null + paneKey: string | null + /** The bare Orca session id the party is addressed by; mail spells it `orca_session_id:`. */ + orcaSessionId: OrcaSessionId | null +}>