diff --git a/src/main/ai-vault-search/session-search-index-consumer.test.ts b/src/main/ai-vault-search/session-search-index-consumer.test.ts index 9b42910bdc1..53e29e13e3b 100644 --- a/src/main/ai-vault-search/session-search-index-consumer.test.ts +++ b/src/main/ai-vault-search/session-search-index-consumer.test.ts @@ -10,7 +10,7 @@ import { userMessages, type SessionSearchIndexFile } from './session-search-staged-write-test-fixture' -import { SessionSearchStore } from './session-search-store' +import { SessionSearchStore, STALE_PATH_LIMIT } from './session-search-store' let index: SessionSearchIndexFile let store: SessionSearchStore @@ -182,13 +182,45 @@ it('ignores a candidate older than the retention cutoff', async () => { expect(store.takeStale()).toEqual([]) }) -it('stops writing while the store refuses writes', async () => { +it('stops writing while the store refuses writes, but remembers what it skipped', async () => { store.setAcceptingWrites(false) replayTranscriptRead({ messages: userMessages('paused', 3) }) await store.settled() expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) expect(errors).toEqual([]) + // A pause is exactly the window in which every read is declined. Forgetting + // them would leave the whole paused span unindexed with nothing to replay it. + expect(store.takeStale().map((candidate) => candidate.file.path)).toEqual([SYNTHETIC_TRANSCRIPT]) +}) + +it('keeps the paused re-read set when the retention window is reconfigured', async () => { + store.setAcceptingWrites(false) + replayTranscriptRead({ messages: userMessages('paused', 2) }) + await store.settled() + expect(store.pendingFileCount).toBe(1) + + // Retention is the only thing allowed to prune this set, and this candidate + // is inside the new window. + store.setRetentionCutoffMs(syntheticCandidate().file.mtimeMs - 1000) + expect(store.pendingFileCount).toBe(1) + + // A cutoff that really does exclude it still prunes. + store.setRetentionCutoffMs(Date.now()) + expect(store.pendingFileCount).toBe(0) +}) + +it('drops the oldest record rather than growing without a bound, and says so', () => { + store.setAcceptingWrites(false) + for (let index = 0; index < STALE_PATH_LIMIT + 5; index++) { + store.markStale(syntheticCandidate({ path: `/transcript-${index}.jsonl` })) + } + + expect(store.pendingFileCount).toBe(STALE_PATH_LIMIT) + expect(store.droppedPendingFileCount).toBe(5) + const kept = store.takeStale().map((candidate) => candidate.file.path) + expect(kept).not.toContain('/transcript-0.jsonl') + expect(kept).toContain(`/transcript-${STALE_PATH_LIMIT + 4}.jsonl`) }) it('keeps the session list running when the index write fails', async () => { @@ -294,3 +326,23 @@ it('keeps a proven file identity when a later read cannot stat it', async () => expect(cursor()).toBe(200) expect(store.takeStale()).toHaveLength(1) }) + +it('leaves no batch on disk when a rejected read is the last one before shutdown', async () => { + replayTranscriptRead({ messages: userMessages('indexed', 3), outcome: { byteOffset: 100 } }) + await store.settled() + + // The parser rejects the file, so this read publishes a cursor and tombstones + // its own staged rows. Nothing writes after it. + replayTranscriptRead({ + messages: userMessages('rejected', 4), + outcome: { session: null, byteOffset: 300 } + }) + await store.settled() + + expect(index.db.prepare('SELECT count(*) AS n FROM search_write_batches').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM search_pending_deletes').get()).toEqual({ + n: 0 + }) + expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) +}) diff --git a/src/main/ai-vault-search/session-search-index-consumer.ts b/src/main/ai-vault-search/session-search-index-consumer.ts index 0bfdadb0910..bce9a3097b4 100644 --- a/src/main/ai-vault-search/session-search-index-consumer.ts +++ b/src/main/ai-vault-search/session-search-index-consumer.ts @@ -30,6 +30,10 @@ export class SessionSearchIndexConsumer implements TranscriptConsumer { beginRead(start: TranscriptReadStart): TranscriptReadConsumer | null { const { candidate } = start if (!this.store.acceptsCandidate(candidate)) { + // A pause is a reason not to write now, not a reason to forget the read. + // `markStale` applies the retention rule itself, so a candidate that is + // out of scope rather than merely paused is still dropped here. + this.store.markStale(candidate) return null } // A parser that decodes where the channel cannot reach it reports every read @@ -82,20 +86,27 @@ class SessionSearchReadConsumer implements TranscriptReadConsumer { finish(outcome: TranscriptReadOutcome): void { const { candidate } = this.start + let published = false try { // An incomplete read's rows are not the whole span, so the cursor must not // move past them; the file is re-read whole instead. - if (this.failed || outcome.incomplete || !this.staged.publish(outcome)) { - this.store.markStale(candidate) - return - } - this.store.writePublished(candidate) + published = !this.failed && !outcome.incomplete && this.staged.publish(outcome) } catch (error) { - this.store.markStale(candidate) this.store.reportWriteFailure(error) } finally { + // Before the store is told anything. `discard` is what tombstones the + // staged rows of a read that decoded no session, and the store's cleanup + // lane reads the tombstone table the moment it is scheduled — telling the + // store first left that batch on disk until some later write happened to + // schedule another pass, which for the last read before a shutdown is + // never. this.staged.discard() } + if (published) { + this.store.writePublished(candidate) + return + } + this.store.markStale(candidate) } } diff --git a/src/main/ai-vault-search/session-search-live-transcript.test.ts b/src/main/ai-vault-search/session-search-live-transcript.test.ts index eed7835258a..af787d463d4 100644 --- a/src/main/ai-vault-search/session-search-live-transcript.test.ts +++ b/src/main/ai-vault-search/session-search-live-transcript.test.ts @@ -4,6 +4,7 @@ import { join } from 'node:path' import { afterEach, beforeEach, expect, it } from 'vitest' import { resetSessionParseCacheForTests } from '../ai-vault/session-scanner-parse-cache' import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers' +import { requestWholeTranscriptRead } from '../ai-vault/session-transcript-reader' import SyncDatabase from '../sqlite/sync-database' import { registerSessionSearchIndexConsumer } from './session-search-index-consumer' import { SessionSearchStore } from './session-search-store' @@ -137,3 +138,34 @@ it('never indexes a credential that appeared in tool output', async () => { expect(sessionsMatching(`"${GITHUB_TOKEN}"`)).toEqual([]) expect(sessionsMatching('redacted')).toEqual([SESSION_ID]) }) + +it('indexes a file the session list already read past, once a whole read is asked for', async () => { + const root = await makeTempDir() + const path = join(root, `${SESSION_ID}.jsonl`) + await writeFile(path, `${userRecord(0, 'the opening prompt')}\n`) + + // The state on first enablement inside a running app: the session list has + // read this file, so the parse cache is warm, while the index is empty. + resetTranscriptConsumersForTests() + await parseTranscript(path) + registerSessionSearchIndexConsumer(store) + + await appendFile(path, `${assistantRecord(1, 'a zygomorphic reply')}\n`) + const appended = await parseTranscript(path) + expect(appended.stats).toMatchObject({ incremental: 1, fullParses: 0 }) + // The append continued from a byte offset the index never saw, so it declined. + expect(sessionsMatching('zygomorphic')).toEqual([]) + + const behind = store.takeStale() + expect(behind.map((candidate) => candidate.file.path)).toEqual([path]) + for (const candidate of behind) { + requestWholeTranscriptRead(candidate.file.path) + } + + const reread = await parseTranscript(path) + expect(reread.stats).toMatchObject({ incremental: 0, fullParses: 1 }) + expect(errors).toEqual([]) + expect(sessionsMatching('zygomorphic')).toEqual([SESSION_ID]) + expect(sessionsMatching('opening')).toEqual([SESSION_ID]) + expect(store.takeStale()).toEqual([]) +}) diff --git a/src/main/ai-vault-search/session-search-store.ts b/src/main/ai-vault-search/session-search-store.ts index 514a7d7b7aa..b4f0d3db5e4 100644 --- a/src/main/ai-vault-search/session-search-store.ts +++ b/src/main/ai-vault-search/session-search-store.ts @@ -13,6 +13,11 @@ import { warmSessionSearchPages } from './session-search-page-warmup' import { deleteExpiredSearchFiles } from './session-search-retention-delete' import { openSessionSearchDatabase } from './session-search-schema' +// A paused store keeps recording what it declined, so the set needs a ceiling. +// Above it the oldest record goes and the drop is counted, because a re-read set +// that silently forgets is worse than one that says it is incomplete. +export const STALE_PATH_LIMIT = 20_000 + export type SessionSearchStoreOptions = { /** The WAL backlog a staging write refuses to grow past. Only tests narrow it. */ walBudgetBytes?: number @@ -37,6 +42,7 @@ export class SessionSearchStore { // Files this index knows it is behind on. Filled by a declined or abandoned // read; PR 3's indexer drains it. Nothing here schedules the re-read. private readonly stale = new Map() + private droppedStalePaths = 0 constructor( path: string, @@ -59,18 +65,22 @@ export class SessionSearchStore { setRetentionCutoffMs(cutoffMs: number | null): void { this.retentionCutoffMs = cutoffMs for (const [path, candidate] of this.stale) { - if (!this.acceptsCandidate(candidate)) { + // Only retention prunes the re-read set. Pausing is a reason not to write + // now, never a reason to forget what still has to be read. + if (!this.withinRetention(candidate)) { this.stale.delete(path) } } } + /** Whether this candidate is new enough to be worth holding rows for at all. */ + private withinRetention(candidate: SessionFileCandidate): boolean { + return this.retentionCutoffMs === null || candidate.file.mtimeMs >= this.retentionCutoffMs + } + + /** Whether a write for this candidate may start right now. */ acceptsCandidate(candidate: SessionFileCandidate): boolean { - return ( - !this.closed && - this.acceptingWrites && - (this.retentionCutoffMs === null || candidate.file.mtimeMs >= this.retentionCutoffMs) - ) + return !this.closed && this.acceptingWrites && this.withinRetention(candidate) } indexedFile(path: string, identity: SessionSearchFileIdentity): SessionSearchIndexedFile | null { @@ -112,14 +122,48 @@ export class SessionSearchStore { this.onError(error) } - /** Records a file whose content the index is behind on, for a later whole re-read. */ + /** + * Records a file whose content the index is behind on, for a later whole + * re-read. Recorded while paused too: a pause is exactly the window in which + * reads are declined, so refusing to remember them would lose every file the + * pause covered. + */ markStale(candidate: SessionFileCandidate): void { - if (this.acceptsCandidate(candidate)) { - this.stale.set(candidate.file.path, candidate) + if (this.closed || !this.withinRetention(candidate)) { + return + } + // Re-inserting moves the path to the end, so the oldest record is the one + // dropped when a long pause overruns the bound. + this.stale.delete(candidate.file.path) + this.stale.set(candidate.file.path, candidate) + while (this.stale.size > STALE_PATH_LIMIT) { + const oldest = this.stale.keys().next() + if (oldest.done) { + break + } + this.stale.delete(oldest.value) + this.droppedStalePaths += 1 } } - /** Hands the re-read set to its scheduler and clears it. */ + /** + * Files the index knew it was behind on and could not keep a record of. A + * non-zero count means the re-read set is incomplete, so coverage cannot be + * reported as whole until a full pass runs. + */ + get droppedPendingFileCount(): number { + return this.droppedStalePaths + } + + /** + * Hands the re-read set to its scheduler and clears it. + * + * These paths are behind, not merely dirty: the index declined their last read + * because it covered a span the index never saw. Re-dispatching a scan is not + * enough on its own, because the reader picks `append` from the session list's + * resume point and the consumer will decline again. The caller must pass each + * path to `requestWholeTranscriptRead` first. + */ takeStale(): SessionFileCandidate[] { const candidates = [...this.stale.values()] this.stale.clear() diff --git a/src/main/ai-vault/session-transcript-reader.ts b/src/main/ai-vault/session-transcript-reader.ts index 228b0c832a3..2fb7d0dbc47 100644 --- a/src/main/ai-vault/session-transcript-reader.ts +++ b/src/main/ai-vault/session-transcript-reader.ts @@ -3,7 +3,10 @@ import type { AiVaultSession } from '../../shared/ai-vault-types' import { parseAgentSessionFile, parserPublishesMessages } from './session-scanner-agent-parser' import { consumeCompleteJsonlLines } from './session-scanner-jsonl-reader' import type { ResumableSessionParseState, SessionFileCandidate } from './session-scanner-types' -import type { SessionParseResumePoint } from './session-parse-cache-store' +import { + invalidateSessionParseCacheEntry, + type SessionParseResumePoint +} from './session-parse-cache-store' import { TranscriptMessageChannel } from './session-transcript-channel' const NEWLINE_BYTE = 0x0a @@ -29,6 +32,24 @@ export type ResumableTranscriptRead = { resume: SessionParseResumePoint } +/** + * Ask for the next read of `path` to be a whole-file `replace`. + * + * Why this lives here: a consumer never chooses its own mode. The reader picks + * `append` or `replace` from the resume point the session list left behind, so a + * consumer that declined an append has no way to get the span it missed — with + * an empty index and a warm parse cache, every read arrives as `append`, every + * one is declined, and nothing is ever indexed. Dropping the resume point is the + * one lever that changes the next read's mode, and only the reader's own cache + * owns it. + * + * The cost is a re-parse for the session list too. That is the honest price of a + * second consumer being behind, and it is paid once per file rather than per scan. + */ +export function requestWholeTranscriptRead(path: string): void { + invalidateSessionParseCacheEntry(path) +} + /** * Read an append-only transcript, resuming from `resume` when the file only * grew and the recorded offset still sits on a line boundary. Anything else