diff --git a/src/main/native-chat/agent-session-journal/journal-corruption-repair.test.ts b/src/main/native-chat/agent-session-journal/journal-corruption-repair.test.ts index f0c30087795..978f076d6e9 100644 --- a/src/main/native-chat/agent-session-journal/journal-corruption-repair.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-corruption-repair.test.ts @@ -155,6 +155,50 @@ describe('a sequence gap', () => { expect(reopened.repair).toEqual({ malformedRows: 0 }) expect(reopened.snapshot().items.map((entry) => entry.body)).toEqual([body('m0'), body('m1')]) }) + + // The prefix survives, so there is no emptied epoch to re-anchor and — a gap + // costing no malformed row — no disclosure either. Without a durable marker + // the next probe reads a contiguous anchored prefix and calls it clean, and + // the rows the repair deleted are never asked for again. + it('still reports corrupt on the next probe, with the deleted suffix unrebuilt', async () => { + const journal = await open() + for (let ordinal = 0; ordinal < 5; ordinal += 1) { + await journal.appendItem(item(ordinal), body(`m${ordinal}`), { fence: 1 }) + } + await journal.close() + await withJournalDatabase((db) => { + db.prepare('DELETE FROM journal_rows WHERE seq = ?').run(4) + }) + + const repaired = await open() + await repaired.close() + expect(loadJournal(root, IDENTITY.sessionId)).toMatchObject({ corrupt: true }) + + // Same policy the emptied-epoch repair takes: a session that writes into the + // epoch owns it, and a later import must not replace rows the user has seen. + const writable = await open() + await writable.appendItem(item(9), body('typed after the repair'), { fence: 1 }) + await writable.close() + expect(loadJournal(root, IDENTITY.sessionId)).toMatchObject({ corrupt: false }) + }) + + // The disclosure is the repair talking about itself, not the session writing: + // counting it as content would retire the marker the instant it was raised. + it('is not settled by the repair disclosure it appends for a malformed row', async () => { + const journal = await open() + for (let ordinal = 0; ordinal < 3; ordinal += 1) { + await journal.appendItem(item(ordinal), body(`m${ordinal}`), { fence: 1 }) + } + await journal.close() + await withJournalDatabase((db) => { + db.prepare('UPDATE journal_rows SET row_json = ? WHERE seq = ?').run('}{', 3) + }) + + const repaired = await open() + expect(repaired.repair.malformedRows).toBe(1) + await repaired.close() + expect(loadJournal(root, IDENTITY.sessionId)).toMatchObject({ corrupt: true }) + }) }) describe('a missing epoch row', () => { diff --git a/src/main/native-chat/agent-session-journal/journal-database-schema.ts b/src/main/native-chat/agent-session-journal/journal-database-schema.ts index 02d4561ede7..37675fbcd8e 100644 --- a/src/main/native-chat/agent-session-journal/journal-database-schema.ts +++ b/src/main/native-chat/agent-session-journal/journal-database-schema.ts @@ -2,11 +2,15 @@ // // `journal_rows` is the append-only log; `journal_sessions` is the derived // projection, upserted in the SAME transaction as every row insert so the live -// epoch and the rows that belong to it can never disagree. +// epoch and the rows that belong to it can never disagree. `journal_repairs` +// carries at most one row per session: the standing demand for a rebuild a +// partial repair leaves behind (see journal-repair-marker.ts). /** DB shape version, carried in `PRAGMA user_version`. Independent of the row - * body version (`JournalRow.v`): a newer build can change either alone. */ -export const JOURNAL_DB_SCHEMA_VERSION = 1 + * body version (`JournalRow.v`): a newer build can change either alone. + * v2 added `journal_repairs`; a build without it would replay a partially + * repaired journal as clean, so it must latch read-only rather than write. */ +export const JOURNAL_DB_SCHEMA_VERSION = 2 export function createJournalTablesSql(): string { return ` @@ -23,5 +27,11 @@ CREATE TABLE IF NOT EXISTS journal_sessions ( epoch TEXT NOT NULL, updated_at INTEGER NOT NULL ); +CREATE TABLE IF NOT EXISTS journal_repairs ( + session_id TEXT PRIMARY KEY, + epoch TEXT NOT NULL, + content_from INTEGER NOT NULL, + repaired_at INTEGER NOT NULL +); ` } diff --git a/src/main/native-chat/agent-session-journal/journal-epoch-replacement.ts b/src/main/native-chat/agent-session-journal/journal-epoch-replacement.ts index 68a22bc7200..9512d5e9fa1 100644 --- a/src/main/native-chat/agent-session-journal/journal-epoch-replacement.ts +++ b/src/main/native-chat/agent-session-journal/journal-epoch-replacement.ts @@ -1,7 +1,8 @@ // Republishing a live item set into a fresh epoch. // // One transaction: discard every row, insert the epoch row plus the replacement -// items, and move the session projection. +// items, move the session projection, and retire any repair marker — this +// republished history is exactly what the marker was holding out for. import type { AgentJournalItemBody, @@ -10,6 +11,7 @@ import type { } from '../../../shared/agent-session-journal-types' import type Database from '../../sqlite/sync-database' import type { JournalLoad } from './journal-open' +import { clearJournalRepairMarker } from './journal-repair-marker' import { applyJournalRow, createJournalReducerState } from './journal-reducer' import { buildJournalItemRow, journalRowBase } from './journal-row-builders' import { @@ -64,6 +66,7 @@ export function replaceJournalEpoch(input: { input.db.exec('BEGIN IMMEDIATE') try { deleteAllJournalRows(input.db) + clearJournalRepairMarker(input.db, input.identity.sessionId) for (const row of rows) { insertJournalRow(input.db, input.identity.sessionId, row) } diff --git a/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts b/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts index 83d87970c20..e95ff821058 100644 --- a/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts +++ b/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts @@ -1,13 +1,15 @@ // Opening a new epoch. // // One transaction: discard every row of the superseded epoch, insert the new -// epoch row at sequence 1, and move the session projection onto it. Superseded -// rows are DELETED rather than retained — nothing would ever shed them. +// epoch row at sequence 1, move the session projection onto it, and retire any +// repair marker the superseded epoch was carrying. Superseded rows are DELETED +// rather than retained — nothing would ever shed them. import { AGENT_SESSION_JOURNAL_SCHEMA_VERSION } from '../../../shared/agent-session-journal-types' import type { AgentSessionProviderHandle } from '../../../shared/agent-session-journal-types' import type Database from '../../sqlite/sync-database' import type { JournalLoad } from './journal-open' +import { clearJournalRepairMarker } from './journal-repair-marker' import { applyJournalRow, createJournalReducerState } from './journal-reducer' import { deleteAllJournalRows, @@ -41,6 +43,7 @@ export function publishNewEpoch(input: { input.db.exec('BEGIN IMMEDIATE') try { deleteAllJournalRows(input.db) + clearJournalRepairMarker(input.db, input.sessionId) insertJournalRow(input.db, input.sessionId, row) upsertJournalSessionRow(input.db, input.sessionId, input.epoch, input.now) input.db.exec('COMMIT') diff --git a/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts b/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts index ff6f5c442ff..ed86d81861e 100644 --- a/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts @@ -7,7 +7,10 @@ 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 { AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types' +import type { + AgentSessionJournalIdentity, + AgentSessionProviderHandle +} from '../../../shared/agent-session-journal-types' import { createLegacyIdentityTracker } from './journal-legacy-identity' import { appendLegacyTranscriptMessages, @@ -28,21 +31,29 @@ function tick(): number { return clock } -function identity(agent: 'claude' | 'codex', sessionId: string): AgentSessionJournalIdentity { +type ImportAgent = 'claude' | 'codex' | 'grok' | 'omp' + +function providerHandle(agent: ImportAgent, sessionId: string): AgentSessionProviderHandle { + if (agent === 'claude') { + return { kind: 'claude', sessionId, leafUuid: null } + } + return agent === 'codex' + ? { kind: 'codex', threadId: sessionId } + : { kind: 'opaque', agent, value: sessionId } +} + +function identity(agent: ImportAgent, sessionId: string): AgentSessionJournalIdentity { return { sessionId, workspaceId: 'ws-1', hostId: 'host-1', agent, - providerHandle: - agent === 'claude' - ? { kind: 'claude', sessionId, leafUuid: null } - : { kind: 'codex', threadId: sessionId } + providerHandle: providerHandle(agent, sessionId) } } async function open( - agent: 'claude' | 'codex', + agent: ImportAgent, sessionId: string, overrides: Partial[0]> = {} ): Promise { @@ -553,3 +564,132 @@ describe('import failures', () => { expect(result).toMatchObject({ ok: false }) }) }) + +// A tool call is only the SOLE block of its message when the provider wrote it +// that way. Claude interleaves it with narration, Grok hangs `tool_calls` off a +// row that also has text, and omp's execution cells always pair the invocation +// with its output — so the multi-block path carries untrusted tool input too. +describe('multi-block legacy messages', () => { + const limits = { ...DEFAULT_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 64 } + const oversized = 'x'.repeat(10_000) + + /** The tool-call block of the first imported multi-block message. */ + function importedToolCallBlock(journal: AgentSessionJournal): unknown { + for (const entry of journal.snapshot().items) { + if (entry.body.kind !== 'message') { + continue + } + const block = entry.body.blocks.find((candidate) => candidate.type === 'tool-call') + if (block) { + return block.input + } + } + return null + } + + it('bounds a Claude tool call that shares its message with narration', async () => { + const journal = await open('claude', CLAUDE_SESSION, { + journalDir: join(root, 'claude-mixed-journal') + }) + const filePath = await writeFixture('claude-mixed.jsonl', [ + { + parentUuid: null, + isSidechain: false, + type: 'assistant', + message: { + role: 'assistant', + content: [ + { type: 'text', text: 'Editing the file.' }, + { + type: 'tool_use', + id: 'toolu_mixed', + name: 'Edit', + input: { file_path: 'a.ts', patch: oversized } + } + ] + }, + uuid: 'dd22be00-1111-4222-8333-444455556666', + timestamp: '2026-08-05T10:00:09.000Z', + sessionId: CLAUDE_SESSION + } + ]) + + const result = await importLegacyTranscriptIntoJournal({ + journal, + agent: 'claude', + sessionId: CLAUDE_SESSION, + fence: 1, + options: { filePath, limits } + }) + + expect(result.ok).toBe(true) + expect(importedToolCallBlock(journal)).toMatchObject({ + truncated: true, + byteLength: expect.any(Number), + digest: expect.stringMatching(/^[0-9a-f]{64}$/), + head: expect.any(String) + }) + expect(JSON.stringify(journal.snapshot().items)).not.toContain('x'.repeat(1_000)) + await journal.close() + }) + + it('bounds a Grok tool call that shares its row with assistant text', async () => { + const journal = await open('grok', CODEX_SESSION, { + journalDir: join(root, 'grok-mixed-journal') + }) + const filePath = await writeFixture('grok-mixed.jsonl', [ + { + type: 'assistant', + id: 'asst-mixed', + timestamp: '2026-08-05T10:00:09.000Z', + content: [{ type: 'text', text: 'Searching.' }], + tool_calls: [{ id: 'c1', name: 'grep', arguments: JSON.stringify({ pattern: oversized }) }] + } + ]) + + const result = await importLegacyTranscriptIntoJournal({ + journal, + agent: 'grok', + sessionId: CODEX_SESSION, + fence: 1, + options: { filePath, limits } + }) + + expect(result.ok).toBe(true) + expect(importedToolCallBlock(journal)).toMatchObject({ truncated: true }) + expect(JSON.stringify(journal.snapshot().items)).not.toContain('x'.repeat(1_000)) + await journal.close() + }) + + it('bounds an omp execution cell, whose invocation always ships with its output', async () => { + const journal = await open('omp', CODEX_SESSION, { + journalDir: join(root, 'omp-mixed-journal') + }) + const filePath = await writeFixture('omp-mixed.jsonl', [ + { + type: 'message', + id: 'omp-mixed-1', + timestamp: '2026-08-05T10:00:09.000Z', + message: { + role: 'bashExecution', + command: `echo ${oversized}`, + output: 'done', + exitCode: 0 + } + } + ]) + + const result = await importLegacyTranscriptIntoJournal({ + journal, + agent: 'omp', + sessionId: CODEX_SESSION, + fence: 1, + options: { filePath, limits } + }) + + expect(result.ok).toBe(true) + expect(importedToolCallBlock(journal)).toMatchObject({ truncated: true }) + expect(JSON.stringify(journal.snapshot().items)).not.toContain('x'.repeat(1_000)) + await journal.close() + }) +}) 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 2f89b15f8fe..6a2ff099dbc 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 @@ -253,10 +253,11 @@ function legacyItemBody( } } -/** Inline block text keeps only a bounded head plus an explicit marker. No blob - * is written: the marker carries the digest and byte length, and the source - * transcript remains the full copy — a blob here would be unreferenced by the - * render model and pruned at the next compaction. */ +/** Every block that can carry untrusted bulk is bounded here, tool calls + * included: a provider decoder is free to put one alongside narration, and the + * sole-block path above never sees those. The remainder is discarded rather + * than stored elsewhere — the marker keeps its digest and byte length, and the + * source transcript remains the full copy. */ function boundBlock(block: NativeChatBlock, limits: JournalPayloadLimits): NativeChatBlock { if (block.type === 'text') { return { ...block, text: boundInlineText(block.text, limits).text } @@ -264,5 +265,8 @@ function boundBlock(block: NativeChatBlock, limits: JournalPayloadLimits): Nativ if (block.type === 'tool-result') { return { ...block, output: boundInlineText(block.output, limits).text } } + if (block.type === 'tool-call') { + return { ...block, input: boundToolInput(block.input, limits) } + } return block } diff --git a/src/main/native-chat/agent-session-journal/journal-open.ts b/src/main/native-chat/agent-session-journal/journal-open.ts index 8645e307ad7..350d2e9cfb9 100644 --- a/src/main/native-chat/agent-session-journal/journal-open.ts +++ b/src/main/native-chat/agent-session-journal/journal-open.ts @@ -22,6 +22,7 @@ import { readJournalSessionEpoch } from './journal-row-table' import { JOURNAL_REPAIR_DISCLOSURE_ITEM_ID } from './journal-repair-disclosure' +import { pendingJournalRepairSequence } from './journal-repair-marker' import { parseJournalRow, type JournalRow } from './journal-row-schema' /** Every epoch row is sequence 1, and no compaction moves that floor. */ @@ -59,6 +60,11 @@ export function replayJournal( } const state = createJournalReducerState(sessionId, epoch) const stored = readJournalEpochRows(db, sessionId, epoch) + // A partial repair keeps its prefix, so the surviving rows look contiguous and + // anchored however much of the timeline it deleted. Its marker is what still + // says otherwise, naming the sequence past which the epoch would be its own + // history again. + const repairedFrom = pendingJournalRepairSequence(db, sessionId, epoch) const rows: JournalRow[] = [] let malformedRows = 0 let latched = false @@ -115,7 +121,12 @@ export function replayJournal( return { state, readOnly: latched, - corrupt: Boolean(gap) || malformedRows > 0 || unanchored || awaitsProviderHistory(rows), + corrupt: + Boolean(gap) || + malformedRows > 0 || + unanchored || + (repairedFrom !== null && awaitsRebuild(rows, repairedFrom)) || + awaitsProviderHistory(rows), malformedRows, ...(truncateFrom !== undefined && !latched ? { truncateFrom } : {}) } @@ -124,18 +135,29 @@ export function replayJournal( /** * The epoch a total repair published, still holding nothing but its own anchor * and disclosure. The rows it dropped were never reconstructed, so provider - * history has to be retried rather than this being called a clean timeline. The - * moment the session writes content of its own the epoch is its own history and - * the retry stops. + * history has to be retried rather than this being called a clean timeline. */ function awaitsProviderHistory(rows: readonly JournalRow[]): boolean { const anchor = rows[0] if (anchor?.kind !== 'epoch' || anchor.reason !== 'unreconcilable_prefix') { return false } + // The anchor sits at sequence 1, so content of the epoch's own starts at 2. + return awaitsRebuild(rows, FIRST_JOURNAL_SEQUENCE + 1) +} + +/** + * True while everything at or above `contentFrom` is the repair's own + * bookkeeping: the deleted history was never rebuilt, so the provider has to be + * asked again. The moment the session writes content of its own past that + * sequence the epoch IS its own history, and the retry stops rather than a + * later import replacing rows the user has since seen. + */ +function awaitsRebuild(rows: readonly JournalRow[], contentFrom: number): boolean { return rows.every( (row) => - row === anchor || (row.kind === 'item' && row.itemId === JOURNAL_REPAIR_DISCLOSURE_ITEM_ID) + row.seq < contentFrom || + (row.kind === 'item' && row.itemId === JOURNAL_REPAIR_DISCLOSURE_ITEM_ID) ) } diff --git a/src/main/native-chat/agent-session-journal/journal-repair-marker.ts b/src/main/native-chat/agent-session-journal/journal-repair-marker.ts new file mode 100644 index 00000000000..f9ef09e688d --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-repair-marker.ts @@ -0,0 +1,69 @@ +// The standing demand for a rebuild that a repair leaves behind. +// +// A repair that empties the epoch republishes an `unreconcilable_prefix` anchor, +// and replay reads that back as history still owed. A repair that KEEPS a prefix +// has no such anchor to publish and — for a plain sequence gap — no malformed +// row to disclose either, so nothing on disk would record that the deleted +// suffix was never reconstructed. This marker is that record, written in the +// SAME transaction as the deletion: a crash between the two would otherwise +// leave the rows gone with nothing left asking for them back. +// +// It records the first sequence at which the epoch would hold content of its +// own again, because it retires under exactly the rule the emptied-epoch anchor +// takes: a fresh epoch carries the rebuild, and a session that writes past that +// sequence owns the epoch and stops the retry. + +import type Database from '../../sqlite/sync-database' +import { deleteJournalRowSuffix } from './journal-row-table' + +const SELECT_REPAIR = 'SELECT epoch, content_from FROM journal_repairs WHERE session_id = ?' +const UPSERT_REPAIR = `INSERT INTO journal_repairs (session_id, epoch, content_from, repaired_at) +VALUES (?, ?, ?, ?) +ON CONFLICT(session_id) DO UPDATE SET + epoch = excluded.epoch, content_from = excluded.content_from, repaired_at = excluded.repaired_at` +const DELETE_REPAIR = 'DELETE FROM journal_repairs WHERE session_id = ?' + +/** + * The sequence a pending repair on THIS epoch left free, or null when none is + * pending. Epoch-scoped: a marker raised on an epoch that has since been + * superseded says nothing about the live one. + */ +export function pendingJournalRepairSequence( + db: Database.Database, + sessionId: string, + epoch: string +): number | null { + const row = db.prepare(SELECT_REPAIR).get(sessionId) as + | { epoch?: string; content_from?: number } + | undefined + return row?.epoch === epoch ? (row.content_from ?? null) : null +} + +/** Retires the marker. Called from inside the epoch transactions, whose new + * epoch is the rebuilt history the marker was holding out for. */ +export function clearJournalRepairMarker(db: Database.Database, sessionId: string): void { + db.prepare(DELETE_REPAIR).run(sessionId) +} + +/** Drop the rejected suffix and record that it is owed, atomically. */ +export function deleteJournalRepairedSuffix(input: { + db: Database.Database + sessionId: string + epoch: string + /** First sequence of the rejected suffix. */ + fromSeq: number + /** First sequence left free once the suffix is gone. */ + contentFrom: number + now: number +}): number { + input.db.exec('BEGIN IMMEDIATE') + try { + const deleted = deleteJournalRowSuffix(input.db, input.sessionId, input.epoch, input.fromSeq) + input.db.prepare(UPSERT_REPAIR).run(input.sessionId, input.epoch, input.contentFrom, input.now) + input.db.exec('COMMIT') + return deleted + } catch (error) { + input.db.exec('ROLLBACK') + throw error + } +} diff --git a/src/main/native-chat/agent-session-journal/journal-store-open.ts b/src/main/native-chat/agent-session-journal/journal-store-open.ts index 4f5794388b9..e002048a094 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-open.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-open.ts @@ -18,9 +18,11 @@ export async function openJournalStoreState(input: { journalDir: string loaded: JournalLoad | null | undefined replay: () => JournalLoad | null - /** Drops the rejected suffix. Corruption is not preserved; the load reports - * `corrupt` and recovery rebuilds the epoch from provider history. */ - deleteSuffix: (fromSeq: number) => number + /** Drops the rejected suffix and records the rebuild it owes, in ONE + * transaction. Corruption is not preserved; replay keeps reporting `corrupt` + * until provider history republishes the epoch or the session writes past + * `contentFrom`, the first sequence the repair left free. */ + deleteSuffix: (fromSeq: number, contentFrom: number) => number start: () => void adopt: (loaded: JournalLoad) => void /** Republishes an anchor row for an epoch a repair emptied. */ @@ -42,7 +44,7 @@ export async function openJournalStoreState(input: { } input.adopt(loaded) if (loaded.truncateFrom !== undefined && !loaded.readOnly) { - input.deleteSuffix(loaded.truncateFrom) + input.deleteSuffix(loaded.truncateFrom, loaded.state.lastSequence + 1) } // A repair that took every live row leaves the epoch with no anchor. Publish // one before anything can append into it: an ordinary row at sequence 1 would diff --git a/src/main/native-chat/agent-session-journal/journal-store-restore.ts b/src/main/native-chat/agent-session-journal/journal-store-restore.ts index b41c8ad5b27..53683a797d1 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-restore.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-restore.ts @@ -10,7 +10,7 @@ import type { JournalEpochController } from './journal-epoch-controller' import { replayJournal } from './journal-open' import type { JournalStoreHost } from './journal-store-collaborators' import { openJournalStoreState } from './journal-store-open' -import { deleteJournalRowSuffix } from './journal-row-table' +import { deleteJournalRepairedSuffix } from './journal-repair-marker' export function restoreJournalStore( host: JournalStoreHost, @@ -23,13 +23,15 @@ export function restoreJournalStore( const opened = host.database() return replayJournal(opened.db, opened.readOnly, host.identity.sessionId) }, - deleteSuffix: (fromSeq) => - deleteJournalRowSuffix( - host.database().db, - host.identity.sessionId, - host.state().epoch, - fromSeq - ), + deleteSuffix: (fromSeq, contentFrom) => + deleteJournalRepairedSuffix({ + db: host.database().db, + sessionId: host.identity.sessionId, + epoch: host.state().epoch, + fromSeq, + contentFrom, + now: host.now() + }), start: () => collaborators.epochController.start('session_created', 0), // `unreconcilable_prefix` is the durable statement that this epoch exists // because a repair emptied one: replay reads it back and keeps asking for diff --git a/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.test.ts b/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.test.ts index 361ee1d8e43..214c889b5c9 100644 --- a/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.test.ts +++ b/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.test.ts @@ -245,7 +245,7 @@ describe('openAgentSessionJournalWithRecovery', () => { }) // The rehydrate deletes every live row to publish its replacement epoch, so - // whatever replay rejected has to be in quarantine BEFORE the import runs. + // everything replay rejected is gone for good by the time the import runs. // Orca minted the submission, receipt and lifecycle identities; no provider // transcript can hand them back. it('rebuilds from provider history when the epoch row itself is gone', async () => { @@ -308,6 +308,45 @@ describe('openAgentSessionJournalWithRecovery', () => { expect(opened.journal.snapshot().items.map((entry) => entry.body.kind)).toEqual(['message']) }) + // A repair that KEEPS a prefix has no emptied epoch to anchor, so nothing + // about the surviving rows records that the deleted suffix was never rebuilt. + // Unmarked, the next probe reads a contiguous anchored prefix, calls it clean, + // and the dropped stretch of timeline is gone for good. + it('keeps a partially repaired journal corrupt until provider history replaces it', async () => { + await seedJournal(3) + await deleteRow(3) + const empty = join(root, 'empty.jsonl') + await writeFile(empty, '', 'utf-8') + + const first = await openAgentSessionJournalWithRecovery({ + identity: IDENTITY, + journalDir, + fence: 1, + historyFilePath: empty + }) + journals.track(first.journal) + expect(first.recovery).toMatchObject({ trigger: 'journal_corrupt', imported: 0 }) + expect(first.recovery?.error).toBeTruthy() + // Only the unanchored suffix went; the prefix the repair kept is still live. + expect(first.journal.snapshot().items.map((entry) => entry.body.kind)).toEqual(['message']) + await first.journal.close() + + // The deletion is durable, so the demand for a rebuild has to be too. + expect(await loadJournal(journalDir, CODEX_SESSION)).toMatchObject({ corrupt: true }) + + // A readable transcript rebuilds the epoch, and THAT is what retires it. + const retried = await openAgentSessionJournalWithRecovery({ + identity: IDENTITY, + journalDir, + fence: 1, + historyFilePath + }) + journals.track(retried.journal) + expect(retried.recovery?.imported).toBeGreaterThan(0) + await retried.journal.close() + expect(await loadJournal(journalDir, CODEX_SESSION)).toMatchObject({ corrupt: false }) + }) + // The reproduced path. Deleting sequence 1 leaves every surviving row // unanchored, so the repair drops ALL of them — and provider history is not // there to publish a replacement. A journal in that state used to reopen as diff --git a/src/shared/agent-session-journal-schemas.ts b/src/shared/agent-session-journal-schemas.ts index 2b1ab5404fc..505aa22c82f 100644 --- a/src/shared/agent-session-journal-schemas.ts +++ b/src/shared/agent-session-journal-schemas.ts @@ -1,11 +1,12 @@ // ─── Canonical runtime schemas for the journal render model ───────────────── -// The journal admits JSON it did not just write — snapshot files and log rows -// re-enter from disk and are republished to clients — while the reducer, the +// The journal admits JSON it did not just write — persisted rows re-enter from +// SQLite on replay and are republished to clients — while the reducer, the // shared projection, and the prompt surfaces dereference nested fields without // guards. These schemas are the single deep validators for that render model: // admission must reject a JSON-valid but structurally wrong item (a question -// whose `options` are null, a prompt without its `resolution`) so corruption -// lands in quarantine instead of throwing mid-render. +// whose `options` are null, a prompt without its `resolution`) so the row is +// rejected at replay, where a repair can delete it, instead of throwing +// mid-render. // // Discriminants (`kind`, known block `type`s) are validated deeply. Open string // fields (roles, dispatch/tool states) stay type-checked, never enum-checked, @@ -154,8 +155,8 @@ export function isAdmissibleAgentJournalSubmission( return AgentJournalSubmissionSchema.safeParse(value).success } -/** Compile-time proof that every canonical value is admissible, so admission - * can never quarantine a row a writer in this build produced. The schemas are +/** Compile-time proof that every canonical value is admissible, so replay can + * never reject a row a writer in this build produced. The schemas are * deliberately wider on open string fields, so only this direction holds. */ type Admits = T export type CanonicalJournalShapesAreAdmissible = [ diff --git a/src/shared/agent-session-journal-types.ts b/src/shared/agent-session-journal-types.ts index f5dabfdec23..f37547b7e79 100644 --- a/src/shared/agent-session-journal-types.ts +++ b/src/shared/agent-session-journal-types.ts @@ -58,13 +58,15 @@ export type AgentJournalItemIdentity = // ─── Bounded payloads ─────────────────────────────────────────────────────── -/** A tool output or diff body clipped to a head plus a content-addressed - * remainder. Crossing a bound sets `truncated`; it never silently drops. */ +/** A tool output or diff body clipped to a head. The remainder is DISCARDED, + * never stored: crossing a bound sets `truncated` and the two fields below + * describe what was dropped, so it is marked rather than silently lost. */ export type AgentJournalBoundedPayload = { head: string /** Byte length of the ORIGINAL payload, not of `head`. */ byteLength: number - /** sha256 of the original payload, and the blob store key when `truncated`. */ + /** sha256 of the original payload — identification only; nothing stores or + * retrieves the discarded remainder by it. */ digest: string truncated: boolean }