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 +}>