fix(ai-vault-search): apply the round-1 guards on the cycle path too

Two of the round-1 fixes were written on the sweep and the cycle walked around
them twenty seconds later.

The degraded-root fence now lives inside the retirement function itself rather
than at one call site, so both passes get it from one place and the cycle
cannot delete what the sweep just protected. The cycle also reads the sweep's
root counts, so it can see a tree that went to zero at all; it does not write
them back, because a recent-window discovery is not a census.

The N-to-zero alarm was single-shot: the degraded sweep's own zero became the
baseline, so the next sweep compared zero with zero and retired the tree it had
just spared. A root keeps its last healthy count until one lists it non-empty
again.

Taking a file no longer means the cursor moved. A forced whole re-read of an
unchanged file writes an identical cursor, so an invalidated file that turned
out not to have changed was never settled: owed forever, re-read whole every
interval, status pinned. It means the index now covers the file at this stat,
which is the same question the skip at the top of the pass asks, and it is the
same function.

An aborted cycle handed the store's re-read set to a queue a tenth its size,
which silently discarded the difference. What came from the store goes back to
the store, under its own bound and its own retention rule.

Whether a sweep finished is now an argument rather than a call site's position,
so an aborted one cannot latch `current` by being reported a line too early.
This commit is contained in:
Jinwoo-H
2026-09-10 23:29:12 -04:00
parent 764adc15a6
commit d53c7ffa58
10 changed files with 259 additions and 65 deletions
@@ -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
})
}
@@ -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<SessionSearchRetirement> {
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
@@ -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
@@ -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.
@@ -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
@@ -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()
@@ -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. */
@@ -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([])
})
@@ -44,6 +44,8 @@ export type SessionSearchReconcileArgs = {
previousRecent: ReadonlySet<string>
/** 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<string, number>
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<string, SessionSearchPendingFile>()
// 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<string, { entry: SessionSearchPendingFile; fromStore: boolean }>()
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
}
})
@@ -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<string, number>,
observed: ReadonlyMap<string, number>
): Map<string, number> {
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