diff --git a/src/main/ai-vault-search/session-search-backfill.ts b/src/main/ai-vault-search/session-search-backfill.ts index 21f5fb8f291..c11d487093b 100644 --- a/src/main/ai-vault-search/session-search-backfill.ts +++ b/src/main/ai-vault-search/session-search-backfill.ts @@ -12,7 +12,6 @@ import type { SessionSearchIndexingStatus } from './session-search-indexing-stat import { degradedSessionSearchRoots, rootFileCounts, - underDegradedRoot, type SessionSearchDegradedRoot } from './session-search-root-health' import { @@ -119,12 +118,9 @@ export async function runSessionSearchBackfill( /** * A sweep is the only pass that sees every root, so it is the only one that can - * retire a source deleted while nothing was running. It is also the pass that - * would delete a user's entire searchable history the first time an SSH mount - * or an external drive is not there, because every path under it answers ENOENT - * at once. A degraded root's files are therefore never retired, however loudly - * the filesystem says they are gone - * (docs/reference/ssh-execution-boundary.md). + * retire a source deleted while nothing was running. The degraded-root fence + * that keeps an unmounted volume from taking its history with it belongs to the + * retirement function itself, which the cycle calls too. */ async function retireSweptAwaySources( args: SessionSearchBackfillArgs, @@ -134,9 +130,10 @@ async function retireSweptAwaySources( const undiscovered = args.store .indexedSources() .map((source) => source.path) - .filter((path) => !discoveredPaths.has(path) && !underDegradedRoot(path, degradedRoots)) + .filter((path) => !discoveredPaths.has(path)) return retireDeletedSessionSearchSources(args.store, undiscovered, { signal: args.signal, - limit: RETIREMENT_CHECKS_PER_SWEEP + limit: RETIREMENT_CHECKS_PER_SWEEP, + degradedRoots }) } 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 fd6d227c5fa..cb9f3a87640 100644 --- a/src/main/ai-vault-search/session-search-deleted-sources.ts +++ b/src/main/ai-vault-search/session-search-deleted-sources.ts @@ -1,5 +1,6 @@ import { wslGatedStat } from '../native-chat/wsl-transcript-fs-access' import { splitOpenCodeSqliteCandidate } from '../ai-vault/session-scanner-opencode-sqlite-paths' +import { underDegradedRoot, type SessionSearchDegradedRoot } from './session-search-root-health' import type { SessionSearchStore } from './session-search-store' export type SessionSearchRetirement = { @@ -19,15 +20,27 @@ const PROVEN_GONE = new Set(['ENOENT', 'ENOTDIR']) /** * Retires index rows for sources that are provably gone. Callers pass the paths - * the index holds but this sweep did not discover; everything else is either - * still there or was never in the discovery window in the first place. + * the index holds but did not discover; everything else is either still there + * or was never in the discovery window in the first place. + * + * The degraded-root fence lives here, not at the call sites, because both the + * sweep and the cycle retire and either one alone would delete a user's history + * the first time an SSH mount or an external drive is not there: every path + * under it answers ENOENT at once. A file under a root this pass could not + * trust is `unverifiable`, whatever the filesystem says + * (docs/reference/ssh-execution-boundary.md). */ export async function retireDeletedSessionSearchSources( store: SessionSearchStore, paths: readonly string[], - options: { signal?: AbortSignal; limit?: number } = {} + options: { + signal?: AbortSignal + limit?: number + degradedRoots?: readonly SessionSearchDegradedRoot[] + } = {} ): Promise { const { signal } = options + const degradedRoots = options.degradedRoots ?? [] const retirement: SessionSearchRetirement = { retired: [], unverifiable: [], unchecked: [] } const limit = options.limit ?? Number.POSITIVE_INFINITY for (const [index, path] of paths.entries()) { @@ -38,6 +51,10 @@ export async function retireDeletedSessionSearchSources( retirement.unchecked.push(...paths.slice(index)) break } + if (underDegradedRoot(path, degradedRoots)) { + retirement.unverifiable.push(path) + continue + } // A synthetic OpenCode row names the database it came from, never a file of // its own; statting the candidate path would report every one of them gone. const statPath = splitOpenCodeSqliteCandidate(path)?.dbPath ?? path 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 5a505ad13f6..2692813b0a7 100644 --- a/src/main/ai-vault-search/session-search-index-pass.ts +++ b/src/main/ai-vault-search/session-search-index-pass.ts @@ -6,11 +6,7 @@ import { type SessionParseStats } from '../ai-vault/session-scanner-parse-cache' import type { SessionFileCandidate } from '../ai-vault/session-scanner-types' -import { - fileIdentity, - isSessionSearchFileCurrent, - type SessionSearchIndexedFile -} from './session-search-file-cursor' +import { fileIdentity, isSessionSearchFileCurrent } from './session-search-file-cursor' import type { SessionSearchCycleAllowance } from './session-search-reconcile-budget' import type { SessionSearchStore } from './session-search-store' @@ -76,7 +72,6 @@ export async function runSessionSearchIndexPass( // behind, and a list cursor already at this file's current stat would make // the parse open nothing at all — the state every transcript is in the // first time the index is switched on inside a running app. - const before = indexedRecord(store, candidate) try { await parseAgentSessionFileCached( candidate, @@ -84,10 +79,12 @@ export async function runSessionSearchIndexPass( stats, forced ? 'whole' : 'any' ) - // A parse that returned without throwing is not a parse the index kept: - // the consumer declines a read it cannot use, and counting that as - // indexed is how a status ends up claiming files it does not hold. - if (published(before, indexedRecord(store, candidate))) { + // Took it, not moved: the test is whether the index now covers this file + // at this stat, which is the same question the skip at the top asks. A + // cursor comparison looks equivalent and is not — a forced re-read of an + // unchanged file writes an identical cursor, so it would never settle and + // the path would be re-read whole every interval for good. + if (indexIsCurrent(store, candidate)) { options.onIndexed?.(candidate, bytes) } } catch (error) { @@ -145,29 +142,6 @@ function mustReadWhole( return typeof size === 'number' && stored.byteOffset > size } -function indexedRecord( - store: SessionSearchStore, - candidate: SessionFileCandidate -): SessionSearchIndexedFile | null { - return store.indexedFile(candidate.file.path, fileIdentity(candidate.file)) -} - -/** True when the index's own cursor for this file moved, which only a publish does. */ -function published( - before: SessionSearchIndexedFile | null, - after: SessionSearchIndexedFile | null -): boolean { - if (!after) { - return false - } - return ( - !before || - before.byteOffset !== after.byteOffset || - before.mtimeMs !== after.mtimeMs || - before.sizeBytes !== after.sizeBytes - ) -} - /** What this read will actually cost: the tail past the index's own cursor. */ function unreadBytes(store: SessionSearchStore, candidate: SessionFileCandidate): number { const size = candidate.file.sizeBytes ?? 0 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 2e0e4627919..e48ffa73b1c 100644 --- a/src/main/ai-vault-search/session-search-indexer.test.ts +++ b/src/main/ai-vault-search/session-search-indexer.test.ts @@ -489,6 +489,62 @@ it('reads a stale file at its current stat, not the one it was recorded with', a expect(indexedCursor(path)).toEqual({ mtime_ms: current.mtimeMs, size_bytes: current.size }) }) +// Round 2, item 1: the sweep kept the rows and a cycle twenty seconds later +// deleted them, because the degraded-root fence was on the sweep path only. +it.skipIf(!CAN_DENY_READ)( + 'keeps an unmounted root through the cycles that follow the sweep', + async () => { + await writeClaudeTranscript(transcriptPath(), ['a session on a removable volume'], SESSION_ID) + await newIndexer().start() + expect(sessionsMatching('removable')).toEqual([SESSION_ID]) + + await rm(harness.claudeProjectDir, { recursive: true, force: true }) + await indexer?.reconcile({ full: true }) + expect(sessionsMatching('removable')).toEqual([SESSION_ID]) + + await nextCycle() + expect(sessionsMatching('removable')).toEqual([SESSION_ID]) + expect(indexer?.status().phase).toBe('degraded') + } +) + +// Round 2, item 3: the alarm was single-shot. The degraded sweep's zero became +// the baseline, so the second sweep compared zero with zero and retired. +it.skipIf(!CAN_DENY_READ)('keeps an unmounted root across repeated sweeps', async () => { + await writeClaudeTranscript(transcriptPath(), ['a session on a removable volume'], SESSION_ID) + await newIndexer().start() + + await rm(harness.claudeProjectDir, { recursive: true, force: true }) + await indexer?.reconcile({ full: true }) + await indexer?.reconcile({ full: true }) + expect(sessionsMatching('removable')).toEqual([SESSION_ID]) + expect(indexer?.status().phase).toBe('degraded') + + // Remounted: the root lists transcripts again and the alarm clears. + await writeClaudeTranscript(transcriptPath(), ['a session on a removable volume'], SESSION_ID) + await indexer?.reconcile({ full: true }) + expect(indexer?.status().degradedRoots).toEqual([]) + expect(sessionsMatching('removable')).toEqual([SESSION_ID]) +}) + +// Round 2, item 2: a forced whole re-read of an unchanged file writes an +// identical cursor. Judging by cursor movement, that never settles: the path +// is owed forever and re-read whole on every interval. +it('settles an invalidated file that turned out not to have changed', async () => { + const path = transcriptPath() + await writeClaudeTranscript(path, ['unchanged after all'], SESSION_ID) + await newIndexer().start() + + indexer?.invalidate([path]) + await indexer?.reconcile() + expect(indexer?.status()).toMatchObject({ filesPending: 0, phase: 'current' }) + + // And it stays settled: the next cycle has no reason to open it again. + const bytes = indexer?.status().bytesIndexed + await nextCycle() + expect(indexer?.status()).toMatchObject({ filesPending: 0, bytesIndexed: bytes }) +}) + // Finding 8: a path sits in both queues the moment a read is declined during a // pause and a caller then invalidates the same file. Summing them reports one // transcript as two, and a caller has no way to tell that from two files. diff --git a/src/main/ai-vault-search/session-search-indexer.ts b/src/main/ai-vault-search/session-search-indexer.ts index 4cdab5aa9f6..72814f67d41 100644 --- a/src/main/ai-vault-search/session-search-indexer.ts +++ b/src/main/ai-vault-search/session-search-indexer.ts @@ -21,6 +21,7 @@ import { sessionSearchHistoryCutoffMs, widensSessionSearchHistory } from './session-search-retention-policy' +import { withLastHealthyRootCounts } from './session-search-root-health' import { removeSessionSearchDatabase } from './session-search-schema' import type { SessionSearchScanRoots } from './session-search-scan-roots' import { SessionSearchStore } from './session-search-store' @@ -161,6 +162,7 @@ export class SessionSearchIndexer { removeSessionSearchDatabase(this.options.databasePath) this.pending.clear() this.previousRecent = new Set() + this.rootFileCounts = new Map() this.openStore() this.fullSweepDue = true }) @@ -273,8 +275,11 @@ export class SessionSearchIndexer { pace: this.pace, signal }) - this.rootFileCounts = sweep.rootFileCounts + // Last healthy count, not last count: a degraded sweep's zero would + // otherwise become the baseline and the next sweep would retire the tree. + this.rootFileCounts = withLastHealthyRootCounts(this.rootFileCounts, sweep.rootFileCounts) this.indexingStatus.setDegradedRoots(sweep.degradedRoots) + this.indexingStatus.sweepFinished(sweep.completed) if (!sweep.completed) { // A sweep is due until it finishes. Clearing the flag on entry meant a // pause part way through abandoned the rest of the machine's transcripts @@ -286,7 +291,6 @@ export class SessionSearchIndexer { // discovery and the first cycle is invisible to both otherwise. The // cycle's retirement cap keeps that one-off check off the critical path. this.previousRecent = sweep.watchPaths - this.indexingStatus.sweepCompleted() this.indexingStatus.finishWork(this.clock.now()) } @@ -301,6 +305,9 @@ export class SessionSearchIndexer { pending: this.pending.drain(), previousRecent: this.previousRecent, retirementChecksPerCycle: this.options.retirementChecksPerCycle, + // Read but not written: a recent-window discovery is not a census, so it + // can spot a root that went to zero without redefining what healthy was. + previousRootFileCounts: this.rootFileCounts, signal }) // Work that was drained and then not read is a hole in the index, not 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 4c004338a65..1b4efa524ba 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 @@ -15,7 +15,7 @@ it('walks discovering to indexing to current, and reports degraded roots once se // Work in flight is the more useful thing to show; the roots are in the // snapshot either way. expect(status.snapshot().phase).toBe('indexing') - status.sweepCompleted() + status.sweepFinished(true) status.finishWork(123) expect(status.snapshot()).toMatchObject({ phase: 'degraded', lastReconcileAt: 123 }) @@ -29,7 +29,7 @@ it('refuses to call itself current with work queued, or before a sweep finished' // No sweep has ever completed, so nothing is known about the long tail. expect(status.snapshot().phase).toBe('indexing') - status.sweepCompleted() + status.sweepFinished(true) expect(status.snapshot().phase).toBe('current') status.setPending(3, 0) @@ -38,9 +38,23 @@ it('refuses to call itself current with work queued, or before a sweep finished' expect(status.snapshot().phase).toBe('current') }) +it('does not let an aborted sweep count as a finished one', () => { + const status = new SessionSearchIndexingStatus() + status.sweepFinished(false) + status.finishWork(1) + // The outcome is the argument, so reporting an aborted sweep cannot latch it. + expect(status.snapshot().phase).toBe('indexing') + + status.sweepFinished(true) + expect(status.snapshot().phase).toBe('current') + status.sweepFinished(false) + // A later abort does not un-know that a whole sweep once finished. + expect(status.snapshot().phase).toBe('current') +}) + it('reports closed over every other phase', () => { const status = new SessionSearchIndexingStatus() - status.sweepCompleted() + status.sweepFinished(true) status.finishWork(1) expect(status.snapshot().phase).toBe('current') status.setClosed() 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 4e0c1221902..ff93881a860 100644 --- a/src/main/ai-vault-search/session-search-indexing-status.ts +++ b/src/main/ai-vault-search/session-search-indexing-status.ts @@ -98,9 +98,13 @@ export class SessionSearchIndexingStatus { this.working = null } - /** One whole sweep ran to completion; until then nothing is `current`. */ - sweepCompleted(): void { - this.sweptClean = true + /** + * One whole sweep finished, or did not. The outcome is the argument rather + * than the call site's position, so an aborted sweep cannot latch this by + * being reported a line too early. + */ + sweepFinished(completed: boolean): void { + this.sweptClean ||= completed } /** Files the index holds, counted in the store rather than tallied per attempt. */ diff --git a/src/main/ai-vault-search/session-search-reconciler.test.ts b/src/main/ai-vault-search/session-search-reconciler.test.ts new file mode 100644 index 00000000000..22b6b4a2c39 --- /dev/null +++ b/src/main/ai-vault-search/session-search-reconciler.test.ts @@ -0,0 +1,85 @@ +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 type { SessionFileCandidate } from '../ai-vault/session-scanner-types' +import { registerSessionSearchIndexConsumer } from './session-search-index-consumer' +import { + openSessionSearchIndexerHarness, + type SessionSearchIndexerHarness +} from './session-search-indexer-test-fixture' +import { SessionSearchIndexingStatus } from './session-search-indexing-status' +import { SessionSearchCycleAllowance } from './session-search-reconcile-budget' +import { runSessionSearchReconcileCycle } from './session-search-reconciler' +import { SessionSearchStore } from './session-search-store' + +// The store's re-read set holds ten times what the indexer's queue does, so an +// abort that hands the store's overflow to the smaller queue silently discards +// the difference. Whatever came from the store goes back to the store. + +const STALE_FILES = 3_000 +const READS_BEFORE_ABORT = 10 + +let harness: SessionSearchIndexerHarness +let store: SessionSearchStore + +beforeEach(async () => { + resetSessionParseCacheForTests() + resetTranscriptConsumersForTests() + harness = await openSessionSearchIndexerHarness('ss-reconciler') + store = new SessionSearchStore(harness.databasePath) + registerSessionSearchIndexConsumer(store) +}) + +afterEach(async () => { + resetTranscriptConsumersForTests() + store.close() + await harness.cleanup() +}) + +function staleCandidate(index: number): SessionFileCandidate { + const at = 1_740_000_000_000 + index + return { + agent: 'claude', + codexHome: null, + file: { + path: `${harness.claudeProjectDir}/stale-${index}.jsonl`, + mtimeMs: at, + modifiedAt: new Date(at).toISOString(), + sizeBytes: 128 + } + } +} + +it('hands the store back its own overflow when a cycle is aborted', async () => { + for (let index = 0; index < STALE_FILES; index++) { + store.markStale(staleCandidate(index)) + } + expect(store.pendingFileCount).toBe(STALE_FILES) + + const controller = new AbortController() + const status = new SessionSearchIndexingStatus() + let reads = 0 + // None of these paths exist, so every read fails; abort once a few have run. + status.failed = () => { + if (++reads >= READS_BEFORE_ABORT) { + controller.abort() + } + } + + const cycle = await runSessionSearchReconcileCycle({ + store, + roots: harness.roots, + status, + recentPerAgent: 12, + allowance: new SessionSearchCycleAllowance({ files: 10_000, bytes: 1_000_000_000 }), + pending: [], + previousRecent: new Set(), + signal: controller.signal + }) + + expect(cycle.completed).toBe(false) + // Back where it came from, at its own bound, not truncated into a queue a + // tenth the size. + expect(store.pendingFileCount).toBe(STALE_FILES) + expect(cycle.deferred).toEqual([]) +}) diff --git a/src/main/ai-vault-search/session-search-reconciler.ts b/src/main/ai-vault-search/session-search-reconciler.ts index 1afe7863c68..fd9bf691f5c 100644 --- a/src/main/ai-vault-search/session-search-reconciler.ts +++ b/src/main/ai-vault-search/session-search-reconciler.ts @@ -44,6 +44,8 @@ export type SessionSearchReconcileArgs = { previousRecent: ReadonlySet /** Stats a cycle spends proving deletions; the rest stay watched. */ retirementChecksPerCycle?: number + /** What each root listed when it was last healthy, so a tree that went empty is visible. */ + previousRootFileCounts?: ReadonlyMap signal?: AbortSignal } @@ -85,12 +87,16 @@ export async function runSessionSearchReconcileCycle( // Everything drained out of a queue, so an abort can put back exactly what // it did not get to. A drained entry that is never read is a hole in the // index, and dropping it is how a file stays missing until the next sweep. - const owed = new Map() + // Which queue it came from is tracked, because they are not the same size: + // the store holds ten times what this one does, so returning its overflow + // here would quietly discard the difference. + const owed = new Map() for (const entry of args.pending) { - owed.set(entry.path, entry) + owed.set(entry.path, { entry, fromStore: false }) } for (const candidate of stale) { - owed.set(candidate.file.path, { path: candidate.file.path, candidate, forced: true }) + const entry = { path: candidate.file.path, candidate, forced: true } + owed.set(candidate.file.path, { entry, fromStore: true }) } let completed = true @@ -120,18 +126,22 @@ export async function runSessionSearchReconcileCycle( completed = false } for (const candidate of pass.deferred) { - owed.set(candidate.file.path, { - path: candidate.file.path, - candidate, - forced: forced.has(candidate.file.path) - }) + const path = candidate.file.path + const entry = { path, candidate, forced: forced.has(path) } + owed.set(path, { entry, fromStore: owed.get(path)?.fromStore ?? false }) } for (const path of queued.unresolved) { // Never silently dropped: a caller asked for this path and the index has // no record of it, so it stays queued and keeps being counted. - owed.set(path, { path, candidate: null, forced: true }) + owed.set(path, { entry: { path, candidate: null, forced: true }, fromStore: false }) } + // Before the retirement, not after it: the cycle deletes rows too, so it + // needs the same fence the sweep has or one interval undoes the sweep's care. + const degradedRoots = await degradedSessionSearchRoots(swept.discoveries, issues, { + signal, + previousFileCounts: args.previousRootFileCounts + }) const retirement = completed ? await retireDeletedSessionSearchSources( store, @@ -140,13 +150,24 @@ export async function runSessionSearchReconcileCycle( ), { signal, - limit: args.retirementChecksPerCycle ?? DEFAULT_SESSION_SEARCH_RETIREMENT_CHECKS + limit: args.retirementChecksPerCycle ?? DEFAULT_SESSION_SEARCH_RETIREMENT_CHECKS, + degradedRoots } ) : { retired: [], unverifiable: [], unchecked: [] } for (const path of retirement.retired) { owed.delete(path) } + // A record that came from the store goes back to the store, which applies + // its own retention rule and its own, larger bound. + const deferred: SessionSearchPendingFile[] = [] + for (const { entry, fromStore } of owed.values()) { + if (fromStore && entry.candidate) { + store.markStale(entry.candidate) + continue + } + deferred.push(entry) + } for (const refusal of cursorChatMetaRefusals()) { recordSessionScanIssue(issues, { @@ -160,8 +181,8 @@ export async function runSessionSearchReconcileCycle( // reach was never checked at all: both stay watched instead of being // retired or forgotten. recentPaths: new Set([...recentPaths, ...retirement.unverifiable, ...retirement.unchecked]), - deferred: [...owed.values()], - degradedRoots: await degradedSessionSearchRoots(swept.discoveries, issues, { signal }), + deferred, + degradedRoots, completed } }) diff --git a/src/main/ai-vault-search/session-search-root-health.ts b/src/main/ai-vault-search/session-search-root-health.ts index 9ccceae9fb2..cc0b7c8d9fa 100644 --- a/src/main/ai-vault-search/session-search-root-health.ts +++ b/src/main/ai-vault-search/session-search-root-health.ts @@ -27,6 +27,25 @@ export function rootFileCounts(discoveries: readonly SessionFileDiscovery[]): Ma return counts } +/** + * Carries a root's last healthy count forward across a sweep that listed it + * empty. Without this the alarm is single-shot: the degraded sweep's zero + * becomes the baseline, the next sweep compares zero against zero, and the + * unmounted tree is retired on the second pass instead of the first. + */ +export function withLastHealthyRootCounts( + previous: ReadonlyMap, + observed: ReadonlyMap +): Map { + const merged = new Map(previous) + for (const [root, count] of observed) { + if (count > 0) { + merged.set(root, count) + } + } + return merged +} + /** * Classifies the roots a sweep just walked. Only roots that yielded nothing are * probed: a root that returned files is readable by construction, which keeps