fix(ai-vault-search): give a behind consumer a way to get the read it needs

`markStale` promised a whole re-read that nothing could deliver. The reader
picks append or replace from the session list's resume point, so with an empty
index and a warm parse cache every read arrives as `append`, the consumer
declines every one, and nothing is ever indexed. That is the state on first
enablement inside a running app.

`requestWholeTranscriptRead` in the reader drops that path's resume point, which
is the one lever that changes the next read's mode, and it lives with the cache
that owns it rather than with the consumer that wants it. `takeStale` documents
that its paths need it before a scan is re-dispatched.

Pausing no longer empties the re-read set. `acceptsCandidate` answers whether a
write may start now and was being used for both jobs, so reconfiguring retention
while paused dropped every entry and a declined read during a pause was never
recorded at all. Retention alone prunes; a paused decline is recorded, bounded,
with the oldest dropped and the drops counted, because a set that silently
forgets is worse than one that says it is incomplete.

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