mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 08:01:56 +00:00
fix(native-chat): treat a record without its journal as recovery, not a fresh session
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.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<string | undefined> {
|
||||
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()
|
||||
|
||||
|
||||
@@ -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<AgentSessionJournalOpened> {
|
||||
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 }
|
||||
}
|
||||
|
||||
@@ -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 } : {}) }
|
||||
}
|
||||
@@ -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: (
|
||||
|
||||
+7
-1
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
+53
-8
@@ -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<string> {
|
||||
|
||||
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()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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.
|
||||
|
||||
+38
-1
@@ -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: {
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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 ────────────────────────────────────────────────────────────
|
||||
|
||||
Reference in New Issue
Block a user