feat(native-chat): persist the dispatch wire id with the submission

performSend now mints the id that names a send on the provider wire and
writes it into the write-ahead submission row before anything is sent,
then hands it to the adapter so the Claude dispatch stamps that exact id
on the frame instead of minting its own. Delivery can then be decided by
identity against the provider's own record, rather than inferred from
absence in a content window.

The journal submission row and AgentJournalSubmission gain
providerWireUuid. Rows written before this change project null. The wire
field is optional as well as nullable: a client paired with an older host
sees it absent, which means "this host never evaluated it" and must never
be read as "no id was recorded". No schema version bump is needed or
wanted, because a submission row still serialises at v2 and an older host
keeps reading it instead of latching the journal read-only.

Echo matching and the dispatch waiters already keyed on this id and are
unchanged. Internal sends that carry no durable submission still mint
their own id.
This commit is contained in:
Brennan Benson
2026-09-14 23:34:30 -07:00
parent e6f789b542
commit b5b863c5ec
11 changed files with 286 additions and 6 deletions
@@ -0,0 +1,227 @@
// The dispatch identity contract: the id persisted with a submission IS the id
// that reaches Claude's wire frame, and it is durable BEFORE that frame is
// written. Delivery can then be decided by identity rather than inferred from
// absence in a content window.
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 { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
import { AgentJournalSubmissionSchema } from '../../shared/agent-session-journal-schemas'
import { createTrackedJournalOpener } from '../native-chat/agent-session-journal/journal-store-test-open'
import type { AgentSessionJournal } from '../native-chat/agent-session-journal/journal-store'
import {
applyJournalRow,
createJournalReducerState,
renderJournalState
} from '../native-chat/agent-session-journal/journal-reducer'
import { parseJournalRow } from '../native-chat/agent-session-journal/journal-row-schema'
import type { StructuredAgentSessionAdapter } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import {
performSend,
type AgentSessionTurnContext
} from '../native-chat/agent-session-wire/structured-agent-session-turns'
import { dispatchClaudeTurn, resolveClaudeReplayWaiter } from './claude-structured-dispatch'
import { sessionFor, userMessage, userReplayFrame } from './claude-structured-dispatch-test-support'
const journals = createTrackedJournalOpener()
let root: string
let journal: AgentSessionJournal
beforeEach(async () => {
root = await mkdtemp(join(tmpdir(), 'orca-dispatch-identity-'))
journal = await journals.open({
identity: {
sessionId: 'session-1',
workspaceId: 'workspace-1',
hostId: 'local',
agent: 'claude',
providerHandle: { kind: 'claude', sessionId: 'provider-session', leafUuid: null }
},
journalDir: root
})
})
afterEach(async () => {
await journals.closeAll()
await rm(root, { recursive: true, force: true })
})
const BODY: AgentJournalMessageItem = {
kind: 'message',
role: 'user',
blocks: [{ type: 'text', text: 'deliver me once' }]
}
/** Drives a real journal through `performSend` into the real Claude dispatch,
* capturing what actually reached the provider connection. */
async function sendThroughClaude(): Promise<{
frames: Record<string, unknown>[]
persistedWhenFrameWasWritten: string | null | undefined
}> {
const frames: Record<string, unknown>[] = []
let persistedWhenFrameWasWritten: string | null | undefined
const send = vi.fn(async (frame: Record<string, unknown>) => {
// Read the durable row at the instant of the wire write: this is the
// ordering claim, not just the final value.
persistedWhenFrameWasWritten = journal.submissions()[0]?.providerWireUuid
frames.push(frame)
})
const session = sessionFor(send)
const context: AgentSessionTurnContext = {
sessionId: 'session-1',
journal,
fence: 1,
adapter: {
dispatch: (input) => dispatchClaudeTurn(session, input)
} as unknown as StructuredAgentSessionAdapter,
persistOptions: async () => undefined,
resolvedBy: 'caller',
publish: vi.fn(),
now: () => 1
}
await performSend(context, {
clientMessageId: 'client-1',
payloadFingerprint: 'fingerprint',
body: BODY
})
return { frames, persistedWhenFrameWasWritten }
}
describe('claude dispatch identity persistence', () => {
it('persists the id that reaches the provider frame', async () => {
const { frames } = await sendThroughClaude()
expect(frames).toHaveLength(1)
const wireUuid = frames[0]!.uuid
expect(typeof wireUuid).toBe('string')
// The receipt is only worth anything if it names the frame Claude saw.
expect(journal.submissions()[0]!.providerWireUuid).toBe(wireUuid)
})
it('has the id durable before the frame is written', async () => {
const { frames, persistedWhenFrameWasWritten } = await sendThroughClaude()
expect(persistedWhenFrameWasWritten).toBe(frames[0]!.uuid)
})
it('reopens the journal with the dispatched id intact', async () => {
const { frames } = await sendThroughClaude()
await journal.close()
const reopened = await journals.open({
identity: {
sessionId: 'session-1',
workspaceId: 'workspace-1',
hostId: 'local',
agent: 'claude',
providerHandle: { kind: 'claude', sessionId: 'provider-session', leafUuid: null }
},
journalDir: root
})
// Surviving a restart is the whole point: this is the read the reattach does.
expect(reopened.submissions()[0]!.providerWireUuid).toBe(frames[0]!.uuid)
})
it('still settles the echo that carries the persisted id', async () => {
const frames: Record<string, unknown>[] = []
const send = vi.fn(async (frame: Record<string, unknown>) => {
frames.push(frame)
})
const session = sessionFor(send)
const settled = vi.fn()
const dispatched = dispatchClaudeTurn(session, {
clientMessageId: 'client-1',
body: BODY,
providerWireUuid: 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeee0001'
})
await expect(dispatched).resolves.toEqual({ state: 'admitted' })
// The supplied id goes on the wire unchanged, and echo matching keys on it.
expect(frames[0]!.uuid).toBe('aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeee0001')
expect(session.dispatchWaiters[0]!.sentUuid).toBe('aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeee0001')
expect(
resolveClaudeReplayWaiter(
session,
userReplayFrame('aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeee0001', 'deliver me once'),
settled
)
).toBe(true)
expect(settled).toHaveBeenCalledWith({
clientMessageId: 'client-1',
providerIdentity: {
provider: 'claude',
sessionId: 'provider-session',
uuid: 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeee0001'
}
})
})
it('still dispatches an internal send that carries no submission', async () => {
const frames: Record<string, unknown>[] = []
const send = vi.fn(async (frame: Record<string, unknown>) => {
frames.push(frame)
})
const session = sessionFor(send)
await expect(
dispatchClaudeTurn(session, { body: userMessage([{ type: 'text', text: '/compact' }]) })
).resolves.toEqual({ state: 'admitted' })
expect(typeof frames[0]!.uuid).toBe('string')
})
})
describe('submission rows written before the field existed', () => {
it('parses and projects a missing id as null, never undefined', () => {
// A v2 submission row exactly as an older build wrote it: no such key.
const legacy = {
v: 2,
kind: 'submission',
epoch: 'epoch-1',
seq: 1,
fence: 1,
ts: 1_000,
clientMessageId: 'm-1',
payloadFingerprint: 'a'.repeat(64),
providerHandle: { kind: 'claude', sessionId: 'provider-session', leafUuid: null },
body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'old' }] }
}
const parsed = parseJournalRow(JSON.stringify(legacy))
if (!parsed.ok) {
throw new Error('a row an older build wrote must stay readable')
}
const state = createJournalReducerState('session-1', 'epoch-1')
applyJournalRow(state, parsed.row)
// This host read the row, so "none recorded" is a known answer: null.
expect(renderJournalState(state).submissions[0]!.providerWireUuid).toBeNull()
})
})
describe('the wire projection of the dispatch identity', () => {
const base = {
clientMessageId: 'm-1',
fence: 1,
payloadFingerprint: 'fingerprint',
dispatchState: 'pending',
providerItemId: null,
reason: null,
submittedAt: 1,
resolvedAt: null
}
it('keeps absent and null distinguishable', () => {
// An old host omits the key; a current host answers null. Collapsing the two
// would let "this host never evaluated it" read as "no id was recorded".
const fromOldHost = AgentJournalSubmissionSchema.parse({ ...base })
const recordedNone = AgentJournalSubmissionSchema.parse({ ...base, providerWireUuid: null })
const recorded = AgentJournalSubmissionSchema.parse({ ...base, providerWireUuid: 'wire-1' })
expect('providerWireUuid' in fromOldHost).toBe(false)
expect(recordedNone.providerWireUuid).toBeNull()
expect(recorded.providerWireUuid).toBe('wire-1')
})
})
+10 -2
View File
@@ -255,7 +255,13 @@ export function retireClaudeDispatchWaiters(session: ClaudeSession): void {
export async function dispatchClaudeTurn(
session: ClaudeSession,
input: { clientMessageId?: string; body: AgentJournalMessageItem }
input: {
clientMessageId?: string
body: AgentJournalMessageItem
/** Already persisted with the submission. Minting a second id here would
* leave the durable row naming a message the provider never saw. */
providerWireUuid?: string
}
): Promise<AgentSessionDispatchOutcome> {
let content: unknown[]
try {
@@ -270,7 +276,9 @@ export async function dispatchClaudeTurn(
// Read the sent content, not the journal blocks: only the mapped trailing prompt decides
// whether Claude runs a command, so the two cannot disagree about which frame settles this.
const acceptsResult = claudeDispatchInvokesSlashCommand(content)
const sentUuid = randomUUID()
// Caller-supplied whenever a durable submission backs this dispatch; internal
// sends (compaction) still mint their own, which nothing needs to recover.
const sentUuid = input.providerWireUuid ?? randomUUID()
const replay = waitForReplay(
session,
acceptsResult,
@@ -248,7 +248,10 @@ function applySubmission(
providerItemId: null,
reason: null,
submittedAt: row.ts,
resolvedAt: null
resolvedAt: null,
// A row written before the field existed carries no id, which is `null` —
// this host read the row, so the answer is known to be "none".
providerWireUuid: row.providerWireUuid ?? null
})
const itemId = agentJournalSubmissionKey(row.clientMessageId)
upsertItem(state, itemId, 0, {
@@ -58,6 +58,7 @@ export function journalSubmissionRowBuilder(
payloadFingerprint: string
body: AgentJournalMessageItem
fence: number
providerWireUuid?: string | null
}
): RowBuilder<JournalSubmissionRow> {
return (seq, ts) =>
@@ -206,6 +207,7 @@ export function buildJournalSubmissionRow(input: {
seq: number
fence: number
ts: number
providerWireUuid?: string | null
}): JournalSubmissionRow {
return {
kind: 'submission',
@@ -213,6 +215,7 @@ export function buildJournalSubmissionRow(input: {
payloadFingerprint: input.payloadFingerprint,
providerHandle: input.providerHandle,
body: input.body,
providerWireUuid: input.providerWireUuid ?? null,
...journalRowBase(input.state.epoch, input.seq, input.fence, input.ts)
}
}
@@ -69,6 +69,10 @@ export type JournalSubmissionRow = JournalRowBase & {
payloadFingerprint: string
providerHandle: AgentSessionProviderHandle
body: AgentJournalMessageItem
/** Id stamped on the dispatched provider frame, written here BEFORE the wire
* write so reattach can match delivery by identity. Absent on rows written
* before the field existed; those project `null`. */
providerWireUuid?: string | null
}
export type JournalDispatchRow = JournalRowBase & {
@@ -205,6 +209,7 @@ function isJournalRow(record: Record<string, unknown>): record is JournalRow {
record.clientMessageId.length > 0 &&
typeof record.payloadFingerprint === 'string' &&
isPlainObject(record.providerHandle) &&
isOptionalWireUuid(record.providerWireUuid) &&
isAdmissibleAgentJournalMessageBody(record.body)
)
}
@@ -232,6 +237,12 @@ function isJournalRow(record: Record<string, unknown>): record is JournalRow {
return typeof record.reason === 'string' && isPlainObject(record.providerHandle)
}
/** Absent on every row predating the field, null when the dispatch recorded no
* wire id. Neither is a malformed row. */
function isOptionalWireUuid(value: unknown): boolean {
return value === undefined || value === null || typeof value === 'string'
}
function isLifecycleMutation(value: unknown): value is JournalLifecycleMutation {
if (!isPlainObject(value) || typeof value.itemId !== 'string') {
return false
@@ -54,6 +54,8 @@ export type JournalSubmissionInput = {
payloadFingerprint: string
body: AgentJournalMessageItem
fence: number
/** Id the caller will stamp on the provider frame, recorded before it is sent. */
providerWireUuid?: string | null
}
export type JournalItemAppendInput = {
@@ -160,6 +160,10 @@ export type StructuredAgentSessionAdapter = {
clientMessageId: string
body: AgentJournalMessageItem
fence: number
/** Already durable in the submission row. An adapter that puts an id on the
* provider frame must use THIS one, not mint its own, or the persisted
* receipt names a message the provider never saw. */
providerWireUuid?: string
}): Promise<AgentSessionDispatchOutcome>
rewindSupport?(sessionId: string): AgentSessionRewindSupport
recoverRewind?(input: {
@@ -84,6 +84,9 @@ export function sendPlan(params: {
reason: DISPATCH_DOUBT_SUBMISSION_MISSING,
submittedAt: resolvedAt,
resolvedAt,
// The row is gone, so this host has no id to report — `null`, which is
// a recorded answer, not the absence an older host would send.
providerWireUuid: null,
recovered: true
}
}
@@ -6,6 +6,7 @@
// row the next attach settles as `unknown`, whereas the reverse would lose a
// turn the provider already accepted.
import { randomUUID } from 'node:crypto'
import type {
AgentJournalMessageItem,
AgentJournalSubmission
@@ -50,13 +51,15 @@ function invalid(message: string): { ok: false; refusal: AgentSessionWireRefusal
async function dispatchSafely(
ctx: AgentSessionTurnContext,
clientMessageId: string,
body: AgentJournalMessageItem
body: AgentJournalMessageItem,
providerWireUuid: string
): Promise<AgentSessionDispatchOutcome> {
try {
return await ctx.adapter.dispatch({
sessionId: ctx.sessionId,
clientMessageId,
body,
providerWireUuid,
fence: ctx.fence
})
} catch (error) {
@@ -105,14 +108,19 @@ export async function performSend(
value: { clientMessageId: input.clientMessageId, submission: existing }
}
}
// Minted HERE and durable before the dispatch below, so the id naming this
// send on the wire is recoverable after a crash. Reattach can then match the
// provider's own record by identity instead of inferring delivery from
// absence in a content window.
const providerWireUuid = randomUUID()
try {
await ctx.journal.appendSubmission({ ...input, fence: ctx.fence })
await ctx.journal.appendSubmission({ ...input, providerWireUuid, fence: ctx.fence })
} catch {
return invalid('The message could not be recorded and was not sent.')
}
ctx.publish()
const outcome = await dispatchSafely(ctx, input.clientMessageId, input.body)
const outcome = await dispatchSafely(ctx, input.clientMessageId, input.body, providerWireUuid)
// An admission needs no dispatch row: the submission is already pending.
if (outcome.state === 'admitted') {
ctx.publish()
@@ -212,6 +212,9 @@ export const AgentJournalSubmissionSchema = z.object({
reason: z.string().nullable(),
submittedAt: z.number(),
resolvedAt: z.number().nullable(),
// Optional AND nullable on purpose: an old host omits the key entirely, which
// is not the same answer as a host that recorded no id.
providerWireUuid: z.string().nullable().optional(),
recovered: z.literal(true).optional()
})
@@ -246,6 +246,14 @@ export type AgentJournalSubmission = {
reason: string | null
submittedAt: number
resolvedAt: number | null
/** Id stamped on the dispatched provider frame, durable before the wire write.
* Three states, and collapsing them loses the only thing it is good for:
* a string is the id this send was dispatched with; `null` is this host
* recording that the row carries none; ABSENT is a host that predates the
* field, which knows nothing either way. Absence must never be read as
* "no id recorded, therefore decidable". Whether the provider adopted the id
* is a separate observed fact, not implied by its presence here. */
providerWireUuid?: string | null
/** Set when crash reconciliation resolved the dispatch, not the provider. A live
* `unknown` is a send still outstanding; a recovered one outlived its writer. */
recovered?: true