diff --git a/src/main/ai-vault-search/session-search-backfill-remainder.ts b/src/main/ai-vault-search/session-search-backfill-remainder.ts index d55cc5a8bd2..a95befe3b44 100644 --- a/src/main/ai-vault-search/session-search-backfill-remainder.ts +++ b/src/main/ai-vault-search/session-search-backfill-remainder.ts @@ -49,6 +49,11 @@ export class SessionSearchBackfillRemainder { * at a time. Rolling the remainder over rather than re-discovering keeps a * first run off the 33 s cold discovery every pass. * + * This runs inside the reconcile cycle, so `overdue` is what keeps the + * recent-N-per-interval promise: without it the pacer's load back-off can + * hold one pass for far longer than the interval and every recency check + * queues behind it. + * * Returns true once a plan too large for the queue has been read through: * discovery still owes the files the queue could not hold, and one sweep * settles that debt. @@ -58,7 +63,8 @@ export class SessionSearchBackfillRemainder { allowance: SessionSearchCycleAllowance, status: SessionSearchIndexingStatus, pace: (signal?: AbortSignal) => Promise, - signal: AbortSignal + signal: AbortSignal, + overdue?: () => boolean ): Promise { if (this.entries.size === 0) { return false @@ -71,6 +77,7 @@ export class SessionSearchBackfillRemainder { signal, pace, allowance, + overdue, onIndexed: (_candidate, bytes) => status.indexed(bytes), onFailed: () => status.failed() }) diff --git a/src/main/ai-vault-search/session-search-backfill.ts b/src/main/ai-vault-search/session-search-backfill.ts index 9d104af82b4..4e30c28ec51 100644 --- a/src/main/ai-vault-search/session-search-backfill.ts +++ b/src/main/ai-vault-search/session-search-backfill.ts @@ -26,6 +26,7 @@ import { type SessionSearchScanRoots } from './session-search-scan-roots' import type { SessionSearchStore } from './session-search-store' +import { sessionSearchEnumeratedContainers } from './session-search-synthetic-sources' // One walk each, for paths the sweep did not discover. Normally near zero; the // cap is there for the case that is not normal, an unmounted tree, where the @@ -42,6 +43,8 @@ export type SessionSearchBackfillArgs = { previousRootsWithFiles?: ReadonlySet /** This pass's reading allowance; what does not fit comes back as `deferred`. */ allowance?: SessionSearchCycleAllowance + /** True once the pass is out of wall time; the rest comes back as `deferred`. */ + overdue?: () => boolean listings: SessionSearchDirectoryReader pace?: (signal?: AbortSignal) => Promise signal?: AbortSignal @@ -102,6 +105,7 @@ export async function runSessionSearchBackfill( signal, pace: args.pace, allowance: args.allowance, + overdue: args.overdue, onIndexed: (_candidate, bytes) => status.indexed(bytes), onFailed: () => status.failed() }) @@ -149,6 +153,9 @@ export async function runSessionSearchBackfill( store, paths: undiscovered, roots, + // Only a sweep enumerates without a per-agent limit, so only a sweep + // may prove a synthetic row's container holds it no longer. + enumeratedContainers: sessionSearchEnumeratedContainers(swept.candidates, issues), emptiedRoots: previousRootsWithFiles ? sessionSearchEmptiedRoots(previousRootsWithFiles, rootsWithFiles) : new Set(), diff --git a/src/main/ai-vault-search/session-search-deleted-sources.test.ts b/src/main/ai-vault-search/session-search-deleted-sources.test.ts index 418f70fe1e7..9f3df735858 100644 --- a/src/main/ai-vault-search/session-search-deleted-sources.test.ts +++ b/src/main/ai-vault-search/session-search-deleted-sources.test.ts @@ -1,6 +1,7 @@ import { chmod, mkdir, rm, writeFile } from 'node:fs/promises' import { join } from 'node:path' import { afterEach, beforeEach, expect, it } from 'vitest' +import { parserPublishesMessages } from '../ai-vault/session-scanner-agent-parser' import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers' import { retireDeletedSessionSearchSources } from './session-search-deleted-sources' import { @@ -54,6 +55,7 @@ function retire( options: { roots?: readonly string[] emptiedRoots?: ReadonlySet + enumeratedContainers?: ReadonlyMap> listings?: SessionSearchDirectoryReader limit?: number } = {} @@ -63,6 +65,7 @@ function retire( paths, roots: options.roots ?? [harness.roots.claudeProjectsDir ?? ''], emptiedRoots: options.emptiedRoots, + enumeratedContainers: options.enumeratedContainers, listings: options.listings ?? new SessionSearchDirectoryListings(), limit: options.limit }) @@ -226,14 +229,71 @@ it('judges a row under no configured root by its own directory', async () => { expect(result.degradedRoots).toEqual([]) }) -it('stats the database behind a synthetic OpenCode row, not the row itself', async () => { +// I8. A synthetic row names a container and an entry inside it. Walking the +// row's own path would report every one of them gone, and walking only the +// container proves nothing about the entry: a session deleted inside a database +// that is still there would never be retired at all. +it('proves a synthetic row against its container, not against its own path', async () => { const db = join(harness.root, 'opencode.db') await writeFile(db, '') - const alive = `${db}#session-1` - await expect(retire([alive], { roots: [] })).resolves.toMatchObject({ retired: [] }) + const kept = `${db}#session-1` + const deleted = `${db}#session-2` + const enumeratedContainers = new Map([[db, new Set(['session-1'])]]) + + const result = await retire([kept, deleted], { roots: [], enumeratedContainers }) + expect(result.retired).toEqual([deleted]) + expect(result.unverifiable).toEqual([]) +}) + +it('keeps a synthetic row when this pass did not enumerate its container', async () => { + const db = join(harness.root, 'opencode.db') + await writeFile(db, '') + const row = `${db}#session-1` + + // A cycle asks for the newest N per agent, so a row it did not return may be + // the one after them. It enumerates nothing and therefore proves nothing. + await expect(retire([row], { roots: [] })).resolves.toMatchObject({ + retired: [], + unverifiable: [row] + }) + + // An enumeration that returned nothing at all is not evidence either: a + // database whose schema this scanner no longer recognises reads as empty + // with no error, and believing it would retire every session in one pass. + await expect( + retire([row], { roots: [], enumeratedContainers: new Map([[db, new Set()]]) }) + ).resolves.toMatchObject({ retired: [], unverifiable: [row] }) +}) + +it('retires a synthetic row when the container it came from is gone', async () => { + const db = join(harness.root, 'opencode.db') + await writeFile(db, '') + const row = `${db}#session-1` + const enumeratedContainers = new Map([[db, new Set(['session-1'])]]) + await expect(retire([row], { roots: [], enumeratedContainers })).resolves.toMatchObject({ + retired: [] + }) await rm(db) - await expect(retire([alive], { roots: [] })).resolves.toMatchObject({ retired: [alive] }) + await expect(retire([row], { roots: [], enumeratedContainers })).resolves.toMatchObject({ + retired: [row] + }) +}) + +// Nothing in this PR can hold a synthetic row: the index pass refuses a source +// whose parser decodes its messages where the message channel cannot reach +// them, and OpenCode's SQLite sessions are read on a worker thread. The rule +// above is the guard for the day that changes -- without it the walk would read +// `#` as a filename and retire every such row the moment it appeared. +it('does not index a source whose messages the channel cannot reach', () => { + const db = join(harness.root, 'opencode.db') + expect( + parserPublishesMessages({ + agent: 'opencode', + codexHome: null, + file: { path: `${db}#session-1`, mtimeMs: 1, modifiedAt: '', sizeBytes: 0 } + }) + ).toBe(false) }) it('caps the walks one pass spends and leaves the rest to be checked again', async () => { diff --git a/src/main/ai-vault-search/session-search-deleted-sources.ts b/src/main/ai-vault-search/session-search-deleted-sources.ts index 700c84d6a03..facb3479c2f 100644 --- a/src/main/ai-vault-search/session-search-deleted-sources.ts +++ b/src/main/ai-vault-search/session-search-deleted-sources.ts @@ -1,8 +1,8 @@ import { basename, dirname } from 'node:path' -import { splitOpenCodeSqliteCandidate } from '../ai-vault/session-scanner-opencode-sqlite-paths' import type { SessionSearchDegradedRoot } from './session-search-degraded-roots' import type { SessionSearchDirectoryReader } from './session-search-directory-listings' import { isUnderScanRoot } from './session-search-scan-roots' +import { splitSyntheticSessionSource } from './session-search-synthetic-sources' import type { SessionSearchStore } from './session-search-store' /* @@ -22,6 +22,10 @@ import type { SessionSearchStore } from './session-search-store' * above it. * I4. A file, or a project directory, the user really deleted retires on the * first pass that proves it. There is no waiting period and no census. + * I8. A row whose path names an entry inside a container rather than a file of + * its own is proven the same way, one level up: the container must be + * present, and the pass must have enumerated it in full and successfully. + * A listing is a listing whether it comes from readdir or from a database. * * What I3 costs, stated rather than hidden: a volume mounted at exactly a * configured root, unmounted so that the mountpoint stays present and lists @@ -56,6 +60,12 @@ export type SessionSearchRetirementArgs = { roots: readonly string[] /** Roots that listed transcripts on the previous pass and none on this one. */ emptiedRoots?: ReadonlySet + /** + * Containers this pass enumerated in full, with the ids each holds. Only a + * census builds it; see session-search-synthetic-sources.ts for the bar a + * container has to meet before it appears here. + */ + enumeratedContainers?: ReadonlyMap> /** One readdir per directory per pass, shared with the rest of the pass. */ listings: SessionSearchDirectoryReader /** Directories walked before the pass moves on; the rest stay watched. */ @@ -100,15 +110,19 @@ export async function retireDeletedSessionSearchSources( retirement.unchecked.push(...paths.slice(index)) break } - // A synthetic OpenCode row names the database it came from, never a file of - // its own; walking the candidate path would report every one of them gone. - const filePath = splitOpenCodeSqliteCandidate(path)?.dbPath ?? path + // A synthetic row names a container and an entry inside it, never a file of + // its own; walking the row's own path would report every one of them gone. + const synthetic = splitSyntheticSessionSource(path) + const filePath = synthetic?.container ?? path const root = configuredRootFor(filePath, args.roots) - const proof = await proveSource(filePath, root ?? dirname(filePath), { + const containerProof = await proveSource(filePath, root ?? dirname(filePath), { listings: args.listings, emptiedRoots, signal }) + const proof = synthetic + ? proveSyntheticSource(synthetic, containerProof, args.enumeratedContainers) + : containerProof if (proof.verdict === 'gone') { store.removeFile(path) retirement.retired.push(path) @@ -179,6 +193,39 @@ async function proveSource( return { verdict: 'unverifiable', reason: `${root} could not be listed.` } } +/** + * A synthetic row is proven by its container's own enumeration, one level above + * where the filesystem walk stops. + * + * The container has to be present first: a database on a volume that is not + * there proves nothing about the sessions inside it, and a database that is + * gone takes its sessions with it. Only then does the enumeration decide, and + * only when this pass made one that was exhaustive and successful -- a cycle + * asks for the newest N per agent, so an id it did not return may just be the + * one after them. + */ +function proveSyntheticSource( + synthetic: { container: string; id: string }, + containerProof: SessionSearchSourceVerdict, + enumerated?: ReadonlyMap> +): SessionSearchSourceVerdict { + if (containerProof.verdict !== 'present') { + return containerProof + } + const ids = enumerated?.get(synthetic.container) + // An enumeration that returned nothing at all is not evidence that the + // container holds nothing: a source whose schema this scanner no longer + // recognises reads as empty with no error to see, and believing it would + // retire every entry in one pass. + if (!ids || ids.size === 0) { + return { + verdict: 'unverifiable', + reason: `${synthetic.container} was not enumerated in full this pass.` + } + } + return ids.has(synthetic.id) ? { verdict: 'present' } : { verdict: 'gone' } +} + /** Longest configured root containing the path, or null for a row under none. */ function configuredRootFor(path: string, roots: readonly string[]): string | null { let owner: string | null = null diff --git a/src/main/ai-vault-search/session-search-index-pass.ts b/src/main/ai-vault-search/session-search-index-pass.ts index 983ca027b30..5fac880fc69 100644 --- a/src/main/ai-vault-search/session-search-index-pass.ts +++ b/src/main/ai-vault-search/session-search-index-pass.ts @@ -14,6 +14,16 @@ export type SessionSearchIndexPassOptions = { signal?: AbortSignal /** Cycle allowance; work that does not fit comes back as `deferred`. */ allowance?: SessionSearchCycleAllowance + /** + * True when the pass has run out of wall time and must hand the rest back. + * + * Why a second bound at all: the allowance counts files and bytes, and the + * pacer sleeps for up to 15 s a batch when the host is loaded, so a pass that + * never exceeds its byte budget can still hold the loop for a quarter of an + * hour. Checked before the allowance is spent, so nothing is charged for work + * this pass will not do. + */ + overdue?: () => boolean /** Paths whose stored cursor must not be trusted, so the read is forced whole. */ forced?: ReadonlySet /** Sleeps between batches so an unasked backfill never owns the CPU. */ @@ -66,7 +76,7 @@ export async function runSessionSearchIndexPass( // cycle's writers appended, and re-statting every file to close it would // cost more than the number is worth. const bytes = forced ? (candidate.file.sizeBytes ?? 0) : unreadBytes(store, candidate) - if (options.allowance && !options.allowance.spend(bytes)) { + if (options.overdue?.() === true || (options.allowance && !options.allowance.spend(bytes))) { deferred.push(...candidates.slice(index)) break } diff --git a/src/main/ai-vault-search/session-search-indexer.test.ts b/src/main/ai-vault-search/session-search-indexer.test.ts index cd81b22be6a..9894ba3be2e 100644 --- a/src/main/ai-vault-search/session-search-indexer.test.ts +++ b/src/main/ai-vault-search/session-search-indexer.test.ts @@ -1,3 +1,4 @@ +import { existsSync } from 'node:fs' import { appendFile, chmod, mkdir, rename, rm, stat, utimes } from 'node:fs/promises' import { join } from 'node:path' import { afterEach, beforeEach, expect, it } from 'vitest' @@ -21,6 +22,7 @@ const INTERVAL_MS = 20_000 const CAN_DENY_READ = process.platform !== 'win32' && process.getuid?.() !== 0 const SESSION_ID = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee' const OTHER_SESSION_ID = 'bbbbbbbb-cccc-4ddd-8eee-ffffffffffff' +const SETTLED_SESSION_ID = 'dddddddd-cccc-4ddd-8eee-ffffffffffff' let harness: SessionSearchIndexerHarness let clock: FakeSessionSearchClock @@ -97,6 +99,23 @@ function transcriptPath(name = SESSION_ID): string { return join(harness.claudeProjectDir, `${name}.jsonl`) } +/** + * Starts the indexer over a root that already holds one indexed transcript, so + * the opening sweep is behind us and `reconcile()` runs a cycle. It is dated + * ahead of everything the caller writes afterwards, so it stays inside any + * recency window and is skipped rather than read. + */ +async function startAfterASweep( + overrides: Partial[0]> = {} +): Promise { + const settled = transcriptPath(SETTLED_SESSION_ID) + await writeClaudeTranscript(settled, ['a conversation from before'], SETTLED_SESSION_ID) + // Wall time, not the fake clock: recency is decided by real file mtimes. + const ahead = new Date(Date.now() + 3_600_000) + await utimes(settled, ahead, ahead) + await newIndexer(overrides).start() +} + /** Advances one reconcile interval and waits for the cycle it fires. */ async function nextCycle(): Promise { clock.advance(INTERVAL_MS) @@ -285,6 +304,9 @@ it.skipIf(!CAN_DENY_READ)( ) it('spends a cycle budget and rolls the rest into the next cycle', async () => { + // The sweep is behind us, so this is the reconciler having to fit four files + // into a two-file allowance. + await startAfterASweep({ budget: { files: 2, bytes: 64 * 1024 } }) for (let index = 0; index < 4; index++) { await writeClaudeTranscript( transcriptPath(`0000000${index}-bbbb-4ccc-8ddd-eeeeeeeeeeee`), @@ -292,11 +314,8 @@ it('spends a cycle budget and rolls the rest into the next cycle', async () => { `0000000${index}-bbbb-4ccc-8ddd-eeeeeeeeeeee` ) } - // Start with the index off so the first pass through the reconciler is the - // one that has to fit four files into a two-file allowance. - newIndexer({ budget: { files: 2, bytes: 64 * 1024 } }) await indexer?.reconcile() - expect(indexer?.status().filesIndexed).toBe(2) + expect(indexer?.status().filesIndexed).toBe(3) expect(indexer?.status().filesPending).toBe(2) await indexer?.reconcile() @@ -308,13 +327,13 @@ it('spends a cycle budget and rolls the rest into the next cycle', async () => { // session list has been scanning since launch, so every transcript already has // a cursor sitting at its current stat and the index has nothing at all. it('fills an empty index over a warm session-list cache on the first reconcile', async () => { + await startAfterASweep() const path = transcriptPath() await writeClaudeTranscript(path, ['scanned before the index existed'], SESSION_ID) // An ordinary parse now reuses its cached fold and opens no file, so no // consumer is asked and there is nothing for a decline to record. await parseTranscript(path) - newIndexer() await indexer?.reconcile() expect(sessionsMatching('scanned')).toEqual([SESSION_ID]) @@ -403,6 +422,8 @@ it('moves the retention window with the clock instead of freezing it at construc // Finding 3: an invalidated path outside the recency window was resolved only // through rows the index already held, so anything else was dropped unread. it('reads an invalidated transcript from outside the recency window', async () => { + // Newest-one per root, and the index has never seen either file below. + await startAfterASweep({ recentPerAgent: 1 }) const older = transcriptPath(OTHER_SESSION_ID) await writeClaudeTranscript(older, ['the older conversation'], OTHER_SESSION_ID) await writeClaudeTranscript(transcriptPath(), ['the newer conversation'], SESSION_ID) @@ -410,8 +431,6 @@ it('reads an invalidated transcript from outside the recency window', async () = const ahead = new Date(newer.mtimeMs + 60_000) await utimes(transcriptPath(), ahead, ahead) - // Newest-one per root, and the index has never seen either file. - newIndexer({ recentPerAgent: 1 }) indexer?.invalidate([older]) await indexer?.reconcile() @@ -444,14 +463,14 @@ it('reports closed once it is closed, whatever it was doing before', async () => // that stat writes a cursor describing a file that no longer looks like this, // so the next cycle distrusts it and re-reads it, forever. it('reads a queued file at its current stat, not the one it was queued with', async () => { + // One file per cycle, so the older one is deferred carrying this stat. + await startAfterASweep({ budget: { files: 1, bytes: 64 * 1024 } }) const older = transcriptPath(OTHER_SESSION_ID) await writeClaudeTranscript(older, ['the deferred conversation'], OTHER_SESSION_ID) await writeClaudeTranscript(transcriptPath(), ['the newer conversation'], SESSION_ID) const ahead = new Date((await stat(transcriptPath())).mtimeMs + 60_000) await utimes(transcriptPath(), ahead, ahead) - // One file per cycle, so the older one is deferred carrying this stat. - newIndexer({ budget: { files: 1, bytes: 64 * 1024 } }) await indexer?.reconcile() expect(indexer?.status().filesPending).toBe(1) @@ -727,23 +746,27 @@ it('counts a file queued in both places once', async () => { // queue running. A `clear()` queued a moment earlier would then delete the // database, open a new one and register a consumer against it, all behind an // indexer whose caller had finished with it. -it('lets nothing queued before close reopen the store', async () => { +// +// F2 is the other half of the same seam: the queued task never running is what +// `close()` is for, but `clear()` had already resolved, so a caller who asked +// for the index to be thrown away was left with it on disk. The close finishes +// the removal itself rather than reopening anything. +it('finishes a clear that was still queued when it was closed', async () => { await writeClaudeTranscript(transcriptPath(), ['indexed before the close'], SESSION_ID) await newIndexer().start() + expect(existsSync(harness.databasePath)).toBe(true) // `clear()` queues the work that removes the database and opens a new one. const clearing = indexer?.clear() indexer?.close() await clearing - expect(sessionsMatching('before')).toEqual([SESSION_ID]) - - // And nothing is registered as a consumer any more, so a transcript written - // after the close gets no row however many times it is parsed. + // Removed, and not reopened: a close leaves no store and no consumer behind. + expect(existsSync(harness.databasePath)).toBe(false) const after = transcriptPath(OTHER_SESSION_ID) await writeClaudeTranscript(after, ['written after the close'], OTHER_SESSION_ID) await parseTranscript(after) - expect(sessionsMatching('after')).toEqual([]) + expect(existsSync(harness.databasePath)).toBe(false) }) // I7: the backfill reads transcript bytes, so it is budgeted like every other @@ -904,3 +927,81 @@ it('never reports more files indexed than the total behind them', async () => { await nextCycle() expect(indexer?.status().filesTotal).toBeNull() }) + +// F1: `fullSweepDue` stayed set across the sweep's await and was cleared on the +// way out, so a request raised while a sweep was running was erased by the +// sweep it arrived during. The pass takes the flag on entry now, and an +// unfinished sweep is what puts it back. +it('indexes a transcript the history window widened in during a sweep', async () => { + for (let index = 0; index < 20; index++) { + const session = `0000${String(index).padStart(4, '0')}-bbbb-4ccc-8ddd-eeeeeeeeeeee` + await writeClaudeTranscript(transcriptPath(session), [`recent session ${index}`], session) + } + const old = transcriptPath(OTHER_SESSION_ID) + await writeClaudeTranscript(old, ['an ancient conversation'], OTHER_SESSION_ID) + const longAgo = new Date(Date.now() - 120 * 86_400_000) + await utimes(old, longAgo, longAgo) + + let paced = 0 + newIndexer({ + historyDays: 30, + // Widen part way through the sweep, which is when a user flips the setting. + pace: async () => { + if (paced++ === 0) { + void indexer?.setHistoryDays(null) + } + } + }) + await indexer?.start() + await indexer?.settled() + + expect(sessionsMatching('ancient')).toEqual([OTHER_SESSION_ID]) +}) + +it('runs another sweep when one is asked for during a sweep', async () => { + for (let index = 0; index < 20; index++) { + const session = `0000${String(index).padStart(4, '0')}-bbbb-4ccc-8ddd-eeeeeeeeeeee` + await writeClaudeTranscript(transcriptPath(session), [`recent session ${index}`], session) + } + const late = transcriptPath(OTHER_SESSION_ID) + + let paced = 0 + // Newest-one per root, so nothing but a second sweep can reach the file that + // appears after this sweep's discovery has already run. + newIndexer({ + recentPerAgent: 1, + pace: async () => { + if (paced++ > 0) { + return + } + await writeClaudeTranscript(late, ['a late conversation'], OTHER_SESSION_ID) + const backdated = new Date(Date.now() - 86_400_000) + await utimes(late, backdated, backdated) + void indexer?.reconcile({ full: true }) + } + }) + await indexer?.start() + await indexer?.settled() + + expect(sessionsMatching('late')).toEqual([OTHER_SESSION_ID]) +}) + +// F4: the allowance counts files and bytes, and the pacer sleeps for up to 15 s +// a batch on a loaded host, so a pass that never came near its byte budget +// could hold the loop for a quarter of an hour. Since the backfill drains +// inside the reconcile cycle, that is the recency promise gone. +it('hands the rest of a pass back when it runs out of wall time', async () => { + for (let index = 0; index < 20; index++) { + const session = `0000${String(index).padStart(4, '0')}-bbbb-4ccc-8ddd-eeeeeeeeeeee` + await writeClaudeTranscript(transcriptPath(session), [`paced session ${index}`], session) + } + // Half an interval per pacing point, which is what a loaded host's back-off + // costs. Eight files a batch, so the third batch is over the line. + await newIndexer({ pace: async () => clock.advance(INTERVAL_MS / 2) }).start() + + expect(indexer?.status()).toMatchObject({ filesIndexed: 16, filesPending: 4 }) + + // And the pass after it picks up exactly what was handed back. + await nextCycle() + expect(indexer?.status()).toMatchObject({ filesIndexed: 20, filesPending: 0 }) +}) diff --git a/src/main/ai-vault-search/session-search-indexer.ts b/src/main/ai-vault-search/session-search-indexer.ts index b60c752f9ef..61f9bdfbd92 100644 --- a/src/main/ai-vault-search/session-search-indexer.ts +++ b/src/main/ai-vault-search/session-search-indexer.ts @@ -23,13 +23,8 @@ import { SessionSearchCycleAllowance } from './session-search-reconcile-budget' import { runSessionSearchReconcileCycle } from './session-search-reconciler' -import { - narrowsSessionSearchHistory, - sessionSearchHistoryCutoffMs, - widensSessionSearchHistory -} from './session-search-retention-policy' +import { SessionSearchRetentionWindow } from './session-search-retention-policy' import { SessionSearchRootRecovery } from './session-search-root-recovery' -import { removeSessionSearchDatabase } from './session-search-schema' import { SessionSearchRegisteredStore } from './session-search-registered-store' import type { SessionSearchStore } from './session-search-store' import { SessionSearchWorkLoop } from './session-search-work-loop' @@ -62,11 +57,11 @@ export class SessionSearchIndexer { private readonly pace: (signal?: AbortSignal) => Promise private readonly loop: SessionSearchWorkLoop - private readonly registered = new SessionSearchRegisteredStore() + private readonly registered: SessionSearchRegisteredStore private previousRecent = new Set() private readonly rootRecovery = new SessionSearchRootRecovery() private readonly backfillRemainder = new SessionSearchBackfillRemainder() - private historyDays: number | null + private readonly retention: SessionSearchRetentionWindow private started = false private paused = false private purgeDue = false @@ -88,7 +83,8 @@ export class SessionSearchIndexer { ) this.onError = options.onError ?? ((error) => console.warn('[ai-vault-search]', error)) this.pace = options.pace ?? pauseBackfill - this.historyDays = options.historyDays + this.registered = new SessionSearchRegisteredStore(options.databasePath, this.onError) + this.retention = new SessionSearchRetentionWindow(options.historyDays) this.loop = new SessionSearchWorkLoop({ clock: this.clock, intervalMs: this.intervalMs, @@ -108,6 +104,7 @@ export class SessionSearchIndexer { } this.started = true this.fullSweepDue = true + this.indexingStatus.setStarted() // Started while paused: `resume()` is what arms the timer. Queueing a pass // that returns immediately would resolve this call as though one had run. return this.paused ? this.loop.settled : this.tick() @@ -149,9 +146,12 @@ export class SessionSearchIndexer { } this.loop.disarm() this.loop.abort() + // Recorded before it is queued: `close()` performs whatever is still owed, + // because a caller who asked for the index to be thrown away and got a + // resolved promise back must not be left with the database on disk. + this.registered.requestRemoval() this.loop.queue(async () => { this.registered.close() - removeSessionSearchDatabase(this.options.databasePath) this.pending.clear() this.previousRecent = new Set() this.rootRecovery.reset() @@ -163,9 +163,16 @@ export class SessionSearchIndexer { return this.started && !this.paused ? this.tick() : this.loop.settled } - /** Runs one pass now, off the timer. A full pass sweeps every root. */ + /** + * Runs one pass now, off the timer. A full pass sweeps every root. + * + * Refused before `start()`. A pass run against an indexer nobody started + * writes the index once and then leaves it to go stale, because there is no + * timer to arm and nothing to notice the next change; a caller that wants one + * pass wants `start()`. + */ reconcile(options: { full?: boolean } = {}): Promise { - if (this.closed) { + if (this.closed || !this.started) { return this.loop.settled } this.fullSweepDue ||= options.full === true @@ -196,11 +203,10 @@ export class SessionSearchIndexer { if (this.closed) { return this.loop.settled } - const previous = this.historyDays - this.historyDays = historyDays + const change = this.retention.moveTo(historyDays) const cutoffMs = this.cutoffMs() this.store?.setRetentionCutoffMs(cutoffMs) - if (narrowsSessionSearchHistory(previous, historyDays)) { + if (change === 'purge') { if (this.paused) { this.purgeDue = true return this.loop.settled @@ -209,7 +215,7 @@ export class SessionSearchIndexer { await this.store?.purgeOlderThan(cutoffMs, signal) }) } - if (!widensSessionSearchHistory(previous, historyDays)) { + if (change !== 'resweep') { return this.loop.settled } this.fullSweepDue = true @@ -235,6 +241,8 @@ export class SessionSearchIndexer { // otherwise still run, and `clear()`'s task reopens the store. Closing the // loop makes every queued task a no-op, so nothing can register a consumer // or open a database against an indexer the caller has finished with. + // The store, its registration, and a removal a queued `clear()` will now + // never perform, because closing the loop is what stops that task running. this.loop.close() this.registered.close() } @@ -266,22 +274,33 @@ export class SessionSearchIndexer { this.purgeDue = false // One readdir per directory for the whole pass, shared by everything in it. const listings = new SessionSearchDirectoryListings() + // A pass may hold the loop for one interval and no longer. The allowance + // bounds bytes; this bounds wall time, which is what the pacer's load + // back-off spends without spending a single byte of budget. + const startedAt = this.clock.now() + const overdue = (): boolean => this.clock.now() - startedAt >= this.intervalMs if (this.fullSweepDue) { + // Cleared before the sweep runs, not after: a widening or an explicit + // `reconcile({ full: true })` raised while this one is in flight sets the + // flag again, and clearing it on the way out would erase that request + // along with this pass's own. An unfinished sweep sets it back itself. + this.fullSweepDue = false // A sweep opens with a purge of its own; running one here first would // compact the database twice for the same narrowing. - await this.sweep(store, cutoffMs, listings, signal) + await this.sweep(store, cutoffMs, listings, overdue, signal) return } if (purgeDue) { await store.purgeOlderThan(cutoffMs, signal) } - await this.cycle(store, listings, signal) + await this.cycle(store, listings, overdue, signal) } private async sweep( store: SessionSearchStore, cutoffMs: number | null, listings: SessionSearchDirectoryListings, + overdue: () => boolean, signal: AbortSignal ): Promise { const sweep = await runSessionSearchBackfill({ @@ -292,11 +311,13 @@ export class SessionSearchIndexer { previousRootsWithFiles: this.rootRecovery.previousRootsWithFiles, allowance: this.backfillAllowance, listings, + overdue, pace: this.pace, signal }) this.indexingStatus.sweepFinished(sweep.completed) if (!sweep.completed) { + this.fullSweepDue = true // An aborted sweep saw part of the machine, so it learned nothing about // root health or orphans. Publishing its empty findings would clear a // live alarm. @@ -306,7 +327,6 @@ export class SessionSearchIndexer { // transcripts until something else happened to ask for a full sweep. return } - this.fullSweepDue = false // Only what this sweep could not settle. A file it discovered and proved // present needs no watching: the recency window covers the ones that // change, and an old file deleted later is the next sweep's to find. @@ -320,6 +340,7 @@ export class SessionSearchIndexer { private async cycle( store: SessionSearchStore, listings: SessionSearchDirectoryListings, + overdue: () => boolean, signal: AbortSignal ): Promise { this.allowance.reset() @@ -351,13 +372,14 @@ export class SessionSearchIndexer { this.backfillAllowance, this.indexingStatus, this.pace, - signal + signal, + overdue ) this.indexingStatus.finishWork(this.clock.now()) } private cutoffMs(): number | null { - return sessionSearchHistoryCutoffMs(this.historyDays, this.clock.now()) + return this.retention.cutoffMs(this.clock.now()) } private publishPending(): void { @@ -374,12 +396,7 @@ export class SessionSearchIndexer { } private openStore(): void { - const store = this.registered.open({ - databasePath: this.options.databasePath, - onError: this.onError, - cutoffMs: this.cutoffMs(), - acceptingWrites: !this.paused - }) + const store = this.registered.open(this.cutoffMs(), !this.paused) this.indexingStatus.setRecoveredRows(store.recoveredWrites) } } diff --git a/src/main/ai-vault-search/session-search-indexing-status.test.ts b/src/main/ai-vault-search/session-search-indexing-status.test.ts index 1b4efa524ba..d4746f2dfb9 100644 --- a/src/main/ai-vault-search/session-search-indexing-status.test.ts +++ b/src/main/ai-vault-search/session-search-indexing-status.test.ts @@ -1,8 +1,37 @@ import { expect, it } from 'vitest' import { SessionSearchIndexingStatus } from './session-search-indexing-status' -it('walks discovering to indexing to current, and reports degraded roots once settled', () => { +/** Every phase below `idle` describes work, and work begins at `start()`. */ +function startedStatus(): SessionSearchIndexingStatus { const status = new SessionSearchIndexingStatus() + status.setStarted() + return status +} + +it('reports idle until it is started, rather than describing work nobody asked for', () => { + const status = new SessionSearchIndexingStatus() + expect(status.snapshot().phase).toBe('idle') + // Not `current` either: an index nobody built is not up to date. + status.setFilesIndexed(0) + expect(status.snapshot()).toMatchObject({ phase: 'idle', filesIndexed: 0 }) + + status.setStarted() + expect(status.snapshot().phase).toBe('indexing') +}) + +// `clear()` throws the index away, so the sweep that covered it no longer +// covers anything; without this the emptied index reports itself current. +it('stops calling itself swept once the index it swept has been cleared', () => { + const status = startedStatus() + status.sweepFinished(true) + expect(status.snapshot().phase).toBe('current') + + status.forgetSweep() + expect(status.snapshot().phase).toBe('indexing') +}) + +it('walks discovering to indexing to current, and reports degraded roots once settled', () => { + const status = startedStatus() status.beginSweep() expect(status.snapshot()).toMatchObject({ phase: 'discovering', filesTotal: null }) @@ -24,7 +53,7 @@ it('walks discovering to indexing to current, and reports degraded roots once se }) it('refuses to call itself current with work queued, or before a sweep finished', () => { - const status = new SessionSearchIndexingStatus() + const status = startedStatus() status.finishWork(1) // No sweep has ever completed, so nothing is known about the long tail. expect(status.snapshot().phase).toBe('indexing') @@ -39,7 +68,7 @@ it('refuses to call itself current with work queued, or before a sweep finished' }) it('does not let an aborted sweep count as a finished one', () => { - const status = new SessionSearchIndexingStatus() + const status = startedStatus() status.sweepFinished(false) status.finishWork(1) // The outcome is the argument, so reporting an aborted sweep cannot latch it. @@ -53,7 +82,7 @@ it('does not let an aborted sweep count as a finished one', () => { }) it('reports closed over every other phase', () => { - const status = new SessionSearchIndexingStatus() + const status = startedStatus() status.sweepFinished(true) status.finishWork(1) expect(status.snapshot().phase).toBe('current') @@ -62,7 +91,7 @@ it('reports closed over every other phase', () => { }) it('reports paused over everything, and keeps the counters it had', () => { - const status = new SessionSearchIndexingStatus() + const status = startedStatus() status.beginSweep() status.planned(3, 1) status.indexed(10) @@ -80,7 +109,7 @@ it('reports paused over everything, and keeps the counters it had', () => { }) it('starts a new sweep from zero but keeps what only an open can know', () => { - const status = new SessionSearchIndexingStatus() + const status = startedStatus() status.setRecoveredRows(7) status.beginSweep() status.planned(1, 0) diff --git a/src/main/ai-vault-search/session-search-indexing-status.ts b/src/main/ai-vault-search/session-search-indexing-status.ts index 1dd6b7538c1..2548bb09b5c 100644 --- a/src/main/ai-vault-search/session-search-indexing-status.ts +++ b/src/main/ai-vault-search/session-search-indexing-status.ts @@ -6,8 +6,13 @@ import type { SessionSearchDegradedRoot } from './session-search-degraded-roots' * `current` is the honest ceiling: the reconciler promises the newest N per * agent within one interval, not every transcript on the machine, so nothing * here ever claims the whole index is up to date. + * + * `idle` is the other end of it: an indexer nobody has started is not behind on + * anything, because it never promised to index. Reporting that as `indexing` + * described work that no timer was going to do. */ export type SessionSearchIndexPhase = + | 'idle' | 'discovering' | 'indexing' | 'current' @@ -46,6 +51,7 @@ export type SessionSearchIndexStatus = { /** Observes the backfill and the reconciler; owns no work and no timers. */ export class SessionSearchIndexingStatus { private working: 'discovering' | 'indexing' | null = null + private started = false private paused = false private closed = false private sweptClean = false @@ -91,6 +97,11 @@ export class SessionSearchIndexingStatus { if (this.closed) { return 'closed' } + if (!this.started) { + // Before `start()` nothing runs and nothing is owed; `reconcile()` is + // refused here too, so there is no work in flight to describe. + return 'idle' + } if (this.paused) { return 'paused' } @@ -129,6 +140,11 @@ export class SessionSearchIndexingStatus { this.filesIndexed = files } + /** `start()` was called; from here the phases describe work. */ + setStarted(): void { + this.started = true + } + setPaused(paused: boolean): void { this.paused = paused } diff --git a/src/main/ai-vault-search/session-search-registered-store.ts b/src/main/ai-vault-search/session-search-registered-store.ts index 2f7e48d985a..1db24a11b8e 100644 --- a/src/main/ai-vault-search/session-search-registered-store.ts +++ b/src/main/ai-vault-search/session-search-registered-store.ts @@ -1,41 +1,55 @@ import { registerSessionSearchIndexConsumer } from './session-search-index-consumer' +import { removeSessionSearchDatabase } from './session-search-schema' import { SessionSearchStore } from './session-search-store' /** * The index store an indexer currently owns, together with its registration as * a transcript consumer. * - * The two have one lifetime between them, which is why they are one object. A - * store closed while still registered is handed reads it cannot take; a + * The three have one lifetime between them, which is why they are one object. + * A store closed while still registered is handed reads it cannot take; a * registration that outlives its store is how work queued before a `close()` - * ends up indexing into a database nobody asked for. + * ends up indexing into a database nobody asked for; and a removal owed but not + * performed is a caller told their history was thrown away while it is still on + * disk. */ export class SessionSearchRegisteredStore { private store: SessionSearchStore | null = null private unregister: (() => void) | null = null + private removalOwed = false + + constructor( + private readonly databasePath: string, + private readonly onError: (error: unknown) => void + ) {} get current(): SessionSearchStore | null { return this.store } - open(args: { - databasePath: string - onError: (error: unknown) => void - cutoffMs: number | null - acceptingWrites: boolean - }): SessionSearchStore { - const store = new SessionSearchStore(args.databasePath, args.onError) - store.setRetentionCutoffMs(args.cutoffMs) - store.setAcceptingWrites(args.acceptingWrites) + open(cutoffMs: number | null, acceptingWrites: boolean): SessionSearchStore { + const store = new SessionSearchStore(this.databasePath, this.onError) + store.setRetentionCutoffMs(cutoffMs) + store.setAcceptingWrites(acceptingWrites) this.store = store this.unregister = registerSessionSearchIndexConsumer(store) return store } + /** A removal has been asked for; whoever performs it settles the debt. */ + requestRemoval(): void { + this.removalOwed = true + } + + /** Closes, then removes the database if one was still owed. */ close(): void { this.unregister?.() this.unregister = null this.store?.close() this.store = null + if (this.removalOwed) { + this.removalOwed = false + removeSessionSearchDatabase(this.databasePath) + } } } diff --git a/src/main/ai-vault-search/session-search-retention-policy.ts b/src/main/ai-vault-search/session-search-retention-policy.ts index 28454746bd4..1d12ac63de8 100644 --- a/src/main/ai-vault-search/session-search-retention-policy.ts +++ b/src/main/ai-vault-search/session-search-retention-policy.ts @@ -37,3 +37,35 @@ export function widensSessionSearchHistory(previous: number | null, next: number const to = normalizeSessionSearchHistoryDays(next) return from !== null && (to === null || to > from) } + +/** What a change to the window asks of whoever owns the index. */ +export type SessionSearchRetentionChange = 'purge' | 'resweep' | 'nothing' + +/** + * The retention window an indexer is currently keeping, and what moving it + * costs. The arithmetic and the decision live together because the two have to + * agree: a purge that uses one cutoff while the accept check holds another + * deletes rows the very next candidate re-indexes. + */ +export class SessionSearchRetentionWindow { + constructor(private historyDays: number | null) {} + + /** The oldest transcript mtime worth indexing right now, or null for all history. */ + cutoffMs(nowMs: number): number | null { + return sessionSearchHistoryCutoffMs(this.historyDays, nowMs) + } + + /** + * Moves the window and says what it asks for: narrowing purges the rows now + * outside it, widening needs a sweep because the files beyond the old bound + * were never read at all. + */ + moveTo(historyDays: number | null): SessionSearchRetentionChange { + const previous = this.historyDays + this.historyDays = historyDays + if (narrowsSessionSearchHistory(previous, historyDays)) { + return 'purge' + } + return widensSessionSearchHistory(previous, historyDays) ? 'resweep' : 'nothing' + } +} diff --git a/src/main/ai-vault-search/session-search-synthetic-sources.ts b/src/main/ai-vault-search/session-search-synthetic-sources.ts new file mode 100644 index 00000000000..86222b9c5fe --- /dev/null +++ b/src/main/ai-vault-search/session-search-synthetic-sources.ts @@ -0,0 +1,56 @@ +import type { AiVaultScanIssue } from '../../shared/ai-vault-types' +import { splitOpenCodeSqliteCandidate } from '../ai-vault/session-scanner-opencode-sqlite-paths' +import type { SessionFileCandidate } from '../ai-vault/session-scanner-types' + +/** + * A row whose path names a container and an entry inside it rather than a file + * of its own. OpenCode's SQLite sessions are the one shape today + * (`#`), which is why this reads through that source's + * own splitter rather than reinventing the encoding. + */ +export type SessionSearchSyntheticSource = { container: string; id: string } + +export function splitSyntheticSessionSource(path: string): SessionSearchSyntheticSource | null { + const openCode = splitOpenCodeSqliteCandidate(path) + return openCode ? { container: openCode.dbPath, id: openCode.sessionId } : null +} + +/** + * Which containers a pass enumerated in full, and every id each of them held. + * + * This is the synthetic equivalent of a directory listing, and it has to meet + * the same bar before the retirement walk may prove anything from it: + * + * - **Exhaustive.** Only a sweep enumerates without a per-agent limit. A cycle + * asks for the newest N, so an id it did not return may simply be the N+1th. + * Callers that are not a census do not build this at all. + * - **Successful.** A container a scan issue names could not be read, and a + * read that failed returns no ids rather than an error the walk can see. A + * named container is left out, so its rows stay unverifiable. + * - **Non-empty.** A container that returned nothing is not evidence that it + * holds nothing: a database whose schema this scanner no longer recognises + * returns an empty list with no error at all, and believing it would retire + * every session in one pass. The cost is one stale row per container whose + * last entry the user deletes, until the container gains an entry or goes. + */ +export function sessionSearchEnumeratedContainers( + candidates: readonly SessionFileCandidate[], + issues: readonly AiVaultScanIssue[] +): Map> { + const containers = new Map>() + for (const candidate of candidates) { + const synthetic = splitSyntheticSessionSource(candidate.file.path) + if (!synthetic) { + continue + } + const ids = containers.get(synthetic.container) ?? new Set() + ids.add(synthetic.id) + containers.set(synthetic.container, ids) + } + for (const issue of issues) { + if (issue.kind !== 'notice') { + containers.delete(issue.path) + } + } + return containers +}