From b5b863c5ece8ae228ca109fdb7c0be10bdce96ed Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Mon, 14 Sep 2026 23:29:21 -0700 Subject: [PATCH] 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. --- ...laude-structured-dispatch-identity.test.ts | 227 ++++++++++++++++++ src/main/claude/claude-structured-dispatch.ts | 12 +- .../agent-session-journal/journal-reducer.ts | 5 +- .../journal-row-builders.ts | 3 + .../journal-row-schema.ts | 11 + .../journal-store-contracts.ts | 2 + .../structured-agent-session-adapter.ts | 4 + ...structured-agent-session-mutation-plans.ts | 3 + .../structured-agent-session-turns.ts | 14 +- src/shared/agent-session-journal-schemas.ts | 3 + src/shared/agent-session-journal-types.ts | 8 + 11 files changed, 286 insertions(+), 6 deletions(-) create mode 100644 src/main/claude/claude-structured-dispatch-identity.test.ts diff --git a/src/main/claude/claude-structured-dispatch-identity.test.ts b/src/main/claude/claude-structured-dispatch-identity.test.ts new file mode 100644 index 00000000000..8ecc74d41cd --- /dev/null +++ b/src/main/claude/claude-structured-dispatch-identity.test.ts @@ -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[] + persistedWhenFrameWasWritten: string | null | undefined +}> { + const frames: Record[] = [] + let persistedWhenFrameWasWritten: string | null | undefined + const send = vi.fn(async (frame: Record) => { + // 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[] = [] + const send = vi.fn(async (frame: Record) => { + 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[] = [] + const send = vi.fn(async (frame: Record) => { + 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') + }) +}) diff --git a/src/main/claude/claude-structured-dispatch.ts b/src/main/claude/claude-structured-dispatch.ts index 45dc3ebc0e3..c1fb66d410a 100644 --- a/src/main/claude/claude-structured-dispatch.ts +++ b/src/main/claude/claude-structured-dispatch.ts @@ -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 { 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, diff --git a/src/main/native-chat/agent-session-journal/journal-reducer.ts b/src/main/native-chat/agent-session-journal/journal-reducer.ts index 5c4917bec46..af1e53a6043 100644 --- a/src/main/native-chat/agent-session-journal/journal-reducer.ts +++ b/src/main/native-chat/agent-session-journal/journal-reducer.ts @@ -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, { diff --git a/src/main/native-chat/agent-session-journal/journal-row-builders.ts b/src/main/native-chat/agent-session-journal/journal-row-builders.ts index be2c2552775..4d1176f818d 100644 --- a/src/main/native-chat/agent-session-journal/journal-row-builders.ts +++ b/src/main/native-chat/agent-session-journal/journal-row-builders.ts @@ -58,6 +58,7 @@ export function journalSubmissionRowBuilder( payloadFingerprint: string body: AgentJournalMessageItem fence: number + providerWireUuid?: string | null } ): RowBuilder { 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) } } diff --git a/src/main/native-chat/agent-session-journal/journal-row-schema.ts b/src/main/native-chat/agent-session-journal/journal-row-schema.ts index 7dc02dd197b..cfb021e4c87 100644 --- a/src/main/native-chat/agent-session-journal/journal-row-schema.ts +++ b/src/main/native-chat/agent-session-journal/journal-row-schema.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): 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): 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 diff --git a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts index 80c806b02e3..60124d96dd0 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts @@ -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 = { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts index 54d3c15ae0e..67fb2419a72 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter.ts @@ -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 rewindSupport?(sessionId: string): AgentSessionRewindSupport recoverRewind?(input: { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts index 96c50036701..7aab5c7c9ba 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts @@ -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 } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts index 4c275ecd738..bf8400f2dd6 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts @@ -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 { 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() diff --git a/src/shared/agent-session-journal-schemas.ts b/src/shared/agent-session-journal-schemas.ts index eaee6bb9a63..7125c38063c 100644 --- a/src/shared/agent-session-journal-schemas.ts +++ b/src/shared/agent-session-journal-schemas.ts @@ -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() }) diff --git a/src/shared/agent-session-journal-types.ts b/src/shared/agent-session-journal-types.ts index 7ba5b277e6d..548f3eead79 100644 --- a/src/shared/agent-session-journal-types.ts +++ b/src/shared/agent-session-journal-types.ts @@ -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