From edf0a41a73ce0bc3b7b29cf3c2d731b2bbc1a168 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Fri, 11 Sep 2026 01:56:23 -0700 Subject: [PATCH] fix(native-chat): treat a record without its journal as recovery, not a fresh session MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A persisted record whose journal is gone answered the load probe exactly like a session that has none yet: read restore dropped it, and the next attach founded a blank `session_created` epoch that then stood as the session's authoritative history. A journal that is absent is now its own recovery trigger. The open founds the `journal_missing` anchor instead of a blank epoch and rebuilds from provider history; when that rebuild fails the demand survives a later write, because there is no kept prefix here to make a new message the session's own past. Read restore publishes such a record rather than pruning its tab. A pre-SQLite remnant is excluded — that history is on disk, in a format this build explains rather than rebuilds. The reset frame gains an optional `resetCause` naming which failure it came out of. The reset reason itself stays inside the released union, so a paired or mobile decoder that has never seen the cause still reloads. --- .../agent-session-journal/journal-open.ts | 29 +++- .../journal-row-schema.ts | 2 + .../agent-session-journal-recovery.test.ts | 131 ++++++++++++++++++ .../agent-session-journal-recovery.ts | 65 ++++++++- .../agent-session-reset-frame.ts | 21 +++ ...structured-agent-session-attach-context.ts | 8 +- ...ured-agent-session-attach-orchestration.ts | 8 +- .../structured-agent-session-attach.ts | 5 +- ...uctured-agent-session-read-restore.test.ts | 61 ++++++-- .../structured-agent-session-read-restore.ts | 46 ++++-- ...ructured-agent-session-subscribers.test.ts | 39 +++++- .../structured-agent-session-subscribers.ts | 7 +- src/shared/agent-session-journal-types.ts | 10 ++ src/shared/agent-session-wire.ts | 13 +- 14 files changed, 405 insertions(+), 40 deletions(-) create mode 100644 src/main/native-chat/agent-session-wire/agent-session-reset-frame.ts 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 e4cc5e0f73e..94bbc4bf5fc 100644 --- a/src/main/native-chat/agent-session-journal/journal-open.ts +++ b/src/main/native-chat/agent-session-journal/journal-open.ts @@ -40,6 +40,10 @@ export type JournalLoad = { /** Directory-internal: the first sequence of an unusable suffix. The store * deletes from here before it accepts a write; a probe leaves it alone. */ truncateFrom?: number + /** The live epoch stands in for a journal that was gone. Distinguished from + * plain corruption because a caller that DROPS a corrupt session would drop + * this one every time it restored, which is the disappearance being fixed. */ + historyMissing?: true } /** @@ -118,6 +122,8 @@ export function replayJournal( // A latched journal reduces to nothing by design; only a writable one can be // held to the anchor. const unanchored = !latched && rows[0]?.kind !== 'epoch' + const anchor = rows[0] + const historyMissing = anchor?.kind === 'epoch' && anchor.reason === 'journal_missing' return { state, readOnly: latched, @@ -128,18 +134,31 @@ export function replayJournal( (repairedFrom !== null && awaitsRebuild(rows, repairedFrom)) || awaitsProviderHistory(rows), malformedRows, - ...(truncateFrom !== undefined && !latched ? { truncateFrom } : {}) + ...(truncateFrom !== undefined && !latched ? { truncateFrom } : {}), + ...(historyMissing ? { historyMissing: true as const } : {}) } } /** - * 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. + * An epoch published in place of history this host could not keep, still holding + * nothing but its own anchor and disclosure. The rows it stands for were never + * reconstructed, so provider 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') { + if (anchor?.kind !== 'epoch') { + return false + } + // A journal that was GONE keeps its demand past any write. A partial repair + // retires on one because the prefix it kept is still the session's own + // history; here there is no prefix, so a later message is indistinguishable + // from a session that never had a past — and the retry is the only thing that + // would ever tell them apart again. + if (anchor.reason === 'journal_missing') { + return true + } + if (anchor.reason !== 'unreconcilable_prefix') { return false } // The anchor sits at sequence 1, so content of the epoch's own starts at 2. 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..cf22a8a4075 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 @@ -42,6 +42,8 @@ export const AGENT_JOURNAL_EPOCH_REASONS = [ 'legacy_import', 'corruption', 'unreconcilable_prefix', + /** The journal was GONE, and this anchor stands in for the history it held. */ + 'journal_missing', 'handle_forked', 'schema_unreadable' ] as const 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 214c889b5c9..3f26e8a5a6d 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 @@ -443,6 +443,137 @@ describe('openAgentSessionJournalWithRecovery', () => { ) }) + // A record whose journal is GONE used to open as a fresh store: the only + // trigger was damage IN a journal, so nothing here was damaged and the blank + // `session_created` epoch became the session's authoritative history. + describe('a record whose journal is absent', () => { + async function epochReasonOnDisk(): Promise { + const load = loadJournal(journalDir, CODEX_SESSION) + let reason: string | undefined + await withJournalDatabase(journalDir, (db) => { + const rows = readJournalEpochRows(db, CODEX_SESSION, load?.state.epoch ?? '') + reason = JSON.parse(rows[0]?.rowJson ?? '{}').reason + }) + return reason + } + + it('rebuilds from provider history instead of founding a blank session', async () => { + const opened = await openAgentSessionJournalWithRecovery({ + identity: IDENTITY, + journalDir, + fence: 1, + historyFilePath, + recordPredatesCall: true + }) + journals.track(opened.journal) + + expect(opened.recovery).toMatchObject({ trigger: 'journal_missing', reset: 'epoch_changed' }) + expect(opened.recovery?.imported).toBeGreaterThan(0) + expect(JSON.stringify(opened.journal.snapshot().items.map((entry) => entry.body))).toContain( + 'add a retry' + ) + await opened.journal.close() + // The blank epoch is absent: what stands is the rebuilt timeline's own. + expect(await epochReasonOnDisk()).toBe('legacy_import') + expect(loadJournal(journalDir, CODEX_SESSION)).toMatchObject({ corrupt: false }) + }) + + it('holds the rebuild owed when provider history cannot be read', async () => { + const opened = await openAgentSessionJournalWithRecovery({ + identity: IDENTITY, + journalDir, + fence: 1, + historyFilePath: join(root, 'missing.jsonl'), + recordPredatesCall: true + }) + journals.track(opened.journal) + + expect(opened.recovery).toMatchObject({ trigger: 'journal_missing', imported: 0 }) + expect(opened.recovery?.error).toBeTruthy() + expect(opened.journal.snapshot().items).toEqual([]) + await opened.journal.close() + // Not `session_created`: a blank epoch presented as history is the defect. + expect(await epochReasonOnDisk()).toBe('journal_missing') + expect(loadJournal(journalDir, CODEX_SESSION)).toMatchObject({ + corrupt: true, + historyMissing: true + }) + }) + + // A partial repair retires on a write because the prefix it kept is still + // the session's own history. Here there is no prefix, so a write would + // silently make the blank permanent. + it('keeps the rebuild owed across a later write', async () => { + const first = await openAgentSessionJournalWithRecovery({ + identity: IDENTITY, + journalDir, + fence: 1, + historyFilePath: join(root, 'missing.jsonl'), + recordPredatesCall: true + }) + journals.track(first.journal) + await first.journal.appendItem( + item(2), + { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'typed later' }] }, + { fence: 1 } + ) + await first.journal.close() + + expect(loadJournal(journalDir, CODEX_SESSION)).toMatchObject({ + corrupt: true, + historyMissing: true + }) + // And the retry still names the cause it recovered from, not the damage + // the standing anchor would otherwise look like. + const retried = await openAgentSessionJournalWithRecovery({ + identity: IDENTITY, + journalDir, + fence: 1, + historyFilePath, + recordPredatesCall: true + }) + journals.track(retried.journal) + expect(retried.recovery).toMatchObject({ trigger: 'journal_missing' }) + expect(retried.recovery?.imported).toBeGreaterThan(0) + await retried.journal.close() + expect(loadJournal(journalDir, CODEX_SESSION)).toMatchObject({ corrupt: false }) + }) + + it('founds an ordinary session when this call is the one creating it', async () => { + const opened = await openAgentSessionJournalWithRecovery({ + identity: IDENTITY, + journalDir, + fence: 1, + historyFilePath + }) + journals.track(opened.journal) + + expect(opened.recovery).toBeNull() + expect(opened.journal.snapshot().items).toEqual([]) + await opened.journal.close() + expect(await epochReasonOnDisk()).toBe('session_created') + expect(loadJournal(journalDir, CODEX_SESSION)).toMatchObject({ corrupt: false }) + }) + + it('founds an ordinary session when the record names no provider conversation', async () => { + const opened = await openAgentSessionJournalWithRecovery({ + identity: { + ...IDENTITY, + providerHandle: { kind: 'opaque', agent: 'codex', value: 'pending' } + }, + journalDir, + fence: 1, + historyFilePath, + recordPredatesCall: true + }) + journals.track(opened.journal) + + expect(opened.recovery).toBeNull() + await opened.journal.close() + expect(await epochReasonOnDisk()).toBe('session_created') + }) + }) + it('rebuilds the emptied epoch once provider history is readable again', async () => { await seedRepairableSession() diff --git a/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.ts b/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.ts index 6e771a31809..dab6aad2fbf 100644 --- a/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.ts +++ b/src/main/native-chat/agent-session-wire/agent-session-journal-recovery.ts @@ -1,28 +1,34 @@ // Journal recovery: rehydrate the timeline from provider history. // -// Two triggers, and they need different destinations. A journal whose prefix is -// unusable is writable, so it is rebuilt in place on a fresh epoch. A journal +// Three triggers, and they need different destinations. A journal whose prefix +// is unusable is writable, so it is rebuilt in place on a fresh epoch. A journal // written by a NEWER schema is not writable by this host at all — rebuilding it // in place would fork the sequence space a newer host still owns — so the // reconstruction goes to a schema-scoped sibling directory that is only ever -// written by hosts at this version and is never merged back. +// written by hosts at this version and is never merged back. A journal that is +// simply GONE gets the same in-place rebuild, but has to be told apart from a +// session that has no history yet — see `journalHistoryIsMissing`. import type { AgentType } from '../../../shared/agent-status-types' import { AGENT_SESSION_JOURNAL_SCHEMA_VERSION, + type AgentJournalResetCause, type AgentJournalResetReason, type AgentSessionJournalIdentity, type AgentSessionProviderHandle } from '../../../shared/agent-session-journal-types' import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' +import { findJournalFileFormatRemnant } from '../agent-session-journal/journal-file-format-remnant' import { importLegacyTranscriptIntoJournal } from '../agent-session-journal/journal-legacy-import' -import { loadJournal } from '../agent-session-journal/journal-open' +import { loadJournal, type JournalLoad } from '../agent-session-journal/journal-open' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import { openAgentSessionJournal } from '../agent-session-journal/journal-store-factory' export type AgentSessionJournalRecovery = { - trigger: 'journal_corrupt' | 'schema_unreadable' - /** What subscribers are told; both force a clean snapshot reload. */ + trigger: AgentJournalResetCause + /** What subscribers are told; every one forces a clean snapshot reload. The + * value stays inside the released reason union — `trigger` is what names the + * cause for a reader new enough to render it. */ reset: AgentJournalResetReason epoch: string imported: number @@ -50,12 +56,42 @@ export function providerHistoryId(handle: AgentSessionProviderHandle): string { return handle.kind === 'claude' ? handle.sessionId : handle.value } +/** + * A journal that is gone, told apart from a session that has none yet. + * + * Both answer nothing to a probe, and only the caller knows which it is looking + * at: `recordPredatesCall` is false exactly when this call is the one creating + * the session. A record that predates it and names a provider conversation had + * history; an opaque handle names no conversation, so there is nothing to + * rebuild and an empty timeline is the honest one. + */ +export function journalHistoryIsMissing(input: { + probe: JournalLoad | null + identity: AgentSessionJournalIdentity + journalDir: string + recordPredatesCall: boolean +}): boolean { + if (input.probe) { + return input.probe.historyMissing === true + } + return ( + input.recordPredatesCall && + input.identity.providerHandle.kind !== 'opaque' && + // A pre-SQLite remnant IS the history, in a format this build cannot read. + // It already has the disclosure that says so, and a rebuild would delete it. + findJournalFileFormatRemnant(input.journalDir) === null + ) +} + export async function openAgentSessionJournalWithRecovery(input: { identity: AgentSessionJournalIdentity journalDir: string fence: number /** Resolve directly to a transcript instead of discovering it by session id. */ historyFilePath?: string | null + /** See `journalHistoryIsMissing`. Default false: a caller that cannot tell + * must not turn a brand-new session into one owed a rebuild forever. */ + recordPredatesCall?: boolean }): Promise { const probe = loadJournal(input.journalDir, input.identity.sessionId) if (probe?.readOnly) { @@ -65,10 +101,27 @@ export async function openAgentSessionJournalWithRecovery(input: { }) return { journal, recovery: await rehydrateOrClose(input, journal, 'schema_unreadable') } } + const missing = journalHistoryIsMissing({ + probe, + identity: input.identity, + journalDir: input.journalDir, + recordPredatesCall: input.recordPredatesCall === true + }) const journal = await openAgentSessionJournal({ identity: input.identity, journalDir: input.journalDir }) + if (missing) { + // The open above founded a `session_created` epoch, which reads back as a + // session that never had a past. Replacing it with the missing-history + // anchor is what keeps the rebuild owed when the import below fails — and + // an anchor already in place is not re-rolled, or every attach would force + // every reader back to a fresh snapshot. + if (!probe?.historyMissing) { + await journal.rollEpoch('journal_missing', input.fence) + } + return { journal, recovery: await rehydrateOrClose(input, journal, 'journal_missing') } + } if (!probe?.corrupt) { return { journal, recovery: null } } diff --git a/src/main/native-chat/agent-session-wire/agent-session-reset-frame.ts b/src/main/native-chat/agent-session-wire/agent-session-reset-frame.ts new file mode 100644 index 00000000000..6daf20a3454 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/agent-session-reset-frame.ts @@ -0,0 +1,21 @@ +import type { + AgentJournalResetCause, + AgentJournalResetReason +} from '../../../shared/agent-session-journal-types' + +/** The reset half of a reset frame. The cause rides ALONGSIDE the reason rather + * than widening it: the reason reaches paired and mobile decoders that may + * reject a value they have never seen, and a reader that ignores the cause + * still reloads. */ +export type AgentSessionResetFrame = { + type: 'reset' + reset: AgentJournalResetReason + resetCause?: AgentJournalResetCause +} + +export function agentSessionResetFrame( + reset: AgentJournalResetReason, + cause?: AgentJournalResetCause +): AgentSessionResetFrame { + return { type: 'reset', reset, ...(cause ? { resetCause: cause } : {}) } +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts index 4b9b46716cb..1919eb92501 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts @@ -7,7 +7,10 @@ import type { AgentSessionTurnActivity, AgentSessionWireRefusal } from '../../../shared/agent-session-wire' -import type { AgentJournalResetReason } from '../../../shared/agent-session-journal-types' +import type { + AgentJournalResetCause, + AgentJournalResetReason +} from '../../../shared/agent-session-journal-types' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import type { StructuredAgentSessionHostDeps, @@ -25,7 +28,8 @@ export type StructuredAgentSessionAttachContext = { sessionId: string, journal: AgentSessionJournal, reset: AgentJournalResetReason, - fence: number + fence: number, + cause?: AgentJournalResetCause ) => void snapshot: (sessionId: string, journal: AgentSessionJournal, fence: number) => void publish: ( diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts index 85e02f2fea7..17c6b552ccf 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts @@ -150,7 +150,13 @@ export function attachStructuredAgentSession( } await recoverInterruptedCompaction(context.deps.store, sessionId, attached.journal, fence) if (attached.recovery) { - context.subscribers.reset(sessionId, attached.journal, attached.recovery.reset, fence) + context.subscribers.reset( + sessionId, + attached.journal, + attached.recovery.reset, + fence, + attached.recovery.trigger + ) } else if (previousFence !== undefined && previousFence !== fence) { context.subscribers.snapshot(sessionId, attached.journal, fence) } else { 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..66b221d059f 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 @@ -188,7 +188,10 @@ export async function attachJournal(input: { sessionId: identity.sessionId }), fence, - historyFilePath + historyFilePath, + // A null expected fence is this call creating the session, so its journal is + // absent rather than lost; anything else attaches to a record that existed. + recordPredatesCall: input.params.envelope.expectedRuntimeFence !== null }) try { // That await is a WRITE. A failure in it leaves the journal with no caller diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.test.ts index 34e3fb8a4cf..c5217a56b87 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.test.ts @@ -3,14 +3,18 @@ // A chat still in the pre-SQLite format has no `journal.db`, so the probe that // loads one reports nothing. Reading that as "no session" is what removed these // chats: an unpublished session is also what prunes its tab out of the saved -// workspace, so the tab is gone before anything can explain itself. +// workspace, so the tab is gone before anything can explain itself. A record +// whose journal is simply GONE reads identically to the probe and is the same +// disappearance, so it comes back too — carrying the anchor that says its +// history is still owed rather than a blank epoch that claims it had none. import { mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' -import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import type { AgentSessionRecord } from '../../../shared/agent-session-record' import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' +import { loadJournal } from '../agent-session-journal/journal-open' import { journalDirectoryFor } from '../agent-session-journal/journal-paths' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import { restoreStructuredAgentSessionRead } from './structured-agent-session-read-restore' @@ -43,10 +47,21 @@ const RECORD = { lease: { sessionId: SESSION_ID, runtimeKind: 'native', runtimeFence: 1 } } as unknown as AgentSessionRecord +/** The same record before any provider proved a handle for it: no conversation + * is named, so there is no history a rebuild could ask for. */ +const RECORD_WITHOUT_CONVERSATION = { + ...RECORD, + providerHandleChain: [] +} as unknown as AgentSessionRecord + const store = { getRecord: (sessionId: string) => (sessionId === SESSION_ID ? RECORD : null) } as unknown as AgentSessionRecordStore +const storeWithoutConversation = { + getRecord: (sessionId: string) => (sessionId === SESSION_ID ? RECORD_WITHOUT_CONVERSATION : null) +} as unknown as AgentSessionRecordStore + let journalRoot: string const opened: AgentSessionJournal[] = [] @@ -62,9 +77,14 @@ async function writeRemnant(name: string): Promise { beforeEach(async () => { journalRoot = await mkdtemp(join(tmpdir(), 'orca-read-restore-')) + // Recovery discovers the provider transcript by session id; point both roots + // at the empty temp tree so the search stays inside the fixture. + vi.stubEnv('ORCA_USER_DATA_PATH', journalRoot) + vi.stubEnv('CODEX_HOME', journalRoot) }) afterEach(async () => { + vi.unstubAllEnvs() await Promise.allSettled(opened.splice(0).map((journal) => journal.close())) await rm(journalRoot, { recursive: true, force: true }) }) @@ -94,12 +114,6 @@ describe('a session whose journal is still the pre-SQLite format', () => { opened.push(restored!.journal) }) - it('still drops a session with neither a journal nor a remnant', async () => { - const restored = await restoreStructuredAgentSessionRead(store, journalRoot, SESSION_ID) - - expect(restored).toBeNull() - }) - it('still drops a session with no record', async () => { await writeRemnant('log.jsonl') @@ -108,3 +122,34 @@ describe('a session whose journal is still the pre-SQLite format', () => { expect(restored).toBeNull() }) }) + +describe('a record whose journal is gone', () => { + it('comes back owed a rebuild rather than being dropped', async () => { + const restored = await restoreStructuredAgentSessionRead(store, journalRoot, SESSION_ID) + + expect(restored).not.toBeNull() + opened.push(restored!.journal) + // Nothing was imported — no transcript exists — so what must NOT be here is + // a blank epoch standing in as the session's own history. + expect(restored!.journal.snapshot().items).toEqual([]) + await restored!.journal.close() + const journalDir = journalDirectoryFor(journalRoot, { + workspaceId: WORKSPACE_ID, + sessionId: SESSION_ID + }) + expect(loadJournal(journalDir, SESSION_ID)).toMatchObject({ + corrupt: true, + historyMissing: true + }) + }) + + it('is still dropped when the record names no provider conversation', async () => { + const restored = await restoreStructuredAgentSessionRead( + storeWithoutConversation, + journalRoot, + SESSION_ID + ) + + expect(restored).toBeNull() + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.ts index 34fd6452b08..e39f448a5ce 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-read-restore.ts @@ -8,6 +8,10 @@ import { loadJournal } from '../agent-session-journal/journal-open' import { journalDirectoryFor } from '../agent-session-journal/journal-paths' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import { openAgentSessionJournal } from '../agent-session-journal/journal-store-factory' +import { + journalHistoryIsMissing, + openAgentSessionJournalWithRecovery +} from './agent-session-journal-recovery' import { attachFingerprintFields, journalIdentityFor, @@ -40,25 +44,45 @@ export async function restoreStructuredAgentSessionRead( workspaceId: record.location.workspaceId, sessionId }) + const identity = journalIdentityFor(record, params) const loaded = loadJournal(journalDir, sessionId) - if (loaded?.corrupt) { + // Restore is a read, so it cannot repair damage — but a journal that is GONE + // is not damage, and dropping it is the disappearance itself. The record + // predates every restore by definition, so this is always the caller that can + // tell lost history from a session that has none. + const missingHistory = journalHistoryIsMissing({ + probe: loaded, + identity, + journalDir, + recordPredatesCall: true + }) + if (loaded?.corrupt && !missingHistory) { return null } // A session still in the pre-SQLite format has no `journal.db` to load. Dropping // it here leaves it unpublished, which is also what prunes its tab out of the // saved workspace — so the chat disappears with nowhere to explain itself. - if (!loaded && !findJournalFileFormatRemnant(journalDir)) { + if (!loaded && !missingHistory && !findJournalFileFormatRemnant(journalDir)) { return null } - const journal = await openAgentSessionJournal({ - identity: journalIdentityFor(record, params), - journalDir, - // Omitted, not `null`: the store reads `null` as "replay already ran and - // found nothing" and founds a fresh epoch. In process the probe above is the - // previous statement, so the window is zero-width; this holds the line for a - // database another process creates in between. - ...(loaded ? { loaded } : {}) - }) + const journal = missingHistory + ? ( + await openAgentSessionJournalWithRecovery({ + identity, + journalDir, + fence: record.lease.runtimeFence, + recordPredatesCall: true + }) + ).journal + : await openAgentSessionJournal({ + identity, + journalDir, + // Omitted, not `null`: the store reads `null` as "replay already ran and + // found nothing" and founds a fresh epoch. In process the probe above is the + // previous statement, so the window is zero-width; this holds the line for a + // database another process creates in between. + ...(loaded ? { loaded } : {}) + }) // Read restore opens the journal and nothing else: no adapter call, so no // provider child. Opening it can still write — a session whose history is in // the old format founds its epoch and commits the row explaining that here. diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts index 4f22b56397c..2141ebce799 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts @@ -2,7 +2,10 @@ 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 { AGENT_SESSION_JOURNAL_SCHEMA_VERSION } from '../../../shared/agent-session-journal-types' +import { + AGENT_JOURNAL_RESET_REASONS, + AGENT_SESSION_JOURNAL_SCHEMA_VERSION +} from '../../../shared/agent-session-journal-types' import type { AgentSessionHandoffStatus, AgentSessionStatusEvent, @@ -35,6 +38,40 @@ afterEach(async () => { }) describe('AgentSessionSubscribers', () => { + // The reason reaches paired and mobile decoders that may reject a value they + // have never seen, so the recovery that caused the reset rides beside it. + it('names the recovery cause without widening the reset reason', async () => { + const journal = await journals.open({ + identity: { + sessionId: SESSION, + workspaceId: 'workspace-1', + hostId: 'local', + agent: 'codex', + providerHandle: { kind: 'codex', threadId: 'thread-1' } + }, + journalDir: join(root, 'reset-cause-journal') + }) + const events: AgentSessionSubscribeEvent[] = [] + const subscribers = new AgentSessionSubscribers() + subscribers.open({ + id: 'subscriber-1', + sessionId: SESSION, + journal, + fence: 3, + emit: (event) => events.push(event) + }) + + subscribers.reset(SESSION, journal, 'epoch_changed', 3, 'journal_missing') + subscribers.reset(SESSION, journal, 'epoch_changed', 3) + + const resets = events.filter((event) => event.type === 'reset') + expect(resets).toHaveLength(2) + expect(resets[0]).toMatchObject({ reset: 'epoch_changed', resetCause: 'journal_missing' }) + expect(AGENT_JOURNAL_RESET_REASONS).toContain(resets[0]?.reset) + // Absent, not null, when the host has no cause to name. + expect(resets[1] && 'resetCause' in resets[1]).toBe(false) + }) + it('publishes the current fence when a resumed cursor is already caught up', async () => { const journal = await journals.open({ identity: { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts index 24415df1d3c..18d3c9f278b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts @@ -18,6 +18,7 @@ import { } from '../../../shared/agent-session-wire' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import { emptyAgentSessionBatch } from './agent-session-empty-batch' +import { agentSessionResetFrame, type AgentSessionResetFrame } from './agent-session-reset-frame' import { createAgentSessionCatchUpReader, readAgentSessionHydrationPage @@ -142,9 +143,11 @@ export class AgentSessionSubscribers { journal: AgentSessionJournal, reason: AgentJournalResetReason, fence: number, + cause?: AgentSessionResetFrame['resetCause'], backgroundTasks?: AgentSessionBackgroundTaskState | null ): void { - this.replay(sessionId, journal, fence, backgroundTasks, { type: 'reset', reset: reason }) + const frame = agentSessionResetFrame(reason, cause) + this.replay(sessionId, journal, fence, backgroundTasks, frame) } snapshot( @@ -161,7 +164,7 @@ export class AgentSessionSubscribers { journal: AgentSessionJournal, fence: number, backgroundTasks: AgentSessionBackgroundTaskState | null | undefined, - frame: { type: 'snapshot' } | { type: 'reset'; reset: AgentJournalResetReason } + frame: { type: 'snapshot' } | AgentSessionResetFrame ): void { const page = readAgentSessionHydrationPage(journal, fence) const hostNow = this.now() diff --git a/src/shared/agent-session-journal-types.ts b/src/shared/agent-session-journal-types.ts index 7ba5b277e6d..41998d686b6 100644 --- a/src/shared/agent-session-journal-types.ts +++ b/src/shared/agent-session-journal-types.ts @@ -279,3 +279,13 @@ export const AGENT_JOURNAL_RESET_REASONS = [ 'schema_unreadable' ] as const export type AgentJournalResetReason = (typeof AGENT_JOURNAL_RESET_REASONS)[number] + +/** What made the reset necessary. Reported ALONGSIDE the reason rather than as + * another reason value: the reason reaches paired and mobile decoders that may + * reject an unknown one, and a reader that ignores this field still resets. */ +export const AGENT_JOURNAL_RESET_CAUSES = [ + 'journal_missing', + 'journal_corrupt', + 'schema_unreadable' +] as const +export type AgentJournalResetCause = (typeof AGENT_JOURNAL_RESET_CAUSES)[number] diff --git a/src/shared/agent-session-wire.ts b/src/shared/agent-session-wire.ts index b4ef57a6b91..76f2c36ceae 100644 --- a/src/shared/agent-session-wire.ts +++ b/src/shared/agent-session-wire.ts @@ -17,6 +17,7 @@ import type { AgentSessionConversationCommand } from './agent-session-conversati import type { AgentJournalCursor, AgentJournalRenderItem, + AgentJournalResetCause, AgentJournalResetReason, AgentJournalResolution, AgentJournalSubmission @@ -129,13 +130,13 @@ export type AgentSessionHistoryResult = | { ok: true; page: AgentSessionHistoryPage; providerSession?: AgentProviderSessionMetadata } /** Every reset carries a byte-bounded tail page so recovery cannot exceed * remote outbound admission or require another call before resubscribing. */ - | { + | ({ ok: false reset: AgentJournalResetReason page: AgentSessionHistoryPage fence?: number providerSession?: AgentProviderSessionMetadata - } + } & AgentSessionResetCauseField) /** Cursor-qualified incremental publication. Items and submissions carry their * CURRENT reduced state rather than a delta, so applying a batch twice @@ -150,6 +151,11 @@ export type AgentSessionJournalBatch = { /** Host wall clock (ms epoch) stamped once per published frame; see `AgentSessionHistoryPage`. */ type AgentSessionHostClockField = { hostNow?: number } +/** Which failure the reset came out of, when the host knows. Absent from older + * hosts and from a reset a cursor mismatch raised, so a reader that renders it + * must fall back to the reason alone. */ +type AgentSessionResetCauseField = { resetCause?: AgentJournalResetCause } + export type AgentSessionSubscribeEvent = | ({ type: 'snapshot' @@ -187,7 +193,8 @@ export type AgentSessionSubscribeEvent = /** Omitted when unchanged; null clears a previous provider catalog. */ commands?: AgentSessionSlashCommand[] | null activity?: AgentSessionTurnActivity | null - } & AgentSessionHostClockField) + } & AgentSessionHostClockField & + AgentSessionResetCauseField) | { type: 'end' } // ─── Status feed ────────────────────────────────────────────────────────────