From 3cde95ddbf3b9c6dd168fb4c168f854cd461ebee Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Sun, 6 Sep 2026 21:30:17 -0700 Subject: [PATCH] fix: validate history before claiming adopted sessions --- .../journal-legacy-import.ts | 35 +++--- ...tured-agent-session-adopted-import.test.ts | 115 ++++++++++-------- ...structured-agent-session-adopted-import.ts | 77 +++++++++++- .../structured-agent-session-attach-flow.ts | 20 +-- 4 files changed, 172 insertions(+), 75 deletions(-) diff --git a/src/main/native-chat/agent-session-journal/journal-legacy-import.ts b/src/main/native-chat/agent-session-journal/journal-legacy-import.ts index 6a2ff099dbc..0907ebd28f2 100644 --- a/src/main/native-chat/agent-session-journal/journal-legacy-import.ts +++ b/src/main/native-chat/agent-session-journal/journal-legacy-import.ts @@ -90,6 +90,24 @@ export async function importLegacyTranscriptIntoJournal(input: { fence: number options?: LegacyImportOptions }): Promise { + const prepared = await prepareLegacyTranscriptImport(input) + if (!prepared.ok) { + return prepared + } + // An empty import must preserve any existing repair anchor and disclosure. + if (prepared.items.length === 0) { + const current = input.journal.cursor() + return { ok: true, epoch: current.epoch, cursor: current, imported: 0, replaced: false } + } + const cursor = await input.journal.replaceEpochItems('legacy_import', input.fence, prepared.items) + return { ok: true, epoch: cursor.epoch, cursor, imported: prepared.items.length, replaced: true } +} + +export async function prepareLegacyTranscriptImport(input: { + agent: AgentType + sessionId: string + options?: LegacyImportOptions +}): Promise<{ ok: true; items: JournalReplacementItem[] } | { ok: false; error: string }> { const options = input.options ?? {} const limits = options.limits ?? DEFAULT_JOURNAL_PAYLOAD_LIMITS const transcriptAgent = resolveNativeChatTranscriptAgent(input.agent) @@ -145,22 +163,7 @@ export async function importLegacyTranscriptIntoJournal(input: { observedAt: message.timestamp ?? undefined }) } - // A transcript that decodes to nothing reconstructs nothing, and an empty - // replacement is not a harmless no-op: it would delete the repair's anchor and - // its disclosure, leaving nothing to ask for the history again. The epoch - // stands so a later read can still rebuild it. - if (replacement.length === 0) { - const current = input.journal.cursor() - return { ok: true, epoch: current.epoch, cursor: current, imported: 0, replaced: false } - } - const cursor = await input.journal.replaceEpochItems('legacy_import', input.fence, replacement) - return { - ok: true, - epoch: cursor.epoch, - cursor, - imported: decoded.messages.length, - replaced: true - } + return { ok: true, items: replacement } } const TRANSCRIPT_DECODERS = { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts index 05ba75481b1..0248799980a 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts @@ -1,12 +1,6 @@ -// Adopting a history row imports its transcript into the new session's journal. -// -// The import runs AFTER the provider is acquired, because for a create there is no earlier moment -// — the journal does not exist until the child does. That ordering is what makes the failure cases -// here load-bearing: by then the provider has already resumed and holds the conversation in -// context, so an attach that succeeds with an empty journal would show the user a blank chat beside -// an agent that can already answer from history. +// Source validation must finish before a new session claims the provider conversation. -import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { mkdtemp, rm, writeFile, truncate } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import { afterEach, describe, expect, it, vi } from 'vitest' @@ -18,7 +12,9 @@ import { type AgentSessionAttachParams } from './structured-agent-session-attach' import { performAttach, type AttachFlowInput } from './structured-agent-session-attach-flow' +import { AgentSessionJournal } from '../agent-session-journal/journal-store' import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' +import * as legacyImport from '../agent-session-journal/journal-legacy-import' const NOW = 1_800_000_000_000 const SESSION = 'codex_adopting_session' @@ -181,51 +177,74 @@ describe('adopting a provider conversation on create', () => { expect(sessionAdapter.acquire).toHaveBeenCalledTimes(1) }) - it('fails the attach when the adopted transcript cannot be read', async () => { - root = await mkdtemp(join(tmpdir(), 'orca-adopt-missing-')) - const sessionAdapter = adapter() + it.each(['missing', 'oversized', 'empty', 'invalid', 'source-less'] as const)( + 'refuses %s source before claiming a conversation', + async (kind) => { + root = await mkdtemp(join(tmpdir(), 'orca-adopt-preflight-')) + const transcriptPath = join(root, 'rollout.jsonl') + if (kind === 'oversized') { + await writeCodexRollout(transcriptPath, 'original turn') + await truncate(transcriptPath, 16 * 1024 * 1024 + 1) + } else if (kind === 'empty' || kind === 'invalid') { + await writeFile(transcriptPath, kind === 'empty' ? '' : 'not json\n') + } + const sessionAdapter = adapter() + const onAttached = vi.fn() + const result = await attach( + kind === 'source-less' ? undefined : transcriptPath, + sessionAdapter, + onAttached + ) + expect(result).toMatchObject({ + ok: false, + refusal: { code: 'agent_session_identity_required' } + }) + expect(sessionAdapter.acquire).not.toHaveBeenCalled() + expect(sessionAdapter.releaseAcquisition).not.toHaveBeenCalled() + expect(onAttached).not.toHaveBeenCalled() + expect(store?.getRecord(SESSION)).toBeNull() + expect(store?.listOperationRows()).toEqual([]) + if (kind === 'oversized') { + expect(JSON.stringify(result)).toContain('import bound') + } + } + ) + + it('still releases acquisition and closes the provisional journal on an import write failure', async () => { + root = await mkdtemp(join(tmpdir(), 'orca-adopt-write-failure-')) + const transcriptPath = join(root, 'rollout.jsonl') + await writeCodexRollout(transcriptPath, 'valid source') + vi.spyOn(AgentSessionJournal.prototype, 'replaceEpochItems').mockRejectedValueOnce( + new Error('disk write failed') + ) const close = vi.spyOn(agentSessionJournalCloseRetries, 'closeOrRetain') - - // A post-acquisition failure throws rather than answering a refusal — the same path a journal - // failure already takes — so the caller learns the outcome is unknown, not that nothing ran. - await expect(attach(join(root, 'does-not-exist.jsonl'), sessionAdapter)).rejects.toThrow( - /ENOENT|no such file/ - ) - // The provider had already resumed, so its acquisition is released rather than left holding a - // conversation no surface will ever show. - expect(sessionAdapter.releaseAcquisition).toHaveBeenCalled() + const sessionAdapter = adapter() + await expect(attach(transcriptPath, sessionAdapter)).rejects.toThrow('disk write failed') + expect(sessionAdapter.acquire).toHaveBeenCalledTimes(1) + expect(sessionAdapter.releaseAcquisition).toHaveBeenCalledTimes(1) expect(close).toHaveBeenCalledTimes(1) - close.mockRestore() }) - it('refuses a source-less replay identity when the journal was never imported', async () => { - root = await mkdtemp(join(tmpdir(), 'orca-adopt-unimported-replay-')) + it('prepares a valid source once before acquisition and imports those exact items', async () => { + root = await mkdtemp(join(tmpdir(), 'orca-adopt-once-')) + const transcriptPath = join(root, 'rollout.jsonl') + await writeCodexRollout(transcriptPath, 'prepared before acquiring') + const prepare = vi.spyOn(legacyImport, 'prepareLegacyTranscriptImport') const sessionAdapter = adapter() - - await expect(attach(undefined, sessionAdapter)).rejects.toThrow( - 'agent_session_identity_required' + const acquire = sessionAdapter.acquire + sessionAdapter.acquire = vi.fn(async (input) => { + expect(prepare).toHaveBeenCalledTimes(1) + await rm(transcriptPath) + return acquire(input) + }) + const result = await attach(transcriptPath, sessionAdapter, async ({ journal }) => + journal.close() ) - expect(sessionAdapter.releaseAcquisition).toHaveBeenCalled() - }) - - it('fails the attach when the adopted transcript decodes to no messages', async () => { - root = await mkdtemp(join(tmpdir(), 'orca-adopt-empty-')) - const transcriptPath = join(root, 'empty.jsonl') - // Well-formed but conversation-free: the row promised turns and the provider resumed them, so - // an empty journal here is a disagreement, not an empty chat. - await writeFile( - transcriptPath, - `${JSON.stringify({ - type: 'session_meta', - payload: { id: THREAD, timestamp: '2026-09-06T18:00:00.000Z', cwd: '/workspace' } - })}\n`, - 'utf8' - ) - const sessionAdapter = adapter() - - await expect(attach(transcriptPath, sessionAdapter)).rejects.toThrow( - 'agent_session_identity_required' - ) - expect(sessionAdapter.releaseAcquisition).toHaveBeenCalled() + expect(result.ok).toBe(true) + if (!result.ok) { + throw new Error('attach failed') + } + expect(JSON.stringify(result.value.page.items)).toContain('prepared before acquiring') + expect(prepare).toHaveBeenCalledTimes(1) }) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts index 9f0ac061292..459d21b1f3b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts @@ -1,18 +1,91 @@ +import type { AgentSessionWireRefusal } from '../../../shared/agent-session-wire' import type { AgentSessionRecord } from '../../../shared/agent-session-record' import type { AgentSessionAttachParams, AttachedJournal } from './structured-agent-session-attach' -import { importLegacyTranscriptIntoJournal } from '../agent-session-journal/journal-legacy-import' +import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' +import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement' +import { + importLegacyTranscriptIntoJournal, + prepareLegacyTranscriptImport +} from '../agent-session-journal/journal-legacy-import' + +export async function prepareAdoptedTranscript( + params: AgentSessionAttachParams +): Promise< + | { ok: true; items: JournalReplacementItem[] | null } + | { ok: false; refusal: AgentSessionWireRefusal } +> { + try { + return { ok: true, items: await readAdoptedTranscript(params) } + } catch (error) { + return { + ok: false, + refusal: { + code: 'agent_session_identity_required', + message: error instanceof Error ? error.message : String(error) + } + } + } +} + +// Validate source input before a new record can claim the provider conversation. +async function readAdoptedTranscript( + params: AgentSessionAttachParams +): Promise { + const adopt = params.adopt + if (!adopt) { + return null + } + if (!adopt.transcriptPath) { + throw new Error('agent_session_identity_required') + } + const prepared = await prepareLegacyTranscriptImport({ + agent: params.agent, + sessionId: + adopt.providerHandle.kind === 'claude' + ? adopt.providerHandle.sessionId + : adopt.providerHandle.threadId, + options: { filePath: adopt.transcriptPath } + }) + if (!prepared.ok) { + throw new Error(prepared.error) + } + if (prepared.items.length === 0) { + throw new Error('agent_session_identity_required') + } + return prepared.items +} // Import before publication so the first visible chat agrees with the provider's resumed context. export async function importAdoptedTranscript( params: AgentSessionAttachParams, attached: AttachedJournal, - record: AgentSessionRecord + record: AgentSessionRecord, + prepared: JournalReplacementItem[] | null +): Promise { + try { + await applyAdoptedTranscript(params, attached, record, prepared) + } catch (error) { + // Publication has not taken ownership of this provisional journal yet. + await agentSessionJournalCloseRetries.closeOrRetain(attached.journal) + throw error + } +} + +async function applyAdoptedTranscript( + params: AgentSessionAttachParams, + attached: AttachedJournal, + record: AgentSessionRecord, + prepared: JournalReplacementItem[] | null ): Promise { const adopt = params.adopt // A new journal contains only its epoch row; replay must preserve subsequent durable writes. if (!adopt || attached.journal.cursor().sequence > 1) { return } + if (prepared) { + await attached.journal.replaceEpochItems('legacy_import', record.lease.runtimeFence, prepared) + return + } if (!adopt.transcriptPath) { throw new Error('agent_session_identity_required') } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts index aab67ce81a6..8697e76ba3b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts @@ -36,8 +36,10 @@ import type { StructuredAgentSessionEventSink } from './structured-agent-session import { readNativeSessionOptions } from './structured-agent-session-option-restoration' import { resolveAgentSessionReplayOutcome } from './structured-agent-session-replay-outcome' import { readAgentSessionHydrationPage } from './agent-session-history-page' -import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' -import { importAdoptedTranscript } from './structured-agent-session-adopted-import' +import { + importAdoptedTranscript, + prepareAdoptedTranscript +} from './structured-agent-session-adopted-import' export type AttachFlowInput = { store: AgentSessionRecordStore @@ -80,6 +82,12 @@ export async function performAttach( let acquisitionGeneration: string | null = null let reservedRecord: AgentSessionRecord | null = null let replayed = false + const preparedTranscript = store.getRecord(sessionId) + ? { ok: true as const, items: null } + : await prepareAdoptedTranscript(params) + if (!preparedTranscript.ok) { + return preparedTranscript + } try { const reserved = await store.reserveOwner( reserveRequestFor({ @@ -183,13 +191,7 @@ export async function performAttach( journalRoot: input.journalRoot, adapter: input.adapter }) - try { - await importAdoptedTranscript(params, attached, record) - } catch (error) { - // Publication has not taken ownership of this provisional journal yet. - await agentSessionJournalCloseRetries.closeOrRetain(attached.journal) - throw error - } + await importAdoptedTranscript(params, attached, record, preparedTranscript.items) await input.onAttached(attached, acquisitionGeneration) await store.recordOperationOutcome({ callerKey: input.callerKey,