Batch search retention cleanup and preserve runtime ownership across queries

This commit is contained in:
Jinwoo-H
2026-09-06 18:34:14 -04:00
parent 7385c0ce50
commit 5d1d41a8e4
21 changed files with 816 additions and 110 deletions
+5
View File
@@ -880,6 +880,11 @@ jobs:
src/relay/windows-port-scan.win32.test.ts
src/main/ai-vault-search/session-search-fts5-contract.test.ts
src/main/ai-vault-search/session-search-schema.test.ts
src/main/ai-vault-search/session-search-query-regressions.test.ts
src/main/ai-vault-search/session-search-retention-delete.test.ts
src/main/ai-vault-search/session-search-refresh-lane.test.ts
src/main/ai-vault/session-newest-files.test.ts
src/main/ai-vault/session-scanner-discovery-cancellation.test.ts
# Why the :parallel variant: identical to build:release except the three
# electron-vite targets overlap instead of running back to back. The Linux package
@@ -0,0 +1,88 @@
import assert from 'node:assert/strict'
import { mkdtemp, rm, stat } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { setImmediate as yieldToEventLoop } from 'node:timers/promises'
import SyncDatabase from '../../src/main/sqlite/sync-database'
import { SessionSearchStore } from '../../src/main/ai-vault-search/session-search-store'
import { deleteExpiredSearchFiles } from '../../src/main/ai-vault-search/session-search-retention-delete'
import { SessionSearchIndexWriter } from '../../src/main/ai-vault-search/session-search-index-writer'
// Bundle with esbuild --bundle --platform=node, then run on the host under test.
const root = await mkdtemp(join(tmpdir(), 'orca-search-retention-bench-'))
try {
for (const mode of ['whole-file', 'batched', 'batched-pinned-reader']) {
const path = join(root, `${mode}.sqlite`)
const store = new SessionSearchStore(path)
let reader: SyncDatabase | null = null
try {
const db = store.db
db.exec(`INSERT INTO sessions(id,agent,session_id,file_path,title,cwd,cwd_key,resume_command)
VALUES (1,'claude','1','fixture','synthetic benchmark','/fixture','/fixture','');
INSERT INTO files(path,byte_offset,mtime_ms,session_row_id) VALUES ('fixture',1,1,1);
BEGIN;
WITH RECURSIVE n(i) AS (VALUES(1) UNION ALL SELECT i+1 FROM n WHERE i<60000)
INSERT INTO messages(id,session_row_id,role) SELECT i,1,'user' FROM n;
INSERT INTO messages_fts(rowid,user_text) SELECT id,'synthetic benchmark needle ' || id ||
' repeated context for a representative coding conversation with commands and paths src/example.ts'
FROM messages;
INSERT INTO conversation_fts(rowid,user_text) SELECT rowid,user_text FROM messages_fts;
COMMIT; PRAGMA wal_checkpoint(TRUNCATE)`)
assert.equal(store.search({ query: 'needle' }).hits.length, 1)
if (mode === 'batched-pinned-reader') {
reader = new SyncDatabase(path, { readonly: true })
reader.exec('BEGIN')
reader.prepare('SELECT count(*) FROM messages').get()
}
const intervals: number[] = []
let previous = performance.now()
const started = previous
if (mode === 'whole-file') {
new SessionSearchIndexWriter(db).removeFile('fixture')
intervals.push(performance.now() - previous)
} else {
await deleteExpiredSearchFiles(
db,
100,
() => false,
() => {},
async () => {
intervals.push(performance.now() - previous)
assert.equal(store.search({ query: 'needle' }).hits.length, 0)
await yieldToEventLoop()
previous = performance.now()
}
)
}
const wallMs = performance.now() - started
assert.equal(
(db.prepare('SELECT count(*) AS n FROM messages_fts').get() as { n: number }).n,
0
)
assert.equal(
(db.prepare('SELECT count(*) AS n FROM conversation_fts').get() as { n: number }).n,
0
)
const walBytes = (await stat(`${path}-wal`)).size
intervals.sort((a, b) => a - b)
console.log(
JSON.stringify({
mode,
platform: process.platform,
node: process.version,
rows: 60000,
wallMs,
steps: intervals.length,
maxStepMs: intervals.at(-1),
p95StepMs: intervals[Math.floor(intervals.length * 0.95)],
walBytes
})
)
} finally {
reader?.close()
store.close()
}
}
} finally {
await rm(root, { recursive: true, force: true })
}
@@ -125,15 +125,6 @@ export class SessionSearchIndexWriter {
)
}
/** Paths of indexed files last modified before the cutoff, oldest first. */
filesOlderThan(cutoffMs: number): string[] {
return (
this.db
.prepare('SELECT path FROM files WHERE mtime_ms < ? ORDER BY mtime_ms')
.all(cutoffMs) as { path: string }[]
).map((row) => row.path)
}
removeFile(path: string): void {
const existing = this.db
.prepare('SELECT session_row_id FROM files WHERE path = ?')
@@ -0,0 +1,72 @@
import { expect, it, vi } from 'vitest'
import { SessionSearchRefreshLane } from './session-search-refresh-lane'
import { waitForPromiseWithSignal } from '../../shared/abort-signal-reason'
function barrier() {
let resolve!: () => void
return {
promise: new Promise<void>((r) => {
resolve = r
}),
release: () => resolve()
}
}
it('shares concurrent tiers, preserves a remaining reader, and reads fresh after completion', async () => {
const lane = new SessionSearchRefreshLane()
const gate = barrier()
const signals: AbortSignal[] = []
const refresh = vi.fn((signal: AbortSignal) => {
signals.push(signal)
return gate.promise
})
const first = new AbortController()
const a = lane.run({ claudeProjectsDir: '/isolated/a' }, refresh, first.signal)
const b = lane.run({ claudeProjectsDir: '/isolated/a' }, refresh)
first.abort()
await expect(a).rejects.toMatchObject({ name: 'AbortError' })
expect(refresh).toHaveBeenCalledTimes(1)
expect(signals[0].aborted).toBe(false)
gate.release()
await b
await lane.run({ claudeProjectsDir: '/isolated/a' }, refresh)
expect(refresh).toHaveBeenCalledTimes(2)
})
it('isolates roots and cancels abandoned work without poisoning a retry', async () => {
const lane = new SessionSearchRefreshLane()
const started = barrier()
const signals: AbortSignal[] = []
const refresh = vi.fn((signal: AbortSignal) => {
signals.push(signal)
if (signals.length === 2) {
started.release()
}
return waitForPromiseWithSignal(new Promise<void>(() => {}), signal)
})
const controller = new AbortController()
const a = lane.run({ codexSessionsDir: '/a' }, refresh, controller.signal).catch((e) => e)
const b = lane.run({ codexSessionsDir: '/b' }, refresh).catch((e) => e)
await started.promise
controller.abort()
expect(await a).toMatchObject({ name: 'AbortError' })
expect(signals[0].aborted).toBe(true)
expect(signals[1].aborted).toBe(false)
lane.cancel()
expect(await b).toMatchObject({ name: 'AbortError' })
await lane.run({ codexSessionsDir: '/a' }, async () => {})
})
it('does not start already-cancelled work and retries failures', async () => {
const lane = new SessionSearchRefreshLane()
const refresh = vi.fn(async () => {
throw new Error('scan failure')
})
await expect(lane.run({}, refresh, AbortSignal.abort())).rejects.toMatchObject({
name: 'AbortError'
})
expect(refresh).not.toHaveBeenCalled()
await expect(lane.run({}, refresh)).rejects.toThrow('scan failure')
await expect(lane.run({}, refresh)).rejects.toThrow('scan failure')
expect(refresh).toHaveBeenCalledTimes(2)
})
@@ -0,0 +1,54 @@
import { waitForPromiseWithSignal, throwIfSignalAborted } from '../../shared/abort-signal-reason'
import type { SessionSearchScanRoots } from './session-search-service'
type Refresh = { controller: AbortController; promise: Promise<void>; users: number }
/** Share concurrent query refreshes, never completed filesystem snapshots. */
export class SessionSearchRefreshLane {
private readonly runs = new Map<string, Refresh>()
async run(
roots: SessionSearchScanRoots,
refresh: (signal: AbortSignal) => Promise<void>,
signal?: AbortSignal
): Promise<void> {
throwIfSignalAborted(signal)
const key = JSON.stringify(Object.entries(roots).sort(([a], [b]) => a.localeCompare(b)))
let run = this.runs.get(key)
if (!run) {
const controller = new AbortController()
run = { controller, users: 0, promise: Promise.resolve() }
const current = run
run.promise = Promise.resolve()
.then(() => {
throwIfSignalAborted(controller.signal)
return refresh(controller.signal)
})
.finally(() => {
if (this.runs.get(key) === current) {
this.runs.delete(key)
}
})
this.runs.set(key, run)
}
run.users++
try {
await waitForPromiseWithSignal(run.promise, signal)
} finally {
run.users--
if (run.users === 0) {
run.controller.abort()
if (this.runs.get(key) === run) {
this.runs.delete(key)
}
}
}
}
cancel(): void {
for (const run of this.runs.values()) {
run.controller.abort()
}
this.runs.clear()
}
}
@@ -0,0 +1,152 @@
import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { expect, it } from 'vitest'
import { SessionSearchStore } from './session-search-store'
import {
deleteExpiredSearchFiles,
RETENTION_DELETE_ROWS_PER_STEP
} from './session-search-retention-delete'
function seed(store: SessionSearchStore, id: number, rows: number, mtime: number) {
const db = store.db
db.prepare(`INSERT INTO sessions(id,agent,session_id,file_path,title,cwd,cwd_key,resume_command)
VALUES (?, 'claude', ?, ?, 'synthetic retention', '/fixture', '/fixture', '')`).run(
id,
String(id),
String(id)
)
db.prepare('INSERT INTO files(path,byte_offset,mtime_ms,session_row_id) VALUES (?,1,?,?)').run(
String(id),
mtime,
id
)
db.exec('BEGIN')
for (let i = 0; i < rows; i++) {
const row = db
.prepare("INSERT INTO messages(session_row_id,role) VALUES (?,'user')")
.run(id).lastInsertRowid
db.prepare('INSERT INTO messages_fts(rowid,user_text) VALUES (?,?)').run(row, 'retentionneedle')
db.prepare('INSERT INTO conversation_fts(rowid,user_text) VALUES (?,?)').run(
row,
'retentionneedle'
)
}
db.exec('COMMIT')
}
it('yields within a large file while hiding partial rows and preserving unrelated sessions', async () => {
const store = new SessionSearchStore(':memory:')
seed(store, 1, 1025, 1)
seed(store, 2, 1, 200)
let previous = 1025
let steps = 0
try {
await deleteExpiredSearchFiles(
store.db,
100,
() => false,
() => {},
async () => {
const left = Number(
(
store.db.prepare('SELECT count(*) AS n FROM messages WHERE session_row_id=1').get() as {
n: number
}
).n
)
expect(previous - left).toBeLessThanOrEqual(RETENTION_DELETE_ROWS_PER_STEP)
expect(previous - left).toBeGreaterThan(0)
previous = left
steps++
expect(store.search({ query: 'retentionneedle' }).hits.map((h) => h.sessionId)).toEqual([
'2'
])
}
)
expect(steps).toBe(5)
expect(store.db.prepare('SELECT count(*) AS n FROM messages_fts').get()).toEqual({ n: 1 })
expect(store.db.prepare('SELECT count(*) AS n FROM conversation_fts').get()).toEqual({ n: 1 })
expect(store.db.prepare('SELECT count(*) AS n FROM search_pending_deletes').get()).toEqual({
n: 0
})
} finally {
store.close()
}
})
it('finishes an interrupted deletion after reopening even when history becomes unlimited', async () => {
const root = await mkdtemp(join(tmpdir(), 'ss-retention-resume-'))
const path = join(root, 'index.sqlite')
let store = new SessionSearchStore(path)
let closed = false
try {
seed(store, 1, 513, 1)
await deleteExpiredSearchFiles(
store.db,
100,
() => closed,
() => {},
async () => {
store.close()
closed = true
}
)
store = new SessionSearchStore(path)
closed = false
expect(store.search({ query: 'retentionneedle' }).hits).toEqual([])
expect(store.coverage().sessionsIndexed).toBe(0)
await store.purgeOlderThan(null)
expect(store.db.prepare('SELECT count(*) AS n FROM messages_fts').get()).toEqual({ n: 0 })
expect(store.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 })
} finally {
if (!closed) {
store.close()
}
await rm(root, { recursive: true, force: true })
}
})
it('does not orphan a replacement file when resuming an older deletion for the same path', async () => {
const store = new SessionSearchStore(':memory:')
try {
seed(store, 1, 2, 1)
store.db.exec(
"INSERT INTO search_pending_deletes VALUES ('1',1); DELETE FROM files WHERE path='1'"
)
seed(store, 2, 2, 1)
store.db.exec("UPDATE files SET path='1' WHERE path='2'")
await store.purgeOlderThan(100)
for (const table of [
'messages',
'messages_fts',
'conversation_fts',
'sessions',
'files',
'search_pending_deletes'
]) {
expect(store.db.prepare(`SELECT count(*) AS n FROM ${table}`).get()).toEqual({ n: 0 })
}
} finally {
store.close()
}
})
it('cancels retention between batches and resumes without exposing a partial session', async () => {
const store = new SessionSearchStore(':memory:')
try {
seed(store, 1, 1025, 1)
const controller = new AbortController()
const purge = store.purgeOlderThan(100, controller.signal)
setImmediate(() => controller.abort())
await purge
const remaining = store.db.prepare('SELECT count(*) AS n FROM messages').get() as { n: number }
expect(remaining.n).toBeGreaterThan(0)
expect(remaining.n).toBeLessThan(1025)
expect(store.search({ query: 'retentionneedle' }).hits).toEqual([])
await store.purgeOlderThan(null)
expect(store.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ n: 0 })
} finally {
store.close()
}
})
@@ -0,0 +1,84 @@
import { setImmediate as yieldToEventLoop } from 'node:timers/promises'
import type SyncDatabase from '../sqlite/sync-database'
import { inSessionParseFileLane } from '../ai-vault/session-parse-file-lane'
export const RETENTION_DELETE_ROWS_PER_STEP = 256
/** A durable tombstone hides partial deletes and lets a reopened store finish them. */
export async function deleteExpiredSearchFiles(
db: SyncDatabase,
cutoffMs: number | null,
closed: () => boolean,
changed: () => void,
yieldStep: () => Promise<void> = yieldToEventLoop
): Promise<void> {
const pending = db.prepare('SELECT path FROM search_pending_deletes').all() as { path: string }[]
const expired =
cutoffMs === null
? []
: (db
.prepare('SELECT path FROM files WHERE mtime_ms < ? ORDER BY mtime_ms')
.all(cutoffMs) as { path: string }[])
for (const { path } of [...pending, ...expired]) {
if (closed()) {
return
}
await inSessionParseFileLane(path, async () => {
if (closed()) {
return
}
db.exec('BEGIN IMMEDIATE')
try {
const file = db
.prepare(`SELECT session_row_id FROM files WHERE path = ? AND mtime_ms < ?
AND path NOT IN (SELECT path FROM search_pending_deletes)`)
.get(path, cutoffMs ?? -Infinity) as { session_row_id: number | null } | undefined
if (file) {
if (file.session_row_id !== null) {
db.prepare(
'INSERT OR IGNORE INTO search_pending_deletes(path, session_row_id) VALUES (?, ?)'
).run(path, file.session_row_id)
}
db.prepare('DELETE FROM files WHERE path = ?').run(path)
}
db.exec('COMMIT')
changed()
} catch (error) {
db.exec('ROLLBACK')
throw error
}
while (!closed()) {
const pending = db
.prepare('SELECT session_row_id FROM search_pending_deletes WHERE path = ?')
.get(path) as { session_row_id: number } | undefined
if (!pending) {
return
}
db.exec('BEGIN IMMEDIATE')
try {
const ids = db
.prepare('SELECT id FROM messages WHERE session_row_id = ? LIMIT ?')
.all(pending.session_row_id, RETENTION_DELETE_ROWS_PER_STEP) as { id: number }[]
const full = db.prepare('DELETE FROM messages_fts WHERE rowid = ?')
const conversation = db.prepare('DELETE FROM conversation_fts WHERE rowid = ?')
const message = db.prepare('DELETE FROM messages WHERE id = ?')
for (const { id } of ids) {
full.run(id)
conversation.run(id)
message.run(id)
}
if (ids.length < RETENTION_DELETE_ROWS_PER_STEP) {
db.prepare('DELETE FROM sessions WHERE id = ?').run(pending.session_row_id)
db.prepare('DELETE FROM search_pending_deletes WHERE path = ?').run(path)
}
db.exec('COMMIT')
changed()
} catch (error) {
db.exec('ROLLBACK')
throw error
}
await yieldStep()
}
})
}
}
@@ -19,7 +19,10 @@ export function sessionRowFilter(
args: AiVaultSearchArgs,
split: AiVaultSearchQuerySplit
): SessionRowFilter {
const filter: SessionRowFilter = { conditions: [], values: [] }
const filter: SessionRowFilter = {
conditions: ['id NOT IN (SELECT session_row_id FROM search_pending_deletes)'],
values: []
}
if (args.agents && args.agents.length > 0) {
filter.conditions.push(`agent IN (${args.agents.map(() => '?').join(',')})`)
filter.values.push(...args.agents)
@@ -2,7 +2,7 @@ import { rmSync } from 'node:fs'
import SyncDatabase from '../sqlite/sync-database'
// Bump to drop and rebuild: the index is a cache over the transcripts, never a source.
export const SESSION_SEARCH_SCHEMA_VERSION = 7
export const SESSION_SEARCH_SCHEMA_VERSION = 8
// unicode61 keeps `_ . - /` inside tokens so paths and identifiers match exactly;
// the `identifiers` column carries the split form (see session-search-identifier-split).
@@ -43,6 +43,10 @@ CREATE TABLE IF NOT EXISTS files(
size_bytes INTEGER,
session_row_id INTEGER
);
CREATE TABLE IF NOT EXISTS search_pending_deletes(
path TEXT PRIMARY KEY,
session_row_id INTEGER NOT NULL UNIQUE
);
CREATE TABLE IF NOT EXISTS messages(
id INTEGER PRIMARY KEY,
session_row_id INTEGER NOT NULL,
@@ -1,3 +1,4 @@
import { SessionSearchRefreshLane } from './session-search-refresh-lane'
import { recordSearchDiscovered } from './session-search-discovered-counts'
import { withCursorChatMetaScan } from '../ai-vault/session-scanner-cursor-chat-meta'
import { mkdirSync } from 'node:fs'
@@ -56,6 +57,7 @@ export class SessionSearchService {
private policy: AiVaultSearchSettings
private store: SessionSearchStore | null = null
private backfillRun: Promise<void> | null = null
private readonly refreshLane = new SessionSearchRefreshLane()
private searchesInFlight = 0
private releaseBackfill: (() => void) | null = null
private backfillController: AbortController | null = null
@@ -87,8 +89,14 @@ export class SessionSearchService {
this.searchesInFlight += 1
try {
if (args.refresh !== false) {
await this.refreshRecent(roots, signal)
await this.reindexStale(signal)
await this.refreshLane.run(
roots,
async (sharedSignal) => {
await this.refreshRecent(roots, sharedSignal)
await this.reindexStale(sharedSignal)
},
signal
)
}
void backfill
return store.search(args)
@@ -232,6 +240,7 @@ export class SessionSearchService {
}
private closeStore(): void {
this.refreshLane.cancel()
const store = this.store
this.store = null
if (!store) {
@@ -287,9 +296,7 @@ export class SessionSearchService {
store.setBackfillState('running')
try {
const cutoff = aiVaultSearchHistoryCutoffMs(this.policy.historyDays)
if (cutoff !== null) {
await store.purgeOlderThan(cutoff)
}
await store.purgeOlderThan(cutoff, signal)
await ensureSessionParseCacheLoaded()
const issues: AiVaultScanIssue[] = []
const options: AiVaultScanOptions = { ...roots, signal }
@@ -1,3 +1,4 @@
import { deleteExpiredSearchFiles } from './session-search-retention-delete'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import { aiVaultSearchHistoryCutoffMs } from '../../shared/ai-vault-search-settings'
import { setImmediate as yieldToEventLoop } from 'node:timers/promises'
@@ -140,33 +141,31 @@ export class SessionSearchStore implements SessionSearchIndexSink {
}
}
/**
* Drops every file last modified before the cutoff, one transaction per file
* so searches interleave, then hands the freed pages back so the index file
* shrinks (a full VACUUM on a multi-GB index takes half a minute).
*/
async purgeOlderThan(cutoffMs: number): Promise<void> {
let paths: string[]
/** Hides expired sessions immediately, then removes their rows in resumable batches. */
async purgeOlderThan(cutoffMs: number | null, signal?: AbortSignal): Promise<void> {
try {
paths = this.writer.filesOlderThan(cutoffMs)
await deleteExpiredSearchFiles(
this.db,
cutoffMs,
() => this.closed || signal?.aborted === true,
() => {
this.providerCounts = null
}
)
if (!this.closed && !signal?.aborted) {
await this.compact(signal)
}
} catch (error) {
this.onError(error)
return
}
for (const path of paths) {
if (this.closed) {
return
if (!this.closed) {
this.onError(error)
}
this.removeFile(path)
await yieldToEventLoop()
}
await this.compact()
}
private async compact(): Promise<void> {
private async compact(signal?: AbortSignal): Promise<void> {
try {
let freed = Number(this.db.pragma('freelist_count', { simple: true }))
while (!this.closed && freed > 0) {
while (!this.closed && !signal?.aborted && freed > 0) {
this.db.pragma(`incremental_vacuum(${COMPACT_PAGES_PER_STEP})`)
const remaining = Number(this.db.pragma('freelist_count', { simple: true }))
// Why: without auto_vacuum the step is a no-op; never spin on it.
@@ -242,6 +241,7 @@ export class SessionSearchStore implements SessionSearchIndexSink {
.prepare(
`SELECT s.agent AS agent, COUNT(DISTINCT s.id) AS sessions, COUNT(m.id) AS messages
FROM sessions s LEFT JOIN messages m ON m.session_row_id = s.id
WHERE s.id NOT IN (SELECT session_row_id FROM search_pending_deletes)
GROUP BY s.agent ORDER BY s.agent`
)
.all() as { agent: AiVaultAgent; sessions: number; messages: number }[])
@@ -0,0 +1,31 @@
import { expect, it } from 'vitest'
import { SessionNewestFiles } from './session-newest-files'
import type { FileWithMtime } from './session-scanner-types'
function file(i: number): FileWithMtime {
const mtimeMs = (i * 7919) % 997
return { path: String(i), mtimeMs, modifiedAt: new Date(mtimeMs).toISOString() }
}
it('retains at most 12 of 100,000 candidates with stable newest-first ties', () => {
const all = Array.from({ length: 100_000 }, (_, i) => file(i))
const retained = new SessionNewestFiles(12)
let peak = 0
for (const candidate of all) {
retained.add(candidate)
peak = Math.max(peak, retained.size)
}
expect(peak).toBe(12)
expect(retained.newest()).toEqual(all.sort((a, b) => b.mtimeMs - a.mtimeMs).slice(0, 12))
})
it('supports full backfill and empty requests', () => {
const all = new SessionNewestFiles(Infinity)
const none = new SessionNewestFiles(0)
for (let i = 0; i < 100; i++) {
all.add(file(i))
none.add(file(i))
}
expect(all.newest()).toHaveLength(100)
expect(none.newest()).toEqual([])
})
+48
View File
@@ -0,0 +1,48 @@
import type { FileWithMtime } from './session-scanner-types'
/** Retain only the requested newest files, preserving traversal order on ties. */
export class SessionNewestFiles {
private readonly files: FileWithMtime[] = []
private readonly limit: number
constructor(limit: number) {
this.limit = limit === Infinity ? limit : Math.max(0, Math.trunc(limit) || 0)
}
add(file: FileWithMtime): void {
if (!Number.isFinite(this.limit)) {
this.files.push(file)
return
}
if (this.limit <= 0) {
return
}
const last = this.files.at(-1)
if (this.files.length >= this.limit && last && file.mtimeMs <= last.mtimeMs) {
return
}
let low = 0
let high = this.files.length
while (low < high) {
const middle = (low + high) >>> 1
if (this.files[middle].mtimeMs >= file.mtimeMs) {
low = middle + 1
} else {
high = middle
}
}
this.files.splice(low, 0, file)
if (this.files.length > this.limit) {
this.files.pop()
}
}
get size(): number {
return this.files.length
}
newest(): FileWithMtime[] {
return this.files.sort((a, b) => b.mtimeMs - a.mtimeMs)
}
}
@@ -61,3 +61,28 @@ describe('walkSessionFiles directory reader', () => {
).rejects.toBe(cancelled)
})
})
it('visits file contents before descending further without retaining paths', async () => {
tempRoot = await mkdtemp(join(tmpdir(), 'orca-session-stream-'))
await writeFile(join(tempRoot, 'first.jsonl'), '{}\n')
await mkdir(join(tempRoot, 'nested'))
await writeFile(join(tempRoot, 'nested', 'second.jsonl'), '{}\n')
const visited: string[] = []
const readDirectory = vi.fn(async (path: string) => {
if (path.endsWith('nested')) {
expect(visited).toEqual([join(tempRoot!, 'first.jsonl')])
}
return (await readdir(path, { withFileTypes: true })).sort((a, b) =>
a.name.localeCompare(b.name)
)
})
const retained = await walkSessionFiles(tempRoot, 'claude', [], {
extensions: new Set(['.jsonl']),
readDirectory,
onFile: async (path) => {
visited.push(path)
}
})
expect(retained).toEqual([])
expect(visited).toHaveLength(2)
})
+39 -41
View File
@@ -1,10 +1,11 @@
import type { Dirent } from 'node:fs'
import { SessionNewestFiles } from './session-newest-files'
import { extname, join } from 'node:path'
import type { AiVaultAgent, AiVaultScanIssue } from '../../shared/ai-vault-types'
import { wslGatedReaddir, wslGatedStat } from '../native-chat/wsl-transcript-fs-access'
import { WslTranscriptFsError } from '../native-chat/wsl-transcript-fs-gate'
import { recordSessionScanIssue } from './session-scan-issues'
import type { FileWithMtime, SessionFileDiscovery } from './session-scanner-types'
import type { SessionFileDiscovery } from './session-scanner-types'
import { errorMessage } from './session-scanner-values'
export async function discoverFiles(args: {
@@ -18,18 +19,42 @@ export async function discoverFiles(args: {
contentDependencyPath?: (path: string) => string | undefined | Promise<string | undefined>
directoryPredicate?: (name: string, depth: number) => boolean
}): Promise<SessionFileDiscovery> {
let paths: string[]
const files = new SessionNewestFiles(args.limit)
try {
paths = await walkSessionFiles(args.rootDir, args.agent, args.issues, {
await walkSessionFiles(args.rootDir, args.agent, args.issues, {
extensions: new Set(args.extensions),
signal: args.signal,
filePredicate: args.filePredicate,
directoryPredicate: args.directoryPredicate
directoryPredicate: args.directoryPredicate,
onFile: async (path) => {
args.signal?.throwIfAborted()
try {
const fileStat = await wslGatedStat(path, 'scan')
const dependencyStat = await optionalContentDependencyStat(
await args.contentDependencyPath?.(path)
)
args.signal?.throwIfAborted()
const mtimeMs = Math.max(fileStat.mtimeMs, dependencyStat?.mtimeMs ?? 0)
files.add({
path,
mtimeMs,
modifiedAt: new Date(mtimeMs).toISOString(),
sizeBytes: fileStat.size + (dependencyStat?.size ?? 0),
dev: fileStat.dev,
ino: fileStat.ino,
nlink: fileStat.nlink
})
} catch (err) {
args.signal?.throwIfAborted()
recordSessionScanIssue(args.issues, {
agent: args.agent,
path,
message: errorMessage(err)
})
}
}
})
} catch (err) {
// Why: discoverAiVaultSessionSources fans out with Promise.all, so one
// stalled distro would otherwise reject the whole vault scan — including
// every healthy local agent. Contain it to this root.
if (!(err instanceof WslTranscriptFsError)) {
throw err
}
@@ -40,39 +65,7 @@ export async function discoverFiles(args: {
})
return { agent: args.agent, rootDir: args.rootDir, files: [] }
}
const files: FileWithMtime[] = []
for (const path of paths) {
args.signal?.throwIfAborted()
try {
const fileStat = await wslGatedStat(path, 'scan')
const dependencyStat = await optionalContentDependencyStat(
await args.contentDependencyPath?.(path)
)
args.signal?.throwIfAborted()
const mtimeMs = Math.max(fileStat.mtimeMs, dependencyStat?.mtimeMs ?? 0)
files.push({
path,
mtimeMs,
modifiedAt: new Date(mtimeMs).toISOString(),
sizeBytes: fileStat.size + (dependencyStat?.size ?? 0),
dev: fileStat.dev,
ino: fileStat.ino,
nlink: fileStat.nlink
})
} catch (err) {
args.signal?.throwIfAborted()
recordSessionScanIssue(args.issues, {
agent: args.agent,
path,
message: errorMessage(err)
})
}
}
return {
agent: args.agent,
rootDir: args.rootDir,
files: files.sort((left, right) => right.mtimeMs - left.mtimeMs).slice(0, args.limit)
}
return { agent: args.agent, rootDir: args.rootDir, files: files.newest() }
}
async function optionalContentDependencyStat(
@@ -104,6 +97,7 @@ export async function walkSessionFiles(
directoryPredicate?: (name: string, depth: number) => boolean
readDirectory?: (dirPath: string) => Promise<Dirent[]>
signal?: AbortSignal
onFile?: (path: string) => Promise<void>
},
depth = 0
): Promise<string[]> {
@@ -140,7 +134,11 @@ export async function walkSessionFiles(
options.extensions.has(extname(entry.name).toLowerCase()) &&
(options.filePredicate?.(fullPath) ?? true)
) {
files.push(fullPath)
if (options.onFile) {
await options.onFile(fullPath)
} else {
files.push(fullPath)
}
}
}
return files
@@ -104,3 +104,25 @@ describe('useAiVaultSearchCoveragePoll', () => {
expect(result.current).toBeNull()
})
})
it('drops coverage from the previous runtime immediately and polls the new owner', async () => {
const { result, rerender } = renderHook(
({ host }) => useAiVaultSearchCoveragePoll(true, null, host),
{ initialProps: { host: 'runtime:a' }, wrapper }
)
await act(async () => {})
expect(result.current?.sessionsIndexed).toBe(5)
let release!: (value: AiVaultSearchCoverage) => void
searchCoverage.mockImplementation(
() =>
new Promise((resolve) => {
release = resolve
})
)
rerender({ host: 'runtime:b' })
expect(result.current).toBeNull()
await act(async () => {
release({ ...coverage('complete'), sessionsIndexed: 9 })
})
expect(result.current?.sessionsIndexed).toBe(9)
})
@@ -13,9 +13,11 @@ export const AI_VAULT_SEARCH_COVERAGE_POLL_MS = 4_000
*/
export function useAiVaultSearchCoveragePoll(
enabled: boolean,
latest: AiVaultSearchCoverage | null = null
latest: AiVaultSearchCoverage | null = null,
ownerKey = ''
): AiVaultSearchCoverage | null {
const [snapshot, setSnapshot] = useState<{
ownerKey: string
source: AiVaultSearchCoverage | null
value: AiVaultSearchCoverage
} | null>(null)
@@ -41,7 +43,7 @@ export function useAiVaultSearchCoveragePoll(
if (stopped || issued !== generation) {
return
}
setSnapshot({ source: latest, value: next })
setSnapshot({ ownerKey, source: latest, value: next })
})
.catch(() => undefined)
}
@@ -50,7 +52,11 @@ export function useAiVaultSearchCoveragePoll(
stopped = true
clearInterval(interval)
}
}, [enabled, latest])
}, [enabled, latest, ownerKey])
return enabled ? (snapshot?.source === latest ? snapshot.value : latest) : null
return enabled
? snapshot?.ownerKey === ownerKey && snapshot.source === latest
? snapshot.value
: latest
: null
}
@@ -214,3 +214,38 @@ describe('useAiVaultSessionSearchRequest', () => {
expect(searchSessions.mock.calls[0]?.[0]).toMatchObject({ query: 'alpha', tier: 'full' })
})
})
it('retires retained and in-flight results when the execution host changes with the same query', async () => {
const resolvers: ((result: AiVaultSearchResult) => void)[] = []
searchSessions.mockImplementation(
() => new Promise<AiVaultSearchResult>((resolve) => resolvers.push(resolve))
)
const { rerender, result } = renderHook(
({ host }: { host: string }) => useAiVaultSessionSearchRequest(argsFor('alpha'), 0, host),
{ initialProps: { host: 'runtime:a' }, wrapper }
)
act(() => {
vi.advanceTimersByTime(AI_VAULT_SEARCH_TYPING_DELAY_MS)
})
await act(async () => {
resolvers[0](searchResult('host-a'))
})
expect(result.current.result?.repairedTerms).toEqual(['host-a'])
act(() => {
vi.advanceTimersByTime(AI_VAULT_SEARCH_SETTLED_DELAY_MS)
})
rerender({ host: 'runtime:b' })
expect(result.current.result).toBeNull()
await act(async () => {
resolvers[1](searchResult('late-host-a'))
})
expect(result.current.result).toBeNull()
act(() => {
vi.advanceTimersByTime(AI_VAULT_SEARCH_SETTLED_DELAY_MS)
})
await act(async () => {
resolvers.at(-1)!(searchResult('host-b'))
})
expect(result.current.result?.repairedTerms).toEqual(['host-b'])
expect(searchSessions).toHaveBeenCalledTimes(4)
})
@@ -19,6 +19,7 @@ export type AiVaultSearchRequestState = {
}
type SettledSearch = {
ownerKey: string
full: boolean
key: string
result: AiVaultSearchResult | null
@@ -40,7 +41,8 @@ type SettledSearch = {
export function useAiVaultSessionSearchRequest(
args: AiVaultSearchArgs | null,
/** Bumped on Enter: runs the full tier now instead of waiting out the debounce. */
flushSignal = 0
flushSignal = 0,
ownerKey = ''
): AiVaultSearchRequestState {
const [settled, setSettled] = useState<SettledSearch | null>(null)
const sequenceRef = useRef(0)
@@ -48,31 +50,35 @@ export function useAiVaultSessionSearchRequest(
// restart the debounce on every parent render.
const argsKey = args ? JSON.stringify(args) : ''
const issue = useCallback((key: string, overrides: Partial<AiVaultSearchArgs>): void => {
if (!key) {
return
}
const requestArgs = JSON.parse(key) as AiVaultSearchArgs
sequenceRef.current += 1
const sequence = sequenceRef.current
void window.api.aiVault
.searchSessions({ ...requestArgs, ...overrides })
.then((result) => {
if (sequenceRef.current === sequence) {
setSettled({ key, result, error: null, full: overrides.tier === 'full' })
}
})
.catch((error: unknown) => {
if (sequenceRef.current === sequence) {
setSettled({
key,
full: true,
result: null,
error: error instanceof Error ? error.message : String(error)
})
}
})
}, [])
const issue = useCallback(
(key: string, overrides: Partial<AiVaultSearchArgs>): void => {
if (!key) {
return
}
const requestArgs = JSON.parse(key) as AiVaultSearchArgs
sequenceRef.current += 1
const sequence = sequenceRef.current
void window.api.aiVault
.searchSessions({ ...requestArgs, ...overrides })
.then((result) => {
if (sequenceRef.current === sequence) {
setSettled({ ownerKey, key, result, error: null, full: overrides.tier === 'full' })
}
})
.catch((error: unknown) => {
if (sequenceRef.current === sequence) {
setSettled({
ownerKey,
key,
full: true,
result: null,
error: error instanceof Error ? error.message : String(error)
})
}
})
},
[ownerKey]
)
// Why: a flush has to cancel the debounce the args effect armed, and the two
// effects cannot share locals, so the args effect publishes its cancel.
@@ -114,8 +120,9 @@ export function useAiVaultSessionSearchRequest(
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [flushSignal, issue])
const current = settled?.key === argsKey ? settled : null
const previous = current === null && argsKey !== '' ? (settled?.result ?? null) : null
const owned = settled?.ownerKey === ownerKey ? settled : null
const current = owned?.key === argsKey ? owned : null
const previous = current === null && argsKey !== '' ? (owned?.result ?? null) : null
return {
result: current?.result ?? previous,
loading: argsKey !== '' && current === null && previous === null,
@@ -79,12 +79,17 @@ export function useAiVaultSessionSearchResults(input: {
}
}, [agents, enabled, newestFirst, query, scopePaths, supportedHost])
const { error, loading, result, updating } = useAiVaultSessionSearchRequest(args, flushSignal)
const { error, loading, result, updating } = useAiVaultSessionSearchRequest(
args,
flushSignal,
executionHostScope
)
// With an empty box no search runs, so the panel reads coverage directly to
// report what is already searchable while the backfill is still going.
const polledCoverage = useAiVaultSearchCoveragePoll(
enabled && supportedHost,
result?.coverage ?? null
result?.coverage ?? null,
executionHostScope
)
// Desktop search always reads this machine's index; a paired web client's
// reads its runtime host, which is the scope it is pinned to.
@@ -0,0 +1,69 @@
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
import {
installBrowserGlobals,
writeStoredRuntimeEnvironment
} from './web-preload-api-test-harness'
import type { RuntimeRpcResponse } from '../../../shared/runtime-rpc-envelope'
beforeEach(() => vi.resetModules())
afterEach(() => {
vi.unstubAllGlobals()
vi.doUnmock('./web-runtime-client')
})
it('preserves legacy search results and surfaces old-host or disconnected errors without local fallback', async () => {
const calls: { method: string; params: unknown }[] = []
let failure: string | null = null
const legacy = {
hits: [],
route: 'or',
durationMs: 1,
coverage: {
sessionsIndexed: 1,
messagesIndexed: 2,
providers: [],
backfill: 'complete',
filesPending: 0,
lastIndexedAt: null
}
}
vi.doMock('./web-runtime-client', () => ({
WebRuntimeClient: class {
call(method: string, params: unknown): Promise<RuntimeRpcResponse<unknown>> {
calls.push({ method, params })
return Promise.resolve(
failure
? {
id: 'fixture',
_meta: { runtimeId: 'host-a' },
ok: false,
error: { code: failure, message: failure }
}
: { id: 'fixture', _meta: { runtimeId: 'host-a' }, ok: true, result: legacy }
)
}
close(): void {}
}
}))
const globals = installBrowserGlobals('Linux')
writeStoredRuntimeEnvironment(globals.storage, 'host-a')
const { installWebPreloadApi } = await import('./web-preload-api')
installWebPreloadApi()
await expect(globals.window.api.aiVault.searchSessions({ query: 'needle' })).resolves.toEqual(
legacy
)
expect(calls).toEqual([
{
method: 'aiVault.searchSessions',
params: { query: 'needle', executionHostId: 'runtime:host-a' }
}
])
for (const error of ['method_not_found', 'connection_closed']) {
failure = error
await expect(
globals.window.api.aiVault.searchSessions({ query: 'needle' })
).rejects.toMatchObject({ code: error })
}
expect(calls).toHaveLength(3)
expect(calls.every((call) => call.method === 'aiVault.searchSessions')).toBe(true)
})