feat(orchestration): a native chat gets the orchestration pointer a CLI agent gets, through the same send as your messages (#25078)

* feat(orchestration): agent mail to a busy chat waits in the chat's own queue

A mail notice for a busy structured chat used to wait in the orchestration lane's
own invisible "until the chat is free" gate. It now goes through the queue a
person's message uses: sendAgentTurn(..., { delivery: 'queue', source }) makes
the host hold it as a draft card the person sees and can Steer or delete, and
the queue sends it when the turn ends.

- The lane's busy gate (turnRunning / awaitingHuman) is gone; the queue's own
  hold decides. Its parking now only waits for its own send or card to settle.
- A queued card's hand-off (sent under a fresh id) is matched by
  queuedMessageId, so it stamps the mail once and never sends a second notice.
- A card from before an Orca restart is still the lane's: no second card.
- A card the person deleted counts as handled for that batch; newer mail
  notifies again. No stored flag: read off the card and delivered_at.
- `orchestration check` waits for the lane to withdraw a card whose mail it
  read, so a stale notice is never sent; more mail replaces the card.
- Each card records who queued it in queued_messages.source_json (versioned,
  schema-validated): the person, or Orca for agents, with every distinct
  sender as an orchestration party plus host id, and the mailbox, dispatch,
  run and message ids. The restart pause holds only the person's cards.
- sendAgentTurn answers with a snapshot, not the journal's live submission.

* test(orchestration): type the mail fixture as a pointer batch message

* fix(orchestration): the queue judges an agent's mail notice as it sends it

A mail notice queued in a busy chat could go out stale or twice: it was
kept true from outside the queue, by orchestration check withdrawing it,
which missed other readers, raced an attempt in flight, lost the card
when /clear carried it (a second notice), and moved it to the back when
new mail replaced it. An agent's card behind a person's paused card never
sent after a restart, so an unattended coordinator stalled.

- The drain asks the card's sender, in its own serialized step, whether
  it still stands: send, restate (count and senders of the mail owed
  now, written onto the card in the hand-off's transaction), or withdraw
  as the host. A failing judge sends as written. Send-now restates but
  never withdraws. Orchestration registers the judge on the host.
- check no longer touches any queue (back to main); onMailRead,
  notifyOrchestrationMailRead and the lane's withdrawals are gone.
- One unsettled card per agent message, enforced at insert; new mail is
  counted when the card sends, never replaces or moves it.
- The lane finds its card by what it is (a notice for the mailbox in the
  session the mailbox reaches now); a decline is a withdrawal a person's
  operation stamped. /clear moves agent cards as the host. Pointer rows
  last only while a direct send is in flight and end with the session.
- Pauses (restart, Stop, /clear) hold only a person's cards, and a held
  card never traps an agent's card behind it, keyed on the message kind.
- The source is named for the message (agent-session-message-source),
  one payload shape per message kind, no relative host id.

* fix(orchestration): Orca's mail notice waits in the queue out of sight

The notice card read as something the person typed and came and went on
its own, the transient-state message the queue should not show. It is
hidden until B can label who sent it.

- The host leaves mail-notice cards out of the published queue, so the
  list, its count, Steer/Delete/Edit, steer-newest and the paused header
  never see them, on every client. Keyed on the message kind: D6's task
  card will be shown, labelled, in B. They still count toward the queue
  limit (temporary, until B).
- A refusal no longer returns a hidden card for the person to act on:
  the host withdraws it, and the lane points again only after a turn ran
  since, the rule it already had for a refused notice.
- Send-now no longer restates a card: no one can reach a hidden one.
- The stored sender drops the pane key, a mailbox credential.
- The gate covers an abandoned worker: a mailbox that reaches no session
  withdraws its notice (tested).

* fix(orchestration): a hidden mail notice never waits on the person

Hiding the notice left three paths where it waited on someone who
cannot see it.

- A failed hand-off write put a send_failed hold on it, which only a
  person's Send or Delete clears: that mailbox's notices stopped for
  good. The host now withdraws a hidden card instead and hands its
  mailbox back to the lane through the ordinary redrive, since an idle
  chat gets no other edge. A write that keeps failing is retried once
  per edge.
- A person's Stop while the agent started on the notice put it back to
  waiting, where no pause holds it, so it was sent again at once. The
  host now withdraws it; the lane points again after a turn ran or when
  newer mail makes it a different notice, as on main. A restart still
  puts it back.
- A notice queued before the person's message went first. One order
  rule now: the person's cards go in order and never past one waiting
  on them; a card they cannot see never delays one they can, and goes
  only when none of theirs may.

The drain moves to its own module (structured-agent-session-queued-
drain.ts). Comments that still described the notice as visible are
corrected.

* fix(orchestration): an accepted notice hand-off stamps its mail; a judge that cannot look decides nothing

- When the agent opened its mail in the notice's own turn (the normal
  flow), an open batch left nothing owed, so the lane never marked the
  mail the accepted hand-off carried as delivered, unlike an accepted
  direct send. The lane now stamps exactly the ids a handed-off notice
  carried whenever undelivered mail remains, owed or not. A later
  conversation is no longer told again about mail already opened.
- At quit the host registry is cleared before teardown, so the drain's
  judge could not resolve the mailbox and withdrew the notice. A judge
  that cannot read its inputs now defers: the card waits for a step
  that can look. A mailbox that resolves to another session or none is
  still withdrawn.
- The judge, the owed-batch selection and the notice body move to
  structured-mail-notice.ts. A failing hand-back of a dropped card is
  logged.

* refactor(orchestration): a chat receives the agent's mail itself, queued like a person's message

A structured chat used to get a derived "You have N orchestration messages,
run check" notice, which could go stale while it waited in the chat's queue,
so the queue judged, restated or withdrew it as it sent. It now gets the mail
itself as the turn, through sendAgentTurn with delivery 'queue', the path a
person's message takes.

- One turn per mailbox batch: the mailbox's unread mail at delivery, in mail
  order, each message led by "[message from <sender>]", its type, subject,
  body, payload and reply hint (what check prints). Mail arriving while that
  card is still in the queue waits for the next card; a card is never edited.
- An idle chat takes it at once; a busy one queues it as an ordinary card the
  person sees and can Steer or delete. Mail is marked read when the chat
  accepts the turn, so check does not return it again; a deleted card leaves
  its mail unread for check, and it is not pushed again.
- The card stores who it is from (source_json): every distinct sender, and
  each message's id, run and sender. Kind 'mail-notice' becomes 'mail'.
- Removed: the send-time judge (QueuedAgentCardVerdict), structured-mail-
  notice, structured-pointer-notice-cards, queued-message-restatement, the
  hidden-card rules, the agent-card exemptions from the restart, Stop and
  /clear pauses, the dropped-card hand-back, and the drain split. An agent's
  card now follows the same pauses as a person's.
- Terminal agents are unchanged: they keep the typed pointer and check.

* fix(orchestration): a chat's check skips mail already queued to it; restarts hold only the person's cards

- When a structured chat runs `orchestration check` (consuming, --peek or
  --wait), mail an agent's card in that chat's own queue still carries
  (any card not deleted) is left out: it is on its way as a turn, so the
  agent does not read it twice. Derived from the queued rows at read time;
  a deleted card no longer carries it, so it comes back. Terminal callers
  and every other read are unchanged. New batches pass the exclusion to
  getOrCreateMailboxDelivery; peeks filter it.
- A restart's queue pause now holds only the person's cards: an unattended
  coordinator's agent card sends after an app restart without waiting for a
  Resume, since the mailbox is the record and reading is marked. Stop and
  /clear still hold agent cards. One rule, in queuePauseHolding. An agent
  card queued behind the person's held card still waits behind it: the
  queue never reorders.
- The duplicated direct-mailbox snapshot routing in check-run and
  check-worker becomes one helper, which keeps check-worker under its line
  limit.

* fix(orchestration): the mailbox is the only record of read mail; a chat's check takes its waiting cards

A card held an exclusive claim on its mail that nothing reconciled with the
mailbox: queueing it stamped the mail delivered, and a chat's check hid every
card that was not deleted. So a chat's own check --wait in one turn could not
see a result its waiting card held, a hand-off that ended in doubt or came back
"Not sent" stranded its mail, a check racing the card's build read the mail
twice, and a card left behind by run-use still sent and marked it read.

What a card or send holds is now derived each time from the chat's queue and
its sends: an accepted send marks its mail read; a waiting card or an
unanswered send holds its mail; everything else unread is pushed again or
returned by check. A card the person deleted is recorded as not to be pushed,
at the delete. The lane withdraws returned, partly read and moved cards. A
chat's consuming check withdraws the waiting cards holding what it reads, one
at a time with the lane. The sender is stored on the sent submission too
(host-only), so read state survives the card row's prune. A refusal before
anything started ends its operation; a restart pause is raised only by the
person's cards; an unreadable agent source stays an agent's.

* refactor(orchestration): a chat gets the pointer a terminal agent gets, sent through the chat's own queue

A structured chat now receives exactly what a terminal agent is typed: "You have N orchestration
messages ... run `orca orchestration check`" (formatMessagePointer, same CLI name). It goes
through the shared sendAgentTurn with delivery 'queue', the composer's queue-if-active: an idle
chat gets it at once, a busy chat's queue holds it as a normal card and sends it when the turn
ends. The lane's idle gate is gone; the send decides, as for the person's message.

The lane queues no second pointer while one of its cards still waits, read from the queued rows;
mail arriving meanwhile, or after the card is sent, is pointed again as the terminal lane points
new unread mail. A queued pointer counts as delivered, as an accepted one does. `check` and mail
read state are main's: only `check` reads mail.

The card records who it is from (queued_messages.source_json, kind 'agent', message
'mail-notice', with its senders) for the next PR to render. A restart's pause is raised by and
holds only the person's cards; Stop and /clear still hold an agent's.

Removed from the earlier designs: the message-as-turn batching, the check exclusion and taking,
the derived claims, lock and reconciliation, the sender on submissions, and their tests.

* fix(orchestration): a chat's pointer follows a card that leaves its queue unsent

The lane holds new mail while a pointer card waits, and retried it only on the chat's next
status change. A card that leaves the queue without a turn after it (the person deletes it
while Stop holds it, or its hand-off comes back) made none, so that mail sat unpointed until
something unrelated happened. The shared queue wiring now tells its host when an agent's card
stops waiting, read only when the queue changed, and the runtime redrives that chat's parked
mail. The runtime forwards the new host dep like the others.

Tests: that case end to end, and /clear carrying an agent card keeps who it is from. Agent card
bodies in tests are the pointer text; the restart pause's header says why an agent card may send.

* refactor(orchestration): no special handling for agent notices in the chat queue

A structured chat's orchestration notice is now sent through the chat's own send
like any message: an idle chat takes it as a turn, a busy chat's queue holds it
and sends it at turn end, under the same Stop, restart and /clear pauses as the
person's cards. Native chat no longer branches on who a card is from; the card
only records it.

- Restore main's queue pause logic (no restart exemption for agent cards).
- Remove the queue watch that redrove mail when an agent card left the queue,
  and the host's queued-row read the lane used for it.
- The lane reads nothing of the chat's queue. It sends no second notice while
  mail it already pointed at is unread, read off the mailbox alone.

* fix(orchestration): point new mail like a terminal, even while earlier pointed mail is unread

Drops the native-chat-only rule that held back a notice while the agent had not yet read mail it
was already pointed at. Also fixes main's stop-note test, which still built a provider handle in
the shape #24991 replaced.

* test(native-chat): keep the queued-message rig fixture under the line limit

* fix(test): drop main's duplicate codexProviderHandle import (same line as #25713)
This commit is contained in:
Brennan Benson
2026-10-05 17:34:06 -07:00
committed by GitHub
parent 41f1103c30
commit dd39dcba5f
34 changed files with 1103 additions and 629 deletions
@@ -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<QueuedMessageRow> {
@@ -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'
@@ -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(
@@ -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()
@@ -77,7 +77,8 @@ async function handOffAndReject(
messageId: 'draft-1',
body,
fingerprint,
hostInstance: 'p'
hostInstance: 'p',
source: { kind: 'user' }
})
await journal.appendSubmission(
{
@@ -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')
@@ -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)
);
`)
@@ -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')
@@ -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
@@ -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<typeof sendPreparation>[2]
@@ -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<void>): Promise<void> {
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 }
}
@@ -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())
@@ -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
)
@@ -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, {
@@ -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<TurnOutcome<AgentSessionSendResult>>
): Promise<TurnOutcome<AgentSessionSendResult>> {
@@ -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)
@@ -113,7 +113,7 @@ async function renderPreamble(worker: 'chat' | 'terminal'): Promise<string> {
/** The turn text the structured lane sends a chat for one message on `mailbox`. */
async function renderChatPointer(mailbox: string): Promise<string> {
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<string> {
// 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<string> {
})
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]!
}
@@ -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:<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:<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<
@@ -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' }])
})
@@ -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' } })
})
})
@@ -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<StructuredSessionTurnOutcome> {
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 }
@@ -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' }
]
}
})
})
})
@@ -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<MessageRow, 'id' | 'from_handle' | 'run_id'>
export function structuredMailSource(input: {
db: OrchestrationDb | null
mailboxHandle: string
dispatchId: string | null
batch: readonly MailSourceMessage[]
}): AgentMessageSource {
const senders = new Map<string, AgentMessageSender>()
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) }
}
}
@@ -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<string>()
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<string, StructuredPointerOperationRow>()
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<string, string> = {
'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()
@@ -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<StructuredPointerGateFacts | null>
/** `null` when the session cannot be read. */
readSessionFacts: (sessionId: string) => Promise<StructuredPointerSessionFacts | null>
send: (input: {
sessionId: string
dispatchId: string | null
operationId: string
expectedRuntimeFence: number
body: AgentJournalMessageItem
source: AgentMessageSource
}) => Promise<StructuredPointerSendOutcome>
/** 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<string> | undefined
): Promise<void> {
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;
@@ -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 () => {
@@ -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<StructuredPointerGateFacts | null> {
/** What each recorded send settled as. */
async function readPointerSessionFacts(
sessionId: string
): Promise<StructuredPointerSessionFacts | null> {
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<AgentJournalSnapshot | null> {
@@ -77,8 +77,8 @@ async function readSessionJournal(sessionId: string): Promise<AgentJournalSnapsh
export function createStructuredMailboxPointerHost(): StructuredMailboxPointerHost {
return {
readGateFacts(sessionId) {
return readPointerGateFacts(sessionId)
readSessionFacts(sessionId) {
return readPointerSessionFacts(sessionId)
},
currentFence(sessionId) {
@@ -101,7 +101,9 @@ export function createStructuredMailboxPointerHost(): StructuredMailboxPointerHo
: structuredSessionPointerCallerKey(input.sessionId),
turn: {
body: input.body,
delivery: 'now',
// As a person's message is: a busy chat queues it as a card, sent when the turn ends.
delivery: 'queue',
source: input.source,
operationId: input.operationId,
expectedRuntimeFence: input.expectedRuntimeFence
}
@@ -112,9 +114,7 @@ export function createStructuredMailboxPointerHost(): StructuredMailboxPointerHo
? { kind: 'unattached' }
: { kind: 'sent', state: 'rejected' }
case 'queued':
// Never for a `now` send. A draft would hand off under a fresh id, which this lane's
// operation row cannot see, so reading it needs its own rule before this lane queues.
return { kind: 'sent', state: 'unknown' }
return { kind: 'queued' }
case 'sent': {
// `pending` is not yet an acknowledgement; only `accepted` may consume mail. A send still
// pending after the wait parks for the next journal edge.
@@ -1,7 +1,6 @@
import { describe, expect, it } from 'vitest'
import type { AgentJournalRenderItem } from '../../../shared/agent-session-journal-types'
import {
decideStructuredSessionPointerDelivery,
retainReasonForDispatch,
structuredDispatchDelivered,
structuredSessionGateFacts
@@ -80,37 +79,6 @@ describe('structured session gate facts', () => {
})
})
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)
@@ -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.
*
@@ -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<void>> {
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<void> {
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<string> {
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()
})
})
@@ -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<typeof fakeCodex>
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<string, unknown>,
options: { sessionId?: string } = {}
): Parameters<RpcDispatcher['dispatch']>[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<string, unknown>,
options?: { sessionId?: string }
): Promise<Record<string, unknown>> {
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<FakeConnection> {
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<string, string>()
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<void> {
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<string[]> {
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<void> {
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<string> {
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<void> {
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<string[]> {
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 : ''))
)
}
@@ -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<typeof fakeCodex>
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<string, unknown>,
options: { sessionId?: string } = {}
): Parameters<RpcDispatcher['dispatch']>[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<string, unknown>,
options?: { sessionId?: string }
): Promise<Record<string, unknown>> {
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<FakeConnection> {
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<string, string>()
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<void> {
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<string[]> {
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<void> {
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<string> {
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<void> {
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
})
})
})
@@ -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<OrchestrationPartyIdentity, 'paneKey'> }>
/** 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<OrcaSessionId>(
(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
}
@@ -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:<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:<id>`. */
orcaSessionId: OrcaSessionId | null
}>