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