diff --git a/src/main/claude/claude-structured-history-window.test.ts b/src/main/claude/claude-structured-history-window.test.ts new file mode 100644 index 00000000000..4a3884d915e --- /dev/null +++ b/src/main/claude/claude-structured-history-window.test.ts @@ -0,0 +1,189 @@ +// The Claude half of restart reconciliation: which transcript records become +// evidence, and when the read may be called boundary-consistent at all. + +import { describe, expect, it } from 'vitest' +import { structuredAgentSessionSendBody } from '../../shared/structured-agent-session-outbox' +import { structuredAgentSessionPayloadFingerprint } from '../../shared/structured-agent-session-mutation' +import { computeAgentSessionPayloadFingerprint } from '../../shared/agent-session-mutation-envelope' +import { claudeProviderHistoryWindowFromJsonl } from './claude-structured-history-window' + +const PROVIDER_SESSION = 'provider-1' +const ORCA_SESSION = 'session-1' + +type Row = Record + +function prompt(uuid: string, parentUuid: string | null, content: unknown, extra: Row = {}): Row { + return { + type: 'user', + uuid, + parentUuid, + sessionId: PROVIDER_SESSION, + message: { role: 'user', content }, + ...extra + } +} + +function jsonl(rows: Row[], leafUuid: string): string { + const lines = [...rows, { type: 'last-prompt', sessionId: PROVIDER_SESSION, leafUuid }] + return `${lines.map((row) => JSON.stringify(row)).join('\n')}\n` +} + +function read(contents: string, previousLeafUuid: string | null, turnInFlight = false) { + return claudeProviderHistoryWindowFromJsonl({ + contents, + providerSessionId: PROVIDER_SESSION, + previousLeafUuid, + sessionId: ORCA_SESSION, + turnInFlight + }) +} + +/** The digest the submission row carries for a plain typed send. */ +function sendFingerprint(text: string): string { + return structuredAgentSessionPayloadFingerprint({ + method: 'agentSession.send', + sessionId: ORCA_SESSION, + fields: { body: structuredAgentSessionSendBody(text, []) } + }) +} + +const ANCHOR = prompt('anchor', null, 'earlier turn') + +describe('claudeProviderHistoryWindowFromJsonl', () => { + it('pins the renderer and host fingerprint functions to the same digest', () => { + // The renderer computes a send's fingerprint with one, the host admission gate + // validates it with the other, and the window matches with the host's. A + // divergence would refuse every send long before it reached here — but it + // would also silently turn every reconciliation into `not_delivered`. + const input = { + method: 'agentSession.send', + sessionId: ORCA_SESSION, + fields: { body: structuredAgentSessionSendBody('ship it', []) } + } + + expect(structuredAgentSessionPayloadFingerprint(input)).toBe( + computeAgentSessionPayloadFingerprint(input) + ) + }) + + it('fingerprints a prompt after the anchor exactly as the send that produced it', () => { + const contents = jsonl( + [ANCHOR, prompt('u-1', 'anchor', [{ type: 'text', text: 'ship it' }])], + 'u-1' + ) + + const window = read(contents, 'anchor') + + expect(window.boundaryConsistent).toBe(true) + expect(window.items).toEqual([ + { + providerItemId: 'u-1', + clientMessageId: null, + payloadFingerprint: sendFingerprint('ship it'), + identity: { provider: 'claude', sessionId: PROVIDER_SESSION, uuid: 'u-1' } + } + ]) + }) + + it('fingerprints a string-content prompt the same as a block-content one', () => { + const asString = read(jsonl([ANCHOR, prompt('u-1', 'anchor', 'ship it')], 'u-1'), 'anchor') + + expect(asString.items[0]?.payloadFingerprint).toBe(sendFingerprint('ship it')) + }) + + it('excludes everything before the anchor', () => { + const contents = jsonl( + [ + prompt('root', null, 'first'), + prompt('anchor', 'root', 'second'), + prompt('u-1', 'anchor', 'third') + ], + 'u-1' + ) + + expect(read(contents, 'anchor').items.map((item) => item.providerItemId)).toEqual(['u-1']) + }) + + it('reports no window and an inconsistent boundary without a durable anchor', () => { + const contents = jsonl([ANCHOR, prompt('u-1', 'anchor', 'ship it')], 'u-1') + + expect(read(contents, null)).toEqual({ + items: [], + boundaryConsistent: false, + turnInFlight: false + }) + }) + + it('reports an inconsistent boundary when the anchor is gone from the file', () => { + // What a compaction or a fresh session file leaves behind. + const contents = jsonl([prompt('u-1', null, 'ship it')], 'u-1') + + expect(read(contents, 'anchor').boundaryConsistent).toBe(false) + }) + + it('reports an inconsistent boundary when the leaf is on a sibling branch', () => { + const contents = jsonl( + [ + prompt('root', null, 'first'), + prompt('anchor', 'root', 'second'), + prompt('u-1', 'root', 'branched') + ], + 'u-1' + ) + + expect(read(contents, 'anchor').boundaryConsistent).toBe(false) + }) + + it('reports an inconsistent boundary on a torn tail', () => { + const contents = `${jsonl([ANCHOR, prompt('u-1', 'anchor', 'ship it')], 'u-1')}{"type":"user"` + + expect(read(contents, 'anchor').boundaryConsistent).toBe(false) + }) + + it('keeps the boundary consistent and the window empty when nothing followed the anchor', () => { + expect(read(jsonl([ANCHOR], 'anchor'), 'anchor')).toEqual({ + items: [], + boundaryConsistent: true, + turnInFlight: false + }) + }) + + it('excludes harness-injected turns, meta turns, tool results and sidechains', () => { + const contents = jsonl( + [ + ANCHOR, + prompt('u-reminder', 'anchor', [ + { type: 'text', text: 'be careful' } + ]), + prompt('u-meta', 'u-reminder', [{ type: 'text', text: 'injected' }], { isMeta: true }), + prompt('u-tool', 'u-meta', [{ type: 'tool_result', content: 'ok' }]), + prompt('u-real', 'u-tool', 'ship it') + ], + 'u-real' + ) + + expect(read(contents, 'anchor').items.map((item) => item.providerItemId)).toEqual(['u-real']) + }) + + it('excludes a prompt carrying an image, whose path the transcript does not keep', () => { + const contents = jsonl( + [ + ANCHOR, + prompt('u-img', 'anchor', [ + { type: 'text', text: 'look at this' }, + { type: 'image', source: { type: 'base64', media_type: 'image/png', data: 'AAAA' } } + ]) + ], + 'u-img' + ) + + expect(read(contents, 'anchor').items).toEqual([]) + expect(read(contents, 'anchor').boundaryConsistent).toBe(true) + }) + + it('carries the caller-proven turn-in-flight fact through to the window', () => { + const contents = jsonl([ANCHOR, prompt('u-1', 'anchor', 'ship it')], 'u-1') + + expect(read(contents, 'anchor', true).turnInFlight).toBe(true) + }) +}) diff --git a/src/main/claude/claude-structured-history-window.ts b/src/main/claude/claude-structured-history-window.ts new file mode 100644 index 00000000000..9e1a898ba45 --- /dev/null +++ b/src/main/claude/claude-structured-history-window.ts @@ -0,0 +1,256 @@ +// Provider history for restart reconciliation, read from the Claude project JSONL. +// +// Why this file is the source of truth for "did Claude take it": a resume replays +// this transcript by session id, so a message absent from it is absent from the +// conversation Orca is about to resume. Absence here is not an inference about a +// dead child — it is the content of the next turn's context. +// +// The window is anchored on the leaf uuid Orca durably recorded for the session. +// Without that anchor the read has no proven start, and the branch proof is what +// decides whether the file we just read still descends from it: a fork, a +// compaction, a sibling branch, or a torn tail all fail the proof, and every one +// of those makes absence meaningless. Failing it reports an inconsistent +// boundary rather than an empty window, because the two decide opposite things. + +import { readFile, stat } from 'node:fs/promises' +import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types' +import { resolveSessionFilePath } from '../native-chat/session-file-resolver' +import type { + ProviderHistoryItem, + ProviderHistoryWindow +} from '../native-chat/agent-session-journal/journal-submission-reconciler' +import { claudeContentBlocks } from '../native-chat/transcript-record-blocks' +import { isKnownHarnessInjectedUserTurnText } from '../../shared/harness-injected-user-turns' +import { computeAgentSessionPayloadFingerprint } from '../../shared/agent-session-mutation-envelope' +import type { NativeChatBlock } from '../../shared/native-chat-types' +import { proveClaudeTranscriptBranchFromJsonl } from './claude-transcript-branch-proof' + +/** Matches the legacy-import bound: a prefix read would make absence meaningless. */ +const MAX_HISTORY_WINDOW_SOURCE_BYTES = 16 * 1024 * 1024 +const MAX_HISTORY_WINDOW_WALK = 10_000 + +const INCONSISTENT: ProviderHistoryWindow = { + items: [], + boundaryConsistent: false, + turnInFlight: false +} + +type TranscriptRecord = Record + +function stringField(source: unknown, key: string): string | null { + if (!source || typeof source !== 'object') { + return null + } + const value = (source as Record)[key] + return typeof value === 'string' && value.trim() ? value : null +} + +function isPlainTextPart(part: unknown): boolean { + if (typeof part === 'string') { + return true + } + return Boolean(part) && typeof part === 'object' && (part as TranscriptRecord).type === 'text' +} + +/** + * A user record Orca could itself have submitted. Everything the harness injects + * is excluded, because a fingerprint computed over machinery would claim a + * history slot the user's message should have had. + * + * Attachments are excluded too, and deliberately: the transcript keeps a pasted + * image as base64 with no path, so the `image-ref` block the submission was + * fingerprinted from cannot be reconstructed. Measured over 2,764 genuine prompt + * records in 12 local transcripts, every multi-block prompt was exactly + * text + image; none carried injected content. + */ +function claudePromptBlocks(record: TranscriptRecord): NativeChatBlock[] | null { + if ( + record.type !== 'user' || + record.isSidechain === true || + record.parent_tool_use_id != null || + record.isMeta === true || + record.isSynthetic === true || + record.isCompactSummary === true + ) { + return null + } + const message = record.message + const content = + message && typeof message === 'object' ? (message as TranscriptRecord).content : undefined + // Read the RAW parts, not the decoded ones: a base64 image decodes to nothing + // at all, so a prompt with an attachment would otherwise pass as text-only and + // be fingerprinted as if the attachment had never been sent. + if (Array.isArray(content) && !content.every((part) => isPlainTextPart(part))) { + return null + } + const blocks = claudeContentBlocks(content) + if (blocks.length === 0 || blocks.some((block) => block.type !== 'text')) { + return null + } + const [first] = blocks + if (first?.type !== 'text' || isKnownHarnessInjectedUserTurnText(first.text)) { + return null + } + return blocks +} + +/** + * The digest the submission row is GUARANTEED to carry. `admitAndRunAgentSessionMutation` + * recomputes this exact call over the send's own body and refuses the send on a + * mismatch, and `performSend` is the only writer of a submission row — so the + * stored fingerprint is this function's output over the stored body, whoever + * produced the envelope. Matching here is therefore an equality between two runs + * of one function, not a guess about two encodings agreeing. + */ +function promptFingerprint(sessionId: string, blocks: NativeChatBlock[]): string { + return computeAgentSessionPayloadFingerprint({ + method: 'agentSession.send', + sessionId, + fields: { body: { kind: 'message', role: 'user', blocks } } + }) +} + +function indexRecords(contents: string): Map { + const byUuid = new Map() + for (const line of contents.split('\n')) { + if (!line.trim()) { + continue + } + let parsed: unknown + try { + parsed = JSON.parse(line) + } catch { + continue + } + if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) { + continue + } + const record = parsed as TranscriptRecord + const uuid = stringField(record, 'uuid') + if (uuid && !byUuid.has(uuid)) { + byUuid.set(uuid, record) + } + } + return byUuid +} + +/** Records strictly after the anchor, oldest first. The branch proof already + * established that this walk reaches the anchor. */ +function walkFromLeaf( + byUuid: Map, + leafUuid: string, + anchorUuid: string +): TranscriptRecord[] { + const collected: TranscriptRecord[] = [] + let cursor: string | null = leafUuid + for (let depth = 0; cursor !== null && cursor !== anchorUuid; depth += 1) { + if (depth >= MAX_HISTORY_WINDOW_WALK) { + return [] + } + const record = byUuid.get(cursor) + if (!record) { + return [] + } + collected.push(record) + cursor = stringField(record, 'parentUuid') + } + return collected.toReversed() +} + +export function claudeProviderHistoryWindowFromJsonl(input: { + contents: string + providerSessionId: string + previousLeafUuid: string | null + /** Orca session id: the fingerprint a submission carries is scoped to it. */ + sessionId: string + /** The caller must PROVE no provider child can be appending; absence proves + * nothing while a turn is running. */ + turnInFlight: boolean +}): ProviderHistoryWindow { + if (!input.previousLeafUuid) { + return INCONSISTENT + } + let leafUuid: string + try { + leafUuid = proveClaudeTranscriptBranchFromJsonl({ + contents: input.contents, + providerSessionId: input.providerSessionId, + previousLeafUuid: input.previousLeafUuid + }).leafUuid + } catch { + // Every failure mode here — missing ancestor, sibling branch, compacted + // cursor, torn tail — is a boundary we cannot vouch for. + return INCONSISTENT + } + const byUuid = indexRecords(input.contents) + const items: ProviderHistoryItem[] = [] + for (const record of walkFromLeaf(byUuid, leafUuid, input.previousLeafUuid)) { + const blocks = claudePromptBlocks(record) + const uuid = stringField(record, 'uuid') + if (!blocks || !uuid) { + continue + } + items.push({ + providerItemId: uuid, + // Claude echoes no client message id, so identity matching reduces to the + // fingerprint pass; the reconciler treats that as the weakest evidence. + clientMessageId: null, + payloadFingerprint: promptFingerprint(input.sessionId, blocks), + identity: { + provider: 'claude', + sessionId: stringField(record, 'sessionId') ?? input.providerSessionId, + uuid + } + }) + } + return { items, boundaryConsistent: true, turnInFlight: input.turnInFlight } +} + +/** + * The window for one attached session: resolve the provider's transcript, then + * read it against the handle's durable leaf. A live child means a send queued + * behind its running turn is not in the file yet, so liveness is carried in + * rather than assumed — only the adapter's session map can answer it. + */ +export async function resolveClaudeProviderHistoryWindow(input: { + identity: AgentSessionJournalIdentity + hasLiveSession: boolean +}): Promise { + const handle = input.identity.providerHandle + if (handle.kind !== 'claude') { + return null + } + const transcriptPath = await resolveSessionFilePath('claude', handle.sessionId) + if (!transcriptPath) { + return null + } + return readClaudeProviderHistoryWindow({ + transcriptPath, + providerSessionId: handle.sessionId, + previousLeafUuid: handle.leafUuid, + sessionId: input.identity.sessionId, + turnInFlight: input.hasLiveSession + }) +} + +export async function readClaudeProviderHistoryWindow(input: { + transcriptPath: string + providerSessionId: string + previousLeafUuid: string | null + sessionId: string + turnInFlight: boolean +}): Promise { + if (!input.previousLeafUuid) { + return INCONSISTENT + } + let contents: string + try { + if ((await stat(input.transcriptPath)).size > MAX_HISTORY_WINDOW_SOURCE_BYTES) { + return INCONSISTENT + } + contents = await readFile(input.transcriptPath, 'utf8') + } catch { + return INCONSISTENT + } + return claudeProviderHistoryWindowFromJsonl({ ...input, contents }) +} diff --git a/src/main/claude/claude-structured-session-adapter.ts b/src/main/claude/claude-structured-session-adapter.ts index d613a175917..5eab424cab5 100644 --- a/src/main/claude/claude-structured-session-adapter.ts +++ b/src/main/claude/claude-structured-session-adapter.ts @@ -32,6 +32,8 @@ import { } from './claude-structured-session-close' import { readClaudeTranscriptLeafWithReproof } from './claude-transcript-branch-proof' import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire' +import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types' +import { resolveClaudeProviderHistoryWindow } from './claude-structured-history-window' export type { ClaudeStructuredLaunch } from './claude-structured-launch-resolution' export type { @@ -164,6 +166,14 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda return exit.settlementPromise } + /** Restart reconciliation reads the transcript a resume replays; only this map + * can say whether a child is still appending to it. */ + providerHistoryWindow = (input: { identity: AgentSessionJournalIdentity }) => + resolveClaudeProviderHistoryWindow({ + identity: input.identity, + hasLiveSession: this.sessions.has(input.identity.sessionId) + }) + private async persistSessionHandle(sessionId: string, session: ClaudeSession): Promise { try { const transcriptLeaf = this.deps.readTranscriptLeaf diff --git a/src/main/native-chat/agent-session-journal/journal-restart-reconciliation.test.ts b/src/main/native-chat/agent-session-journal/journal-restart-reconciliation.test.ts new file mode 100644 index 00000000000..05b7e7b4cce --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-restart-reconciliation.test.ts @@ -0,0 +1,225 @@ +// Wiring the restart reconciler: what provider history is allowed to decide +// about a submission the crash boundary could only doubt. + +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { agentJournalItemKey } from '../../../shared/agent-session-journal-item-key' +import type { + AgentJournalItemIdentity, + AgentJournalMessageItem, + AgentSessionJournalIdentity +} from '../../../shared/agent-session-journal-types' +import { digestPayload } from './journal-payload-bounds' +import { reconcileJournalSubmissionsAgainstHistory } from './journal-restart-reconciliation' +import type { ProviderHistoryItem, ProviderHistoryWindow } from './journal-submission-reconciler' +import { createTrackedJournalOpener } from './journal-store-test-open' + +const IDENTITY: AgentSessionJournalIdentity = { + sessionId: 'session-1', + workspaceId: 'ws-1', + hostId: 'host-1', + agent: 'claude', + providerHandle: { kind: 'claude', sessionId: 'provider-1', leafUuid: null } +} + +function claudeIdentity(uuid: string): AgentJournalItemIdentity { + return { provider: 'claude', sessionId: 'provider-1', uuid } +} + +let root: string +let clock = 1_000 + +function tick(): number { + clock += 1 + return clock +} + +function userMessage(text: string): AgentJournalMessageItem { + return { kind: 'message', role: 'user', blocks: [{ type: 'text', text }] } +} + +const journals = createTrackedJournalOpener() + +async function open() { + return journals.open({ + identity: IDENTITY, + journalDir: root, + now: tick, + mintEpoch: () => `epoch-${clock}` + }) +} + +function history(uuid: string, text: string): ProviderHistoryItem { + return { + providerItemId: uuid, + clientMessageId: null, + payloadFingerprint: digestPayload(text), + identity: claudeIdentity(uuid) + } +} + +function window( + items: ProviderHistoryItem[], + overrides: Partial = {} +): ProviderHistoryWindow { + return { items, boundaryConsistent: true, turnInFlight: false, ...overrides } +} + +/** A host that wrote the submission row and died before learning its outcome. */ +async function reopenAfterCrash( + body: AgentJournalMessageItem = userMessage('deploy the thing'), + text = 'deploy the thing' +) { + const journal = await open() + await journal.appendSubmission({ + clientMessageId: 'cm_1', + payloadFingerprint: digestPayload(text), + body, + fence: 1 + }) + const restarted = await open() + await restarted.markPendingSubmissionsUnknown(2) + return restarted +} + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-journal-restart-reconcile-')) + clock = 1_000 +}) + +afterEach(async () => { + await journals.closeAll() + await rm(root, { recursive: true, force: true }) +}) + +describe('reconcileJournalSubmissionsAgainstHistory', () => { + it('settles a message found in provider history as accepted on its provider identity', async () => { + const journal = await reopenAfterCrash() + + const settled = await reconcileJournalSubmissionsAgainstHistory({ + journal, + fence: 2, + history: window([history('uuid-1', 'deploy the thing')]) + }) + + expect(settled).toEqual(['cm_1']) + const submission = journal.submissions()[0] + expect(submission?.dispatchState).toBe('accepted') + expect(submission?.providerItemId).toBe(agentJournalItemKey(claudeIdentity('uuid-1'))) + }) + + it('settles a message provably absent from history as rejected: not_delivered', async () => { + const journal = await reopenAfterCrash() + + const settled = await reconcileJournalSubmissionsAgainstHistory({ + journal, + fence: 2, + history: window([]) + }) + + expect(settled).toEqual(['cm_1']) + const submission = journal.submissions()[0] + expect(submission?.dispatchState).toBe('rejected') + expect(submission?.reason).toBe('not_delivered') + }) + + it('leaves a submission unknown while the provider reports a turn in flight', async () => { + const journal = await reopenAfterCrash() + + const settled = await reconcileJournalSubmissionsAgainstHistory({ + journal, + fence: 2, + history: window([], { turnInFlight: true }) + }) + + expect(settled).toEqual([]) + expect(journal.submissions()[0]?.dispatchState).toBe('unknown') + }) + + it('leaves a submission unknown when the history boundary is inconsistent', async () => { + const journal = await reopenAfterCrash() + + const settled = await reconcileJournalSubmissionsAgainstHistory({ + journal, + fence: 2, + history: window([], { boundaryConsistent: false }) + }) + + expect(settled).toEqual([]) + expect(journal.submissions()[0]?.dispatchState).toBe('unknown') + }) + + it('refuses to reject a submission carrying an attachment it cannot fingerprint', async () => { + const journal = await reopenAfterCrash( + { + kind: 'message', + role: 'user', + blocks: [ + { type: 'text', text: 'look at this' }, + { type: 'image-ref', path: '/tmp/shot.png' } + ] + }, + 'look at this' + ) + + const settled = await reconcileJournalSubmissionsAgainstHistory({ + journal, + fence: 2, + history: window([]) + }) + + expect(settled).toEqual([]) + expect(journal.submissions()[0]?.dispatchState).toBe('unknown') + }) + + it('does not let an item the journal already committed stand in for a new send', async () => { + const journal = await open() + // An identical message, delivered and committed BEFORE the one that crashed. + await journal.appendItem(claudeIdentity('uuid-old'), userMessage('deploy the thing'), { + fence: 1 + }) + await journal.appendSubmission({ + clientMessageId: 'cm_1', + payloadFingerprint: digestPayload('deploy the thing'), + body: userMessage('deploy the thing'), + fence: 1 + }) + const restarted = await open() + await restarted.markPendingSubmissionsUnknown(2) + + await reconcileJournalSubmissionsAgainstHistory({ + journal: restarted, + fence: 2, + history: window([history('uuid-old', 'deploy the thing')]) + }) + + expect(restarted.submissions()[0]?.dispatchState).toBe('rejected') + }) + + it('leaves two identical unsettled sends unknown rather than guessing between them', async () => { + const journal = await open() + for (const id of ['cm_1', 'cm_2']) { + await journal.appendSubmission({ + clientMessageId: id, + payloadFingerprint: digestPayload('ping'), + body: userMessage('ping'), + fence: 1 + }) + } + const restarted = await open() + await restarted.markPendingSubmissionsUnknown(2) + + await reconcileJournalSubmissionsAgainstHistory({ + journal: restarted, + fence: 2, + history: window([history('uuid-1', 'ping'), history('uuid-2', 'ping')]) + }) + + expect(restarted.submissions().map((entry) => entry.dispatchState)).toEqual([ + 'unknown', + 'unknown' + ]) + }) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-restart-reconciliation.ts b/src/main/native-chat/agent-session-journal/journal-restart-reconciliation.ts new file mode 100644 index 00000000000..4df3e30060e --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-restart-reconciliation.ts @@ -0,0 +1,106 @@ +// The production caller for `reconcileSubmissions`. +// +// Runs once per journal open, after the crash boundary has already settled every +// survivor to `unknown`. It only ever narrows that answer: `accepted` when the +// provider's own history holds the message, `rejected` when a boundary we can +// vouch for proves it never arrived. Anything the reconciler leaves `unknown` +// is left exactly as the crash boundary wrote it. +// +// Nothing here dispatches. A `rejected` submission becomes re-sendable only +// through the user's Retry, which rotates the client message id; Orca still +// never puts a message back on the wire on the user's behalf. + +import type { + AgentJournalMessageItem, + AgentJournalSubmission +} from '../../../shared/agent-session-journal-types' +import { + agentJournalItemKey, + agentJournalSubmissionKey +} from '../../../shared/agent-session-journal-item-key' +import type { AgentSessionJournal } from './journal-store' +import { reconcileSubmissions, type ProviderHistoryWindow } from './journal-submission-reconciler' + +/** + * Only a text-only body can be compared against provider content. A submission + * carrying an attachment was fingerprinted over an `image-ref` path the + * transcript does not keep, so its absence from history would be an artefact of + * the encoding rather than evidence — and `rejected` is the one outcome that + * costs the user a duplicate if it is wrong. Those stay `unknown`. + */ +function comparableBody(body: AgentJournalMessageItem | undefined): boolean { + return ( + body?.kind === 'message' && + body.role === 'user' && + body.blocks.length > 0 && + body.blocks.every((block) => block.type === 'text') + ) +} + +function comparableSubmissions(journal: AgentSessionJournal): AgentJournalSubmission[] { + const { items, submissions } = journal.snapshot() + const bodies = new Map(items.map((item) => [item.itemId, item.body])) + return submissions.filter((submission) => { + if (submission.dispatchState !== 'pending' && submission.dispatchState !== 'unknown') { + return false + } + const body = bodies.get(agentJournalSubmissionKey(submission.clientMessageId)) + return comparableBody(body?.kind === 'message' ? body : undefined) + }) +} + +/** Items the journal already committed are not new evidence: leaving them + * claimable would let an undelivered message match an older identical one. */ +function unseenHistory( + journal: AgentSessionJournal, + history: ProviderHistoryWindow +): ProviderHistoryWindow { + const committed = new Set(journal.snapshot().items.map((item) => item.itemId)) + return { + ...history, + items: history.items.filter((item) => !committed.has(agentJournalItemKey(item.identity))) + } +} + +/** + * Decide what the crash boundary could only doubt. Returns the client message + * ids this pass settled, so the attach result stops reporting them unconfirmed. + */ +export async function reconcileJournalSubmissionsAgainstHistory(input: { + journal: AgentSessionJournal + fence: number + history: ProviderHistoryWindow +}): Promise { + const submissions = comparableSubmissions(input.journal) + if (submissions.length === 0) { + return [] + } + const settled: string[] = [] + for (const outcome of reconcileSubmissions({ + submissions, + history: unseenHistory(input.journal, input.history) + })) { + if (outcome.outcome === 'unknown') { + continue + } + await input.journal.resolveDispatch( + outcome.outcome === 'accepted' + ? { + clientMessageId: outcome.clientMessageId, + state: 'accepted', + providerIdentity: outcome.identity, + fence: input.fence, + recovered: true + } + : { + clientMessageId: outcome.clientMessageId, + state: 'rejected', + reason: outcome.reason, + fence: input.fence, + recovered: true + } + ) + settled.push(outcome.clientMessageId) + } + return settled +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts index 42f9289783a..2bd3c653a0b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adapter-router.ts @@ -105,6 +105,9 @@ export class StructuredAgentSessionAdapterRouter implements StructuredAgentSessi historyFilePath = (input: { identity: AgentSessionJournalIdentity }) => this.requireAgent(input.identity).historyFilePath?.(input) ?? Promise.resolve(null) + providerHistoryWindow = (input: { identity: AgentSessionJournalIdentity }) => + this.requireAgent(input.identity).providerHistoryWindow?.(input) ?? Promise.resolve(null) + closeSession = (sessionId: string): Promise => this.stopSession(sessionId, (adapter) => adapter.closeSession) 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 3bd14a057aa..057c9b907c2 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 @@ -27,6 +27,7 @@ import type { AgentSessionSlashCommand, AgentSessionWireRefusalCode } from '../../../shared/agent-session-wire' +import type { ProviderHistoryWindow } from '../agent-session-journal/journal-submission-reconciler' import type { StructuredAgentSessionEventSink } from './structured-agent-session-event-sink' export class AgentSessionAcquisitionRefusal extends Error { @@ -220,6 +221,14 @@ export type StructuredAgentSessionAdapter = { /** Transcript path for journal recovery. Omit to let the existing session-file * resolver discover it from the provider session id. */ historyFilePath?(input: { identity: AgentSessionJournalIdentity }): Promise + /** Provider history for restart reconciliation, bounded to what the provider + * recorded after the journal's last committed item. Only the adapter can say + * whether the read has a proven start and whether a turn is still running, so + * it owns both flags. Omit where the provider records no boundary-consistent + * history; an omitted window leaves every unsettled submission `unknown`. */ + providerHistoryWindow?(input: { + identity: AgentSessionJournalIdentity + }): Promise /** Gracefully stops the structured owner after its event stream is drained. */ /** Returns true only after the provider child exit is proven. */ closeSession?(sessionId: string): Promise diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-reconciliation.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-reconciliation.test.ts new file mode 100644 index 00000000000..97f429933d8 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-reconciliation.test.ts @@ -0,0 +1,150 @@ +// Attach is where the restart reconciler runs. These cover the wiring itself: +// that the window the adapter reports reaches the journal, that what it settles +// stops being reported unconfirmed, and that deciding never sends. + +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 { agentSessionRecordFixture } from '../../../shared/agent-session-record.test-fixture' +import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types' +import { digestPayload } from '../agent-session-journal/journal-payload-bounds' +import { journalDirectoryFor } from '../agent-session-journal/journal-paths' +import type { ProviderHistoryWindow } from '../agent-session-journal/journal-submission-reconciler' +import { createTrackedJournalOpener } from '../agent-session-journal/journal-store-test-open' +import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter' +import { + attachJournal, + journalIdentityFor, + type AgentSessionAttachParams +} from './structured-agent-session-attach' + +const RECORD = agentSessionRecordFixture() + +const PARAMS = { + envelope: { + sessionId: RECORD.sessionId, + clientOperationId: 'op-1', + expectedRuntimeFence: RECORD.lease.runtimeFence, + payloadFingerprint: 'fp' + }, + location: RECORD.location, + provider: 'claude', + agent: 'claude', + accountHome: RECORD.accountHome, + runtimeKind: 'native' +} as unknown as AgentSessionAttachParams + +const IDENTITY = journalIdentityFor(RECORD, PARAMS) + +let root: string +const journals = createTrackedJournalOpener() + +function userMessage(text: string): AgentJournalMessageItem { + return { kind: 'message', role: 'user', blocks: [{ type: 'text', text }] } +} + +function window(overrides: Partial = {}): ProviderHistoryWindow { + return { items: [], boundaryConsistent: true, turnInFlight: false, ...overrides } +} + +/** Only the surface `attachJournal` touches; every send-shaped method is a spy + * so a re-delivery would be visible rather than silent. */ +function adapterWith(providerHistoryWindow?: () => Promise): { + adapter: StructuredAgentSessionAdapter + dispatch: ReturnType +} { + const dispatch = vi.fn() + const adapter = { + dispatch, + ...(providerHistoryWindow ? { providerHistoryWindow } : {}) + } as unknown as StructuredAgentSessionAdapter + return { adapter, dispatch } +} + +/** A previous process wrote the submission row and died before its outcome. */ +async function crashedJournal(clientMessageId = 'cm_1', text = 'deploy the thing') { + const journal = await journals.open({ + identity: IDENTITY, + journalDir: journalDirectoryFor(root, { + workspaceId: IDENTITY.workspaceId, + sessionId: IDENTITY.sessionId + }) + }) + await journal.appendSubmission({ + clientMessageId, + payloadFingerprint: digestPayload(text), + body: userMessage(text), + fence: RECORD.lease.runtimeFence + }) + await journal.close() +} + +async function attach(adapter: StructuredAgentSessionAdapter) { + const attached = await attachJournal({ + record: RECORD, + params: PARAMS, + journalRoot: root, + adapter + }) + journals.track(attached.journal) + return attached +} + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-attach-reconcile-')) +}) + +afterEach(async () => { + await journals.closeAll() + await rm(root, { recursive: true, force: true }) +}) + +describe('attachJournal restart reconciliation', () => { + it('settles a provably undelivered submission and stops reporting it unconfirmed', async () => { + await crashedJournal() + const { adapter, dispatch } = adapterWith(async () => window()) + + const attached = await attach(adapter) + + expect(attached.unconfirmedClientMessageIds).toEqual([]) + const submission = attached.journal.submissions()[0] + expect(submission?.dispatchState).toBe('rejected') + expect(submission?.reason).toBe('not_delivered') + // Deciding is not sending: nothing here puts the message back on the wire. + expect(dispatch).not.toHaveBeenCalled() + }) + + it('still reports a submission unconfirmed when the window cannot decide it', async () => { + await crashedJournal() + const { adapter, dispatch } = adapterWith(async () => window({ turnInFlight: true })) + + const attached = await attach(adapter) + + expect(attached.unconfirmedClientMessageIds).toEqual(['cm_1']) + expect(attached.journal.submissions()[0]?.dispatchState).toBe('unknown') + expect(dispatch).not.toHaveBeenCalled() + }) + + it('leaves the crash boundary untouched for an adapter that reports no history', async () => { + await crashedJournal() + const { adapter } = adapterWith() + + const attached = await attach(adapter) + + expect(attached.unconfirmedClientMessageIds).toEqual(['cm_1']) + expect(attached.journal.submissions()[0]?.dispatchState).toBe('unknown') + }) + + it('does not fail the attach when reading provider history throws', async () => { + await crashedJournal() + const { adapter } = adapterWith(async () => { + throw new Error('transcript unreadable') + }) + + const attached = await attach(adapter) + + expect(attached.unconfirmedClientMessageIds).toEqual(['cm_1']) + expect(attached.journal.submissions()[0]?.dispatchState).toBe('unknown') + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach.ts index 38b626b2123..1e13261e546 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-attach.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach.ts @@ -39,6 +39,8 @@ import { agentSessionProviderHandleChainHead } from '../../../shared/agent-sessi import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' import { journalDirectoryFor } from '../agent-session-journal/journal-paths' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' +import { reconcileJournalSubmissionsAgainstHistory } from '../agent-session-journal/journal-restart-reconciliation' +import type { ProviderHistoryWindow } from '../agent-session-journal/journal-submission-reconciler' import { openAgentSessionJournalWithRecovery, type AgentSessionJournalRecovery @@ -161,14 +163,23 @@ export function journalIdentityFor( export type AttachedJournal = { journal: AgentSessionJournal recovery: AgentSessionJournalRecovery | null - /** Submissions the crash boundary settled as `unknown` on this open. */ + /** Submissions still `unknown` after this open: the crash boundary settled + * them there and provider history could not decide them either. */ unconfirmedClientMessageIds: string[] } /** * Open the session's journal, recovering it when the stored one is unusable, - * and settle every submission left in flight by a previous process. Orca never - * re-sends those; they surface as delivery unconfirmed. + * settle every submission left in flight by a previous process, then let + * provider history decide the ones it can prove. + * + * Why the reconciliation belongs HERE and nowhere else: this runs after the + * record store handed this host the lease and before `onAttached` starts a + * provider child, so nothing can be appending to the provider's history while it + * is read, and the window stays valid until the resume consumes it. Every other + * settlement site — a proven child exit, a handoff suspend — runs while the host + * may still start another child, and a read there could be overtaken before it + * is acted on. Orca still never re-sends: this decides state only. */ export async function attachJournal(input: { record: AgentSessionRecord @@ -193,9 +204,16 @@ export async function attachJournal(input: { try { // That await is a WRITE. A failure in it leaves the journal with no caller // holding a reference to close it. + const unconfirmed = await opened.journal.markPendingSubmissionsUnknown(fence) + const settled = await reconcileAgainstProviderHistory({ + adapter: input.adapter, + identity, + journal: opened.journal, + fence + }) return { ...opened, - unconfirmedClientMessageIds: await opened.journal.markPendingSubmissionsUnknown(fence) + unconfirmedClientMessageIds: unconfirmed.filter((id) => !settled.includes(id)) } } catch (error) { // A rejected close leaves the handle open, so the journal is retained for a @@ -205,6 +223,35 @@ export async function attachJournal(input: { } } +/** Reading provider history is best effort: a provider that reports none, or a + * read that fails, leaves every submission exactly as the crash boundary wrote + * it. The journal writes the outcome implies are NOT caught here — a failed + * write must reach the caller that retains the journal handle. */ +async function reconcileAgainstProviderHistory(input: { + adapter: StructuredAgentSessionAdapter + identity: AgentSessionJournalIdentity + journal: AgentSessionJournal + fence: number +}): Promise { + if (!input.adapter.providerHistoryWindow) { + return [] + } + let history: ProviderHistoryWindow | null + try { + history = await input.adapter.providerHistoryWindow({ identity: input.identity }) + } catch { + return [] + } + if (!history) { + return [] + } + return reconcileJournalSubmissionsAgainstHistory({ + journal: input.journal, + fence: input.fence, + history + }) +} + /** * The first link of an adopting session's chain. *