mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 08:02:43 +00:00
feat(ai-vault-search): read a candidate set into the index, and say what happened
One pass over a candidate list, plus the four things a caller has to be told about afterwards: what was discovered, which roots could not be read, which sources are provably gone, and where the whole run stands. Two absence rules are load-bearing. The file walker swallows a readdir failure and returns, so an unreadable root and an uninstalled agent both arrive as "no files"; only a root that yielded nothing is re-probed, and only ENOENT counts as absent. A source is retired on the same evidence and no weaker: an EACCES, an EIO or a stalled distro keeps its rows. Progress lives in the index's own files table, so a pass skips what it already covers and an interrupted run resumes instead of starting over.
This commit is contained in:
@@ -0,0 +1,84 @@
|
||||
import { mkdirSync } from 'node:fs'
|
||||
import { chmod, mkdir, rm, writeFile } from 'node:fs/promises'
|
||||
import { dirname, 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'
|
||||
import {
|
||||
openSessionSearchIndexerHarness,
|
||||
type SessionSearchIndexerHarness
|
||||
} from './session-search-indexer-test-fixture'
|
||||
import { SessionSearchStore } from './session-search-store'
|
||||
|
||||
const CAN_DENY_READ = process.platform !== 'win32' && process.getuid?.() !== 0
|
||||
|
||||
let harness: SessionSearchIndexerHarness
|
||||
let store: SessionSearchStore
|
||||
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.
|
||||
store.removeFile = (path: string) => removed.push(path)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
store.close()
|
||||
await harness.cleanup()
|
||||
})
|
||||
|
||||
it('retires only what is proven gone and re-watches what it could not stat', async () => {
|
||||
const present = join(harness.root, 'present.jsonl')
|
||||
await writeFile(present, '{}')
|
||||
const result = await retireDeletedSessionSearchSources(store, [
|
||||
present,
|
||||
join(harness.root, 'never-existed.jsonl')
|
||||
])
|
||||
expect(result.retired).toEqual([join(harness.root, 'never-existed.jsonl')])
|
||||
expect(removed).toEqual([join(harness.root, 'never-existed.jsonl')])
|
||||
expect(result.unverifiable).toEqual([])
|
||||
})
|
||||
|
||||
it.skipIf(!CAN_DENY_READ)(
|
||||
'keeps rows for an unreadable source rather than calling it deleted',
|
||||
async () => {
|
||||
const blocked = join(harness.root, 'blocked')
|
||||
await mkdir(blocked)
|
||||
const hidden = join(blocked, 'hidden.jsonl')
|
||||
await writeFile(hidden, '{}')
|
||||
await chmod(blocked, 0o000)
|
||||
try {
|
||||
const result = await retireDeletedSessionSearchSources(store, [hidden])
|
||||
expect(result.retired).toEqual([])
|
||||
expect(result.unverifiable).toEqual([hidden])
|
||||
expect(removed).toEqual([])
|
||||
} finally {
|
||||
await chmod(blocked, 0o755)
|
||||
}
|
||||
}
|
||||
)
|
||||
|
||||
it('stats the database behind a synthetic OpenCode row, not the row itself', async () => {
|
||||
const db = join(harness.root, 'opencode.db')
|
||||
await writeFile(db, '')
|
||||
const alive = `${db}#session-1`
|
||||
await expect(retireDeletedSessionSearchSources(store, [alive])).resolves.toMatchObject({
|
||||
retired: []
|
||||
})
|
||||
|
||||
await rm(db)
|
||||
await expect(retireDeletedSessionSearchSources(store, [alive])).resolves.toMatchObject({
|
||||
retired: [alive]
|
||||
})
|
||||
})
|
||||
|
||||
it('caps the stats one cycle spends and leaves the rest to be checked again', async () => {
|
||||
const paths = ['/a', '/b', '/c', '/d'].map((name) => join(harness.root, name))
|
||||
const result = await retireDeletedSessionSearchSources(store, paths, { limit: 2 })
|
||||
expect(result.retired).toEqual(paths.slice(0, 2))
|
||||
expect(result.unchecked).toEqual(paths.slice(2))
|
||||
})
|
||||
@@ -0,0 +1,60 @@
|
||||
import { wslGatedStat } from '../native-chat/wsl-transcript-fs-access'
|
||||
import { splitOpenCodeSqliteCandidate } from '../ai-vault/session-scanner-opencode-sqlite-paths'
|
||||
import type { SessionSearchStore } from './session-search-store'
|
||||
|
||||
export type SessionSearchRetirement = {
|
||||
/** Paths proven gone and dropped from the index. */
|
||||
retired: string[]
|
||||
/** Paths that could not be statted; their rows stay, and they stay watched. */
|
||||
unverifiable: string[]
|
||||
/** Paths the per-cycle cap left for next time. */
|
||||
unchecked: string[]
|
||||
}
|
||||
|
||||
// Why only ENOENT retires a source: docs/reference/ssh-execution-boundary.md —
|
||||
// loss of contact is never evidence of absence. An EACCES, an EIO or a stalled
|
||||
// WSL distro leaves the rows exactly where they are, because the alternative is
|
||||
// erasing a user's searchable history the first time a mount hiccups.
|
||||
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.
|
||||
*/
|
||||
export async function retireDeletedSessionSearchSources(
|
||||
store: SessionSearchStore,
|
||||
paths: readonly string[],
|
||||
options: { signal?: AbortSignal; limit?: number } = {}
|
||||
): Promise<SessionSearchRetirement> {
|
||||
const { signal } = options
|
||||
const retirement: SessionSearchRetirement = { retired: [], unverifiable: [], unchecked: [] }
|
||||
const limit = options.limit ?? Number.POSITIVE_INFINITY
|
||||
for (const [index, path] of paths.entries()) {
|
||||
// Why capped: the sweep hands over every path it discovered, and one stat
|
||||
// per transcript on a 3,600-session machine is not a cycle's worth of work.
|
||||
// What is left keeps being watched, so the check finishes over a few cycles.
|
||||
if (signal?.aborted || index >= limit) {
|
||||
retirement.unchecked.push(...paths.slice(index))
|
||||
break
|
||||
}
|
||||
// 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
|
||||
try {
|
||||
await wslGatedStat(statPath, 'scan', signal)
|
||||
} catch (error) {
|
||||
const code =
|
||||
error && typeof error === 'object' && 'code' in error && typeof error.code === 'string'
|
||||
? error.code
|
||||
: null
|
||||
if (code !== null && PROVEN_GONE.has(code)) {
|
||||
store.removeFile(path)
|
||||
retirement.retired.push(path)
|
||||
continue
|
||||
}
|
||||
retirement.unverifiable.push(path)
|
||||
}
|
||||
}
|
||||
return retirement
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
import { expect, it } from 'vitest'
|
||||
import type { SessionFileDiscovery } from '../ai-vault/session-scanner-types'
|
||||
import { sessionSearchDiscoveredCounts } from './session-search-discovered-counts'
|
||||
|
||||
function discovery(agent: SessionFileDiscovery['agent'], files: number): SessionFileDiscovery {
|
||||
return {
|
||||
agent,
|
||||
rootDir: `/roots/${agent}`,
|
||||
files: Array.from({ length: files }, (_unused, index) => ({
|
||||
path: `/roots/${agent}/${index}.jsonl`,
|
||||
mtimeMs: 0,
|
||||
modifiedAt: new Date(0).toISOString()
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
it('keeps a provider whose every read failed apart from one that is not installed', () => {
|
||||
const counts = sessionSearchDiscoveredCounts(
|
||||
[discovery('claude', 2), discovery('claude', 1), discovery('codex', 0)],
|
||||
[
|
||||
{ agent: 'codex', path: '/roots/codex', message: 'unreadable' },
|
||||
{ agent: 'gemini', path: '/roots/gemini', message: 'unreadable' },
|
||||
{ agent: 'claude', path: '/roots/claude', kind: 'notice', message: 'list truncated' }
|
||||
]
|
||||
)
|
||||
expect(Object.fromEntries(counts)).toEqual({
|
||||
claude: { files: 3, failures: 0 },
|
||||
codex: { files: 0, failures: 1 },
|
||||
gemini: { files: 0, failures: 1 }
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,33 @@
|
||||
import type { AiVaultAgent, AiVaultScanIssue } from '../../shared/ai-vault-types'
|
||||
import type { SessionFileDiscovery } from '../ai-vault/session-scanner-types'
|
||||
|
||||
/** What one sweep saw for one agent, before anything was read. */
|
||||
export type SessionSearchDiscoveredCount = { files: number; failures: number }
|
||||
|
||||
/**
|
||||
* Keeps a provider with discovered transcripts but no indexed rows visible.
|
||||
* Without it a provider whose every file failed to parse is indistinguishable
|
||||
* from one that is not installed, and the difference is the whole question a
|
||||
* coverage report answers.
|
||||
*/
|
||||
export function sessionSearchDiscoveredCounts(
|
||||
discoveries: readonly SessionFileDiscovery[],
|
||||
issues: readonly AiVaultScanIssue[]
|
||||
): Map<AiVaultAgent, SessionSearchDiscoveredCount> {
|
||||
const counts = new Map<AiVaultAgent, SessionSearchDiscoveredCount>()
|
||||
const entry = (agent: AiVaultAgent): SessionSearchDiscoveredCount => {
|
||||
const existing = counts.get(agent) ?? { files: 0, failures: 0 }
|
||||
counts.set(agent, existing)
|
||||
return existing
|
||||
}
|
||||
for (const discovery of discoveries) {
|
||||
entry(discovery.agent).files += discovery.files.length
|
||||
}
|
||||
for (const issue of issues) {
|
||||
// 'notice' rows are scanner commentary, never a failed read.
|
||||
if (issue.kind !== 'notice') {
|
||||
entry(issue.agent).failures += 1
|
||||
}
|
||||
}
|
||||
return counts
|
||||
}
|
||||
@@ -0,0 +1,106 @@
|
||||
import { mkdirSync } from 'node:fs'
|
||||
import { appendFile, rm } from 'node:fs/promises'
|
||||
import { dirname } 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'
|
||||
import { registerSessionSearchIndexConsumer } from './session-search-index-consumer'
|
||||
import { runSessionSearchIndexPass } from './session-search-index-pass'
|
||||
import {
|
||||
claudeLines,
|
||||
openSessionSearchIndexerHarness,
|
||||
writeClaudeTranscript,
|
||||
type SessionSearchIndexerHarness
|
||||
} from './session-search-indexer-test-fixture'
|
||||
import { SessionSearchCycleAllowance } from './session-search-reconcile-budget'
|
||||
import { discoverSessionSearchCandidates } from './session-search-scan-roots'
|
||||
import { SessionSearchStore } from './session-search-store'
|
||||
|
||||
const FIRST = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee'
|
||||
const SECOND = 'bbbbbbbb-cccc-4ddd-8eee-ffffffffffff'
|
||||
|
||||
let harness: SessionSearchIndexerHarness
|
||||
let store: SessionSearchStore
|
||||
let errors: unknown[]
|
||||
|
||||
beforeEach(async () => {
|
||||
resetSessionParseCacheForTests()
|
||||
resetTranscriptConsumersForTests()
|
||||
errors = []
|
||||
harness = await openSessionSearchIndexerHarness('ss-index-pass')
|
||||
await writeClaudeTranscript(transcript(FIRST), ['the first transcript'], FIRST)
|
||||
await writeClaudeTranscript(transcript(SECOND), ['the second transcript'], SECOND)
|
||||
store = openStore()
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
resetTranscriptConsumersForTests()
|
||||
store.close()
|
||||
await harness.cleanup()
|
||||
})
|
||||
|
||||
function transcript(sessionId: string): string {
|
||||
return join(harness.claudeProjectDir, `${sessionId}.jsonl`)
|
||||
}
|
||||
|
||||
function openStore(): SessionSearchStore {
|
||||
mkdirSync(dirname(harness.databasePath), { recursive: true })
|
||||
const opened = new SessionSearchStore(harness.databasePath, (error) => errors.push(error))
|
||||
registerSessionSearchIndexConsumer(opened)
|
||||
return opened
|
||||
}
|
||||
|
||||
async function candidates() {
|
||||
return (
|
||||
await discoverSessionSearchCandidates(harness.roots, {
|
||||
limitPerAgent: Number.POSITIVE_INFINITY
|
||||
})
|
||||
).candidates
|
||||
}
|
||||
|
||||
it('re-reads nothing it already holds, even with a cold session-list cache', async () => {
|
||||
const first = await runSessionSearchIndexPass(store, await candidates())
|
||||
expect(first.stats.fullParses).toBe(2)
|
||||
|
||||
// A restart: the parse cache is gone, the index's `files` table is not.
|
||||
store.close()
|
||||
resetTranscriptConsumersForTests()
|
||||
resetSessionParseCacheForTests()
|
||||
store = openStore()
|
||||
|
||||
const second = await runSessionSearchIndexPass(store, await candidates())
|
||||
expect(second.stats).toMatchObject({ fullParses: 0, incremental: 0, reused: 0, bytesRead: 0 })
|
||||
expect(errors).toEqual([])
|
||||
})
|
||||
|
||||
it('resumes into a grown transcript instead of re-reading it whole', async () => {
|
||||
await runSessionSearchIndexPass(store, await candidates())
|
||||
await appendFile(transcript(FIRST), `${claudeLines(['a later turn'], FIRST, 10).join('\n')}\n`)
|
||||
|
||||
const second = await runSessionSearchIndexPass(store, await candidates())
|
||||
expect(second.stats).toMatchObject({ incremental: 1, fullParses: 0 })
|
||||
})
|
||||
|
||||
it('hands back everything the allowance had no room for, in order', async () => {
|
||||
const allowance = new SessionSearchCycleAllowance({ files: 1, bytes: 1_000_000 })
|
||||
const all = await candidates()
|
||||
const pass = await runSessionSearchIndexPass(store, all, { allowance })
|
||||
|
||||
expect(pass.deferred.map((one) => one.file.path)).toEqual(
|
||||
all.slice(1).map((one) => one.file.path)
|
||||
)
|
||||
expect(
|
||||
harness.read((db) => db.prepare('SELECT count(*) AS n FROM visible_sessions').get())
|
||||
).toEqual({ n: 1 })
|
||||
})
|
||||
|
||||
it('skips a source the reader cannot even open without failing the pass', async () => {
|
||||
const all = await candidates()
|
||||
await rm(transcript(FIRST))
|
||||
const pass = await runSessionSearchIndexPass(store, all)
|
||||
expect(pass.deferred).toEqual([])
|
||||
expect(
|
||||
harness.read((db) => db.prepare('SELECT count(*) AS n FROM visible_sessions').get())
|
||||
).toEqual({ n: 1 })
|
||||
})
|
||||
@@ -0,0 +1,143 @@
|
||||
import { throwIfAiVaultScanCancelled } from '../ai-vault/ai-vault-scan-cancellation'
|
||||
import { parserPublishesMessages } from '../ai-vault/session-scanner-agent-parser'
|
||||
import {
|
||||
createSessionParseStats,
|
||||
invalidateSessionParseCacheEntry,
|
||||
parseAgentSessionFileCached,
|
||||
sessionParseCacheCoversTranscript,
|
||||
type SessionParseStats
|
||||
} from '../ai-vault/session-scanner-parse-cache'
|
||||
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'
|
||||
import type { SessionSearchStore } from './session-search-store'
|
||||
|
||||
export type SessionSearchIndexPassOptions = {
|
||||
signal?: AbortSignal
|
||||
/** Cycle allowance; work that does not fit comes back as `deferred`. */
|
||||
allowance?: SessionSearchCycleAllowance
|
||||
/** Paths whose stored cursor must not be trusted, so the read is forced whole. */
|
||||
forced?: ReadonlySet<string>
|
||||
/** Sleeps between batches so an unasked backfill never owns the CPU. */
|
||||
pace?: (signal?: AbortSignal) => Promise<void>
|
||||
onIndexed?: (candidate: SessionFileCandidate, bytes: number) => void
|
||||
onFailed?: (candidate: SessionFileCandidate) => void
|
||||
}
|
||||
|
||||
export type SessionSearchIndexPassResult = {
|
||||
stats: SessionParseStats
|
||||
/** Candidates the allowance had no room for, in the order they were queued. */
|
||||
deferred: SessionFileCandidate[]
|
||||
}
|
||||
|
||||
const FILES_PER_PACE = 8
|
||||
|
||||
/**
|
||||
* Reads a candidate list through the transcript reader so the registered index
|
||||
* consumer folds it. Candidates the index already covers at their current stat
|
||||
* are skipped outright, which is what makes a restart resume: the `files` table
|
||||
* survives the process, so a second start re-reads nothing it already holds.
|
||||
*/
|
||||
export async function runSessionSearchIndexPass(
|
||||
store: SessionSearchStore,
|
||||
candidates: readonly SessionFileCandidate[],
|
||||
options: SessionSearchIndexPassOptions = {}
|
||||
): Promise<SessionSearchIndexPassResult> {
|
||||
const stats = createSessionParseStats()
|
||||
const deferred: SessionFileCandidate[] = []
|
||||
let sincePace = 0
|
||||
for (const [index, candidate] of candidates.entries()) {
|
||||
throwIfAiVaultScanCancelled(options.signal)
|
||||
const forced = mustReadWhole(store, candidate, options.forced)
|
||||
if (!indexable(store, candidate, forced)) {
|
||||
continue
|
||||
}
|
||||
const bytes = forced ? (candidate.file.sizeBytes ?? 0) : unreadBytes(store, candidate)
|
||||
if (options.allowance && !options.allowance.spend(bytes)) {
|
||||
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.
|
||||
if (forced || sessionParseCacheCoversTranscript(candidate, process.platform)) {
|
||||
invalidateSessionParseCacheEntry(candidate.file.path)
|
||||
}
|
||||
try {
|
||||
await parseAgentSessionFileCached(candidate, process.platform, stats)
|
||||
options.onIndexed?.(candidate, bytes)
|
||||
} catch (error) {
|
||||
throwIfAiVaultScanCancelled(options.signal)
|
||||
options.onFailed?.(candidate)
|
||||
console.warn(
|
||||
'[ai-vault-search] indexing skipped',
|
||||
candidate.agent,
|
||||
error instanceof Error ? error.name : 'ParseError'
|
||||
)
|
||||
}
|
||||
if (++sincePace >= FILES_PER_PACE) {
|
||||
sincePace = 0
|
||||
await options.pace?.(options.signal)
|
||||
}
|
||||
}
|
||||
return { stats, deferred }
|
||||
}
|
||||
|
||||
function indexable(
|
||||
store: SessionSearchStore,
|
||||
candidate: SessionFileCandidate,
|
||||
forced: boolean
|
||||
): boolean {
|
||||
if (!store.acceptsCandidate(candidate)) {
|
||||
return false
|
||||
}
|
||||
// A parser that decodes where the message channel cannot reach it can never
|
||||
// extend the index, so reading it here would be pure cost.
|
||||
if (!parserPublishesMessages(candidate)) {
|
||||
return false
|
||||
}
|
||||
if (forced) {
|
||||
return true
|
||||
}
|
||||
return !isSessionSearchFileCurrent(
|
||||
store.indexedFile(candidate.file.path, fileIdentity(candidate.file)),
|
||||
candidate.file
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* True when the rows under this path describe a file that no longer exists:
|
||||
* the same name now carries a different dev/ino, or it is shorter than the
|
||||
* offset the index read to. Either way an append would splice two files
|
||||
* together, so the read has to start over.
|
||||
*/
|
||||
function mustReadWhole(
|
||||
store: SessionSearchStore,
|
||||
candidate: SessionFileCandidate,
|
||||
forced?: ReadonlySet<string>
|
||||
): boolean {
|
||||
if (forced?.has(candidate.file.path)) {
|
||||
return true
|
||||
}
|
||||
const stored = store.indexedFile(candidate.file.path, null)
|
||||
if (!stored) {
|
||||
return false
|
||||
}
|
||||
if (!store.indexedFile(candidate.file.path, fileIdentity(candidate.file))) {
|
||||
return true
|
||||
}
|
||||
const size = candidate.file.sizeBytes
|
||||
return typeof size === 'number' && stored.byteOffset > size
|
||||
}
|
||||
|
||||
/** 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
|
||||
const indexed = store.indexedFile(candidate.file.path, fileIdentity(candidate.file))
|
||||
if (!indexed || indexed.byteOffset > size) {
|
||||
return size
|
||||
}
|
||||
return size - indexed.byteOffset
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
import { expect, it } from 'vitest'
|
||||
import { SessionSearchIndexingStatus } from './session-search-indexing-status'
|
||||
|
||||
it('walks discovering to indexing to current, and reports degraded roots once settled', () => {
|
||||
const status = new SessionSearchIndexingStatus()
|
||||
status.beginSweep()
|
||||
expect(status.snapshot()).toMatchObject({ phase: 'discovering', filesTotal: null })
|
||||
|
||||
status.planned(2, 0)
|
||||
status.indexed(1_000)
|
||||
expect(status.snapshot()).toMatchObject({ phase: 'indexing', filesIndexed: 1, filesTotal: 2 })
|
||||
|
||||
status.setDegradedRoots([{ root: '/blocked', reason: 'EACCES' }])
|
||||
// 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.finishWork(123)
|
||||
expect(status.snapshot()).toMatchObject({ phase: 'degraded', lastReconcileAt: 123 })
|
||||
|
||||
status.setDegradedRoots([])
|
||||
expect(status.snapshot().phase).toBe('current')
|
||||
})
|
||||
|
||||
it('reports paused over everything, and keeps the counters it had', () => {
|
||||
const status = new SessionSearchIndexingStatus()
|
||||
status.beginSweep()
|
||||
status.planned(3, 1)
|
||||
status.indexed(10)
|
||||
status.setPending(4, 2)
|
||||
status.setPaused(true)
|
||||
expect(status.snapshot()).toMatchObject({
|
||||
phase: 'paused',
|
||||
filesIndexed: 1,
|
||||
filesTotal: 3,
|
||||
failures: 1,
|
||||
filesPending: 4,
|
||||
droppedPending: 2
|
||||
})
|
||||
})
|
||||
|
||||
it('starts a new sweep from zero but keeps what only an open can know', () => {
|
||||
const status = new SessionSearchIndexingStatus()
|
||||
status.setRecoveredRows(7)
|
||||
status.beginSweep()
|
||||
status.planned(1, 0)
|
||||
status.indexed(500)
|
||||
status.failed()
|
||||
status.beginSweep()
|
||||
expect(status.snapshot()).toMatchObject({
|
||||
filesIndexed: 0,
|
||||
bytesIndexed: 0,
|
||||
failures: 0,
|
||||
recoveredRows: 7
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,127 @@
|
||||
import type { AiVaultAgent } from '../../shared/ai-vault-types'
|
||||
import type { SessionSearchDiscoveredCount } from './session-search-discovered-counts'
|
||||
import type { SessionSearchDegradedRoot } from './session-search-root-health'
|
||||
|
||||
/**
|
||||
* `current` is the honest ceiling: the reconciler promises the newest N per
|
||||
* agent within one interval, not every transcript on the machine, so nothing
|
||||
* here ever claims the whole index is up to date.
|
||||
*/
|
||||
export type SessionSearchIndexPhase = 'discovering' | 'indexing' | 'current' | 'paused' | 'degraded'
|
||||
|
||||
export type SessionSearchIndexStatus = {
|
||||
phase: SessionSearchIndexPhase
|
||||
/** Files read into the index since the current backfill began, and its total. */
|
||||
filesIndexed: number
|
||||
filesTotal: number | null
|
||||
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. */
|
||||
droppedPending: number
|
||||
/** Unfinished writes the open tombstoned; a non-zero value means a crash. */
|
||||
recoveredRows: number
|
||||
failures: number
|
||||
degradedRoots: SessionSearchDegradedRoot[]
|
||||
lastReconcileAt: number | null
|
||||
discovered: Record<string, SessionSearchDiscoveredCount>
|
||||
}
|
||||
|
||||
/** Observes the backfill and the reconciler; owns no work and no timers. */
|
||||
export class SessionSearchIndexingStatus {
|
||||
private working: 'discovering' | 'indexing' | null = null
|
||||
private paused = false
|
||||
private filesIndexed = 0
|
||||
private filesTotal: number | null = null
|
||||
private bytesIndexed = 0
|
||||
private filesPending = 0
|
||||
private droppedPending = 0
|
||||
private recoveredRows = 0
|
||||
private failures = 0
|
||||
private degradedRoots: SessionSearchDegradedRoot[] = []
|
||||
private lastReconcileAt: number | null = null
|
||||
private discovered = new Map<AiVaultAgent, SessionSearchDiscoveredCount>()
|
||||
|
||||
snapshot(): SessionSearchIndexStatus {
|
||||
return {
|
||||
phase: this.phase(),
|
||||
filesIndexed: this.filesIndexed,
|
||||
filesTotal: this.filesTotal,
|
||||
bytesIndexed: this.bytesIndexed,
|
||||
filesPending: this.filesPending,
|
||||
droppedPending: this.droppedPending,
|
||||
recoveredRows: this.recoveredRows,
|
||||
failures: this.failures,
|
||||
degradedRoots: this.degradedRoots.map((root) => ({ ...root })),
|
||||
lastReconcileAt: this.lastReconcileAt,
|
||||
discovered: Object.fromEntries(this.discovered)
|
||||
}
|
||||
}
|
||||
|
||||
private phase(): SessionSearchIndexPhase {
|
||||
if (this.paused) {
|
||||
return 'paused'
|
||||
}
|
||||
// Why work outranks degradation: a run in progress is the more useful thing
|
||||
// to show, and the degraded roots are still in the snapshot either way.
|
||||
if (this.working) {
|
||||
return this.working
|
||||
}
|
||||
return this.degradedRoots.length > 0 ? 'degraded' : 'current'
|
||||
}
|
||||
|
||||
setPaused(paused: boolean): void {
|
||||
this.paused = paused
|
||||
}
|
||||
|
||||
setRecoveredRows(rows: number): void {
|
||||
this.recoveredRows = rows
|
||||
}
|
||||
|
||||
/** A full sweep restarts the progress pair; a reconcile cycle only reports work. */
|
||||
beginSweep(): void {
|
||||
this.working = 'discovering'
|
||||
this.filesIndexed = 0
|
||||
this.filesTotal = null
|
||||
this.bytesIndexed = 0
|
||||
this.failures = 0
|
||||
}
|
||||
|
||||
beginCycle(): void {
|
||||
this.working ??= 'indexing'
|
||||
}
|
||||
|
||||
setDiscovered(counts: Map<AiVaultAgent, SessionSearchDiscoveredCount>): void {
|
||||
this.discovered = counts
|
||||
}
|
||||
|
||||
/** Total is what the sweep will actually read, after retention filtering. */
|
||||
planned(total: number, failures: number): void {
|
||||
this.working = 'indexing'
|
||||
this.filesTotal = total
|
||||
this.failures += failures
|
||||
}
|
||||
|
||||
indexed(bytes: number): void {
|
||||
this.filesIndexed += 1
|
||||
this.bytesIndexed += bytes
|
||||
}
|
||||
|
||||
failed(): void {
|
||||
this.failures += 1
|
||||
}
|
||||
|
||||
setPending(pending: number, dropped: number): void {
|
||||
this.filesPending = pending
|
||||
this.droppedPending = dropped
|
||||
}
|
||||
|
||||
setDegradedRoots(roots: SessionSearchDegradedRoot[]): void {
|
||||
this.degradedRoots = roots
|
||||
}
|
||||
|
||||
finishWork(atMs: number): void {
|
||||
this.working = null
|
||||
this.lastReconcileAt = atMs
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,76 @@
|
||||
import { chmod, mkdir, mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, expect, it } from 'vitest'
|
||||
import type { SessionFileDiscovery } from '../ai-vault/session-scanner-types'
|
||||
import { degradedSessionSearchRoots } from './session-search-root-health'
|
||||
|
||||
const CAN_DENY_READ = process.platform !== 'win32' && process.getuid?.() !== 0
|
||||
|
||||
let roots: string[] = []
|
||||
|
||||
afterEach(async () => {
|
||||
await Promise.all(roots.map((root) => rm(root, { recursive: true, force: true })))
|
||||
roots = []
|
||||
})
|
||||
|
||||
async function tempRoot(): Promise<string> {
|
||||
const root = await mkdtemp(join(tmpdir(), 'ss-root-health-'))
|
||||
roots.push(root)
|
||||
return root
|
||||
}
|
||||
|
||||
function discovery(rootDir: string, files: number): SessionFileDiscovery {
|
||||
return {
|
||||
agent: 'claude',
|
||||
rootDir,
|
||||
files: Array.from({ length: files }, (_unused, index) => ({
|
||||
path: join(rootDir, `${index}.jsonl`),
|
||||
mtimeMs: 0,
|
||||
modifiedAt: new Date(0).toISOString()
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
it('leaves an agent that is simply not installed alone', async () => {
|
||||
const root = await tempRoot()
|
||||
expect(await degradedSessionSearchRoots([discovery(join(root, 'never-created'), 0)], [])).toEqual(
|
||||
[]
|
||||
)
|
||||
})
|
||||
|
||||
it('never probes a root that returned files', async () => {
|
||||
// The path does not exist, so a probe would report it degraded; a root that
|
||||
// yielded transcripts is readable by construction and must not be re-checked.
|
||||
expect(await degradedSessionSearchRoots([discovery('/definitely/not/here', 3)], [])).toEqual([])
|
||||
})
|
||||
|
||||
it.skipIf(!CAN_DENY_READ)('names a root that exists but cannot be read', async () => {
|
||||
const root = await tempRoot()
|
||||
const blocked = join(root, 'blocked')
|
||||
await mkdir(blocked)
|
||||
await chmod(blocked, 0o000)
|
||||
try {
|
||||
const degraded = await degradedSessionSearchRoots([discovery(blocked, 0)], [])
|
||||
expect(degraded).toHaveLength(1)
|
||||
expect(degraded[0]?.root).toBe(blocked)
|
||||
expect(degraded[0]?.reason).toContain('EACCES')
|
||||
} finally {
|
||||
await chmod(blocked, 0o755)
|
||||
}
|
||||
})
|
||||
|
||||
it('carries a root-level scan issue through, but not a per-file notice', async () => {
|
||||
const root = await tempRoot()
|
||||
await mkdir(join(root, 'healthy'))
|
||||
const rootDir = join(root, 'healthy')
|
||||
const degraded = await degradedSessionSearchRoots(
|
||||
[discovery(rootDir, 2)],
|
||||
[
|
||||
{ agent: 'claude', path: rootDir, message: 'The distro stopped responding.' },
|
||||
{ agent: 'claude', path: rootDir, kind: 'notice', message: 'issue list truncated' },
|
||||
{ agent: 'claude', path: join(rootDir, 'one.jsonl'), message: 'a single unreadable file' }
|
||||
]
|
||||
)
|
||||
expect(degraded).toEqual([{ root: rootDir, reason: 'The distro stopped responding.' }])
|
||||
})
|
||||
@@ -0,0 +1,62 @@
|
||||
import type { AiVaultScanIssue } from '../../shared/ai-vault-types'
|
||||
import type { SessionFileDiscovery } from '../ai-vault/session-scanner-types'
|
||||
import { wslGatedReaddir } from '../native-chat/wsl-transcript-fs-access'
|
||||
|
||||
/** A scan root the index could not read, and what stopped it. */
|
||||
export type SessionSearchDegradedRoot = { root: string; reason: string }
|
||||
|
||||
// Why the indexer probes at all: the walker swallows a readdir failure and
|
||||
// returns, so an EACCES root and an agent that was never installed both arrive
|
||||
// as "no files". Reporting the first as an empty index would be the
|
||||
// loss-of-contact-as-absence mistake docs/reference/ssh-execution-boundary.md
|
||||
// forbids, so an empty root is re-checked and only ENOENT counts as absent.
|
||||
const ABSENT_ROOT = new Set(['ENOENT', 'ENOTDIR'])
|
||||
|
||||
/**
|
||||
* 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
|
||||
* the per-cycle cost at one readdir per genuinely empty tree.
|
||||
*/
|
||||
export async function degradedSessionSearchRoots(
|
||||
discoveries: readonly SessionFileDiscovery[],
|
||||
issues: readonly AiVaultScanIssue[],
|
||||
signal?: AbortSignal
|
||||
): Promise<SessionSearchDegradedRoot[]> {
|
||||
const degraded = new Map<string, string>()
|
||||
for (const issue of issues) {
|
||||
if (issue.kind !== 'notice' && discoveries.some((one) => one.rootDir === issue.path)) {
|
||||
degraded.set(issue.path, issue.message)
|
||||
}
|
||||
}
|
||||
const empty = [
|
||||
...new Set(
|
||||
discoveries.filter((one) => one.files.length === 0).map((discovery) => discovery.rootDir)
|
||||
)
|
||||
]
|
||||
for (const root of empty) {
|
||||
if (degraded.has(root) || signal?.aborted) {
|
||||
continue
|
||||
}
|
||||
const reason = await unreadableRootReason(root, signal)
|
||||
if (reason) {
|
||||
degraded.set(root, reason)
|
||||
}
|
||||
}
|
||||
return [...degraded].map(([root, reason]) => ({ root, reason }))
|
||||
}
|
||||
|
||||
async function unreadableRootReason(root: string, signal?: AbortSignal): Promise<string | null> {
|
||||
try {
|
||||
await wslGatedReaddir(root, 'scan', signal)
|
||||
return null
|
||||
} catch (error) {
|
||||
const code =
|
||||
error && typeof error === 'object' && 'code' in error && typeof error.code === 'string'
|
||||
? error.code
|
||||
: null
|
||||
if (code !== null && ABSENT_ROOT.has(code)) {
|
||||
return null
|
||||
}
|
||||
return error instanceof Error ? error.message : String(error)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
import type { AiVaultScanIssue } from '../../shared/ai-vault-types'
|
||||
import { sessionCandidatesFromDiscoveries } from '../ai-vault/session-scanner-candidates'
|
||||
import { discoverAiVaultSessionSources } from '../ai-vault/session-scanner-source-discovery'
|
||||
import type {
|
||||
AiVaultScanOptions,
|
||||
SessionFileCandidate,
|
||||
SessionFileDiscovery
|
||||
} from '../ai-vault/session-scanner-types'
|
||||
|
||||
/**
|
||||
* Where the indexer looks. The caller resolves these so the index enumerates
|
||||
* exactly the trees the session list does; the indexer owns the bounds
|
||||
* (`limit`, `limitPerAgent`, `unlimited`) and its own cancellation, so those
|
||||
* are not the caller's to set.
|
||||
*/
|
||||
export type SessionSearchScanRoots = Omit<
|
||||
AiVaultScanOptions,
|
||||
'signal' | 'limit' | 'unlimited' | 'limitPerAgent' | 'scopePaths'
|
||||
>
|
||||
|
||||
export type SessionSearchDiscovery = {
|
||||
/** Newest first, Codex hardlink aliases collapsed, exactly as a list scan sees them. */
|
||||
candidates: SessionFileCandidate[]
|
||||
discoveries: SessionFileDiscovery[]
|
||||
issues: AiVaultScanIssue[]
|
||||
}
|
||||
|
||||
/**
|
||||
* The discovery half of a list scan, without the parse. `limitPerAgent` is the
|
||||
* sidebar's own recency rule (`SessionNewestFiles` keeps the newest N per root);
|
||||
* passing Infinity is what makes a sweep whole.
|
||||
*/
|
||||
export async function discoverSessionSearchCandidates(
|
||||
roots: SessionSearchScanRoots,
|
||||
args: { limitPerAgent: number; signal?: AbortSignal }
|
||||
): Promise<SessionSearchDiscovery> {
|
||||
const issues: AiVaultScanIssue[] = []
|
||||
const options: AiVaultScanOptions = { ...roots, signal: args.signal }
|
||||
const discoveries = await discoverAiVaultSessionSources({
|
||||
options,
|
||||
limitPerAgent: args.limitPerAgent,
|
||||
issues
|
||||
})
|
||||
const candidates = await sessionCandidatesFromDiscoveries(discoveries, options)
|
||||
return { candidates, discoveries, issues }
|
||||
}
|
||||
Reference in New Issue
Block a user