feat(ai-vault-search): take the reader's whole-read seam and PR 2's pause rules

`requestWholeTranscriptRead` replaces the local invalidation: the reader owns
the resume point, so asking it is the honest way to say the index needs a span
the session list has already moved past.

Two callers, and the second is the one no decline can reach. A forced path is
one the store handed back from `takeStale` or a caller invalidated, and it has
to arrive as a `replace` or the consumer declines the same append forever. But
when the list's cursor already sits at a file's current stat the parse opens
nothing at all, so no consumer is asked and there is nothing to record. That is
every transcript on the machine the first time the index is switched on inside
a running app, so it gets a test on both the sweep and the cycle path.

Pausing now keeps the store's re-read set, so `filesPending` and
`droppedPending` add both bounded queues together: a caller cannot act on one
of them alone. The schema creates the index directory, so the indexer no longer
does.
This commit is contained in:
Jinwoo-H
2026-09-10 23:29:11 -04:00
parent d8c456b17d
commit 1fa323dbe1
6 changed files with 50 additions and 23 deletions
@@ -1,6 +1,5 @@
import { mkdirSync } from 'node:fs'
import { chmod, mkdir, rm, writeFile } from 'node:fs/promises'
import { dirname, join } from 'node:path'
import { join } from 'node:path'
import { afterEach, beforeEach, expect, it } from 'vitest'
import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers'
import { retireDeletedSessionSearchSources } from './session-search-deleted-sources'
@@ -19,7 +18,6 @@ let removed: string[]
beforeEach(async () => {
resetTranscriptConsumersForTests()
harness = await openSessionSearchIndexerHarness('ss-deleted-sources')
mkdirSync(dirname(harness.databasePath), { recursive: true })
removed = []
store = new SessionSearchStore(harness.databasePath)
// Only the removal matters here; the store's own removal path has its own tests.
@@ -1,6 +1,5 @@
import { mkdirSync } from 'node:fs'
import { appendFile, rm } from 'node:fs/promises'
import { dirname, join } from 'node:path'
import { join } from 'node:path'
import { afterEach, beforeEach, expect, it } from 'vitest'
import { resetSessionParseCacheForTests } from '../ai-vault/session-scanner-parse-cache'
import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers'
@@ -44,7 +43,6 @@ function transcript(sessionId: string): string {
}
function openStore(): SessionSearchStore {
mkdirSync(dirname(harness.databasePath), { recursive: true })
const opened = new SessionSearchStore(harness.databasePath, (error) => errors.push(error))
registerSessionSearchIndexConsumer(opened)
return opened
@@ -2,11 +2,11 @@ import { throwIfAiVaultScanCancelled } from '../ai-vault/ai-vault-scan-cancellat
import { parserPublishesMessages } from '../ai-vault/session-scanner-agent-parser'
import {
createSessionParseStats,
invalidateSessionParseCacheEntry,
parseAgentSessionFileCached,
sessionParseCacheCoversTranscript,
type SessionParseStats
} from '../ai-vault/session-scanner-parse-cache'
import { requestWholeTranscriptRead } from '../ai-vault/session-transcript-reader'
import type { SessionFileCandidate } from '../ai-vault/session-scanner-types'
import { fileIdentity, isSessionSearchFileCurrent } from './session-search-file-cursor'
import type { SessionSearchCycleAllowance } from './session-search-reconcile-budget'
@@ -63,13 +63,20 @@ export async function runSessionSearchIndexPass(
deferred.push(...candidates.slice(index))
break
}
// Two reasons to drop the session list's cursor, both of which end with
// the reader opening the file: the index has to re-read it whole, or the
// list is already done with it and would otherwise read nothing at all.
// The second is not an edge case: any file the list scanned before the
// index existed is in exactly that state.
// Two reasons to ask the reader for a whole read, both ending with the file
// actually being opened. `forced` covers a path the store handed back from
// `takeStale` or a caller invalidated: the reader would otherwise pick
// `append` from the session list's resume point and the consumer would
// decline it again, every cycle, forever.
//
// The second is the case no decline can reach. When the list's cursor
// already sits at this file's current stat, the parse reuses its cached fold
// and opens nothing at all, so no consumer is ever asked and there is
// nothing to record as stale. Every transcript the list scanned before the
// index existed is in that state, which is what first enablement inside a
// running app looks like.
if (forced || sessionParseCacheCoversTranscript(candidate, process.platform)) {
invalidateSessionParseCacheEntry(candidate.file.path)
requestWholeTranscriptRead(candidate.file.path)
}
try {
await parseAgentSessionFileCached(candidate, process.platform, stats)
@@ -212,7 +212,7 @@ it('restarts the backfill when the history window widens and purges when it narr
expect(sessionsMatching('recent')).toEqual([SESSION_ID])
})
it('refuses writes while paused and bounds the queue it keeps', async () => {
it('refuses writes while paused, remembers what it declined, and bounds both queues', async () => {
const path = transcriptPath()
await writeClaudeTranscript(path, ['indexed before the pause'], SESSION_ID)
await newIndexer({ pendingLimit: 3 }).start()
@@ -225,14 +225,19 @@ it('refuses writes while paused and bounds the queue it keeps', async () => {
await parseTranscript(path)
await nextCycle()
expect(sessionsMatching('paused')).toEqual([])
// A pause is the window in which reads are declined, so forgetting them would
// lose exactly the files the pause covered.
expect(indexer?.status().filesPending).toBe(1)
// A caller can keep invalidating right through the pause; that queue is capped.
indexer?.invalidate(['/a', '/b', '/c', '/d', '/e'])
const paused = indexer?.status()
expect(paused?.filesPending).toBe(3)
expect(paused?.filesPending).toBe(4)
expect(paused?.droppedPending).toBe(2)
await indexer?.resume()
expect(sessionsMatching('paused')).toEqual([SESSION_ID])
expect(indexer?.status()).toMatchObject({ filesPending: 0, droppedPending: 2 })
expect(indexer?.status().phase).not.toBe('paused')
})
@@ -279,14 +284,29 @@ it('spends a cycle budget and rolls the rest into the next cycle', async () => {
expect(sessionsMatching('budgeted')).toHaveLength(4)
})
it('indexes a transcript the session list had already parsed before it existed', async () => {
// First enablement inside a running app is the normal case, not an edge: the
// 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 () => {
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])
})
it('fills an empty index over a warm session-list cache on the first sweep', async () => {
const path = transcriptPath()
await writeClaudeTranscript(path, ['scanned before the index existed'], SESSION_ID)
// The list's cursor now sits at this file's current stat, so an ordinary
// parse reads nothing and the index would stay empty forever.
await parseTranscript(path)
await newIndexer().start()
expect(sessionsMatching('scanned')).toEqual([SESSION_ID])
})
@@ -1,5 +1,3 @@
import { mkdirSync } from 'node:fs'
import { dirname } from 'node:path'
import { runSessionSearchBackfill } from './session-search-backfill'
import { pauseBackfill } from './session-search-backfill-pacing'
import {
@@ -329,14 +327,16 @@ export class SessionSearchIndexer {
}
private publishPending(): void {
// Both queues, because a caller cannot act on one of them: the store's
// re-read set and this indexer's roll-over queue are each bounded, and
// either overrunning means the same thing for coverage.
this.indexingStatus.setPending(
this.pending.size + (this.store?.pendingFileCount ?? 0),
this.pending.droppedCount
this.pending.droppedCount + (this.store?.droppedPendingFileCount ?? 0)
)
}
private openStore(): void {
mkdirSync(dirname(this.options.databasePath), { recursive: true })
const store = new SessionSearchStore(this.options.databasePath, this.onError)
store.setRetentionCutoffMs(this.cutoffMs())
store.setAcceptingWrites(!this.paused)
@@ -17,7 +17,11 @@ export type SessionSearchIndexStatus = {
bytesIndexed: number
/** Queued re-reads: the store's stale set plus whatever the budget rolled over. */
filesPending: number
/** Pending entries dropped at the bound, so a long pause cannot grow memory. */
/**
* Re-reads dropped at a bound, the indexer's queue and the store's re-read set
* together. Non-zero means the queue is knowingly incomplete, so coverage
* cannot be reported as whole until the next full sweep.
*/
droppedPending: number
/** Unfinished writes the open tombstoned; a non-zero value means a crash. */
recoveredRows: number