feat(ai-vault): full-text agent session search index, RPC, and CLI

Phase 0: parse-cache resume guard checks inode; tool_use blocks render in previews.
Phase 1: FTS5 index (node:sqlite) fed by the existing parse cache under an
AsyncLocalStorage capture scope; identifier shadow column, 3 KB tool cap,
fts5vocab typo repair, phrase-first, length prior; background backfill in the
scanner process; aiVault.searchSessions / searchCoverage RPC; orca search CLI.
This commit is contained in:
Jinwoo-H
2026-09-06 04:48:03 -04:00
parent 337433b39a
commit ab7ec480f3
49 changed files with 2781 additions and 207 deletions
+75
View File
@@ -0,0 +1,75 @@
import { basename } from 'node:path'
import { aiVaultAgentLabel } from '../shared/ai-vault-types'
import type { AiVaultSearchHit, AiVaultSearchResult } from '../shared/ai-vault-search-types'
const ROLE_LABEL: Record<AiVaultSearchHit['evidence']['role'], string> = {
user: 'you ',
assistant: 'agent',
tool: 'tool ',
system: 'sys ',
unknown: ' '
}
function relativeAge(iso: string | null, now = Date.now()): string {
if (!iso) {
return 'unknown time'
}
const ms = now - Date.parse(iso)
if (!Number.isFinite(ms) || ms < 0) {
return 'just now'
}
const minutes = Math.round(ms / 60_000)
if (minutes < 60) {
return `${Math.max(1, minutes)} min ago`
}
const hours = Math.round(minutes / 60)
if (hours < 48) {
return `${hours} h ago`
}
const days = Math.round(hours / 24)
if (days < 14) {
return `${days} d ago`
}
const weeks = Math.round(days / 7)
return weeks < 9 ? `${weeks} wk ago` : `${Math.round(days / 30)} mo ago`
}
function projectLabel(hit: AiVaultSearchHit): string {
const cwd = hit.cwd ? basename(hit.cwd) : '—'
return hit.branch ? `${cwd} · ${hit.branch}` : cwd
}
function formatHit(index: number, hit: AiVaultSearchHit): string {
const header = `${String(index + 1).padStart(2)}. ${hit.title}`
const meta = `${aiVaultAgentLabel(hit.agent)} · ${projectLabel(hit)} · ${relativeAge(hit.updatedAt)}`
const evidence = hit.evidence.snippet
? ` ${ROLE_LABEL[hit.evidence.role]} ▸ ${hit.evidence.snippet.replaceAll('\n', ' ')}`
: null
const resume = ` resume: ${hit.resumeCommand}${hit.cwd ? ` (cwd ${hit.cwd})` : ''}`
return [`${header} ${meta}`, evidence, resume].filter(Boolean).join('\n')
}
export function formatAgentSessionSearch(
result: AiVaultSearchResult,
context: { query: string; cwd: string }
): string {
const lines: string[] = []
if (result.hits.length === 0) {
lines.push(`No sessions match "${context.query}".`)
} else {
lines.push(...result.hits.map((hit, index) => formatHit(index, hit)), '')
}
if (result.repairedTerms) {
lines.push(`Searched for: ${result.repairedTerms.join(' ')}`)
}
const { coverage } = result
const scope = `${coverage.sessionsIndexed.toLocaleString()} sessions indexed`
const pending =
coverage.backfill === 'running'
? ', still indexing older sessions'
: coverage.filesPending > 0
? `, ${coverage.filesPending} changed files pending`
: ''
lines.push(`${scope}${pending} · ${result.durationMs.toFixed(0)} ms`)
return lines.join('\n')
}
+12
View File
@@ -0,0 +1,12 @@
const SEARCH_FLAG_HELP: Record<string, string> = {
'agent-session': '--agent-session <query> Text to find in any coding-agent session (required)',
agent: '--agent <agent> Only this agent (claude, codex, cursor, …); repeatable',
path: '--path <dir> Only sessions whose working directory is under this path; repeatable',
since: '--since <iso> Only sessions updated at or after this ISO 8601 timestamp',
newest: '--newest Sort by session time instead of relevance',
host: '--host <host> Search a paired runtime host (runtime:<environment>)'
}
export function formatSearchFlagHelp(command: string, flag: string): string | null {
return command === 'search' ? (SEARCH_FLAG_HELP[flag] ?? null) : null
}
+5
View File
@@ -251,5 +251,10 @@ export const HANDLER_GROUPS: readonly HandlerGroup[] = [
name: 'skills',
keys: ['skills list', 'skills get', 'skills install', 'skills update'],
load: async () => (await import('./handlers/skills.js')).SKILL_HANDLERS
},
{
name: 'search',
keys: ['search'],
load: async () => (await import('./handlers/search.js')).SEARCH_HANDLERS
}
]
+77
View File
@@ -0,0 +1,77 @@
import type { CommandHandler } from '../dispatch'
import { printResult } from '../format'
import {
getOptionalPositiveIntegerFlag,
getOptionalStringFlag,
getRepeatedStringFlag
} from '../flags'
import { parseHostFlag } from '../execution-host-flag'
import { RuntimeClientError } from '../runtime/types'
import { AI_VAULT_AGENTS, type AiVaultAgent } from '../../shared/ai-vault-types'
import type { AiVaultSearchHit, AiVaultSearchResult } from '../../shared/ai-vault-search-types'
import { formatAgentSessionSearch } from '../agent-session-search-format'
function parseAgents(flags: Map<string, string | boolean>): AiVaultAgent[] | undefined {
const values = getRepeatedStringFlag(flags, 'agent')
if (values.length === 0) {
return undefined
}
const agents: AiVaultAgent[] = []
for (const value of values) {
const lowered = value.toLowerCase()
if (!(AI_VAULT_AGENTS as readonly string[]).includes(lowered)) {
throw new RuntimeClientError(
'invalid_argument',
`Unknown --agent ${value}. Expected one of: ${AI_VAULT_AGENTS.join(', ')}.`
)
}
agents.push(lowered as AiVaultAgent)
}
return agents
}
function parseSince(flags: Map<string, string | boolean>): string | undefined {
const value = getOptionalStringFlag(flags, 'since')
if (value === undefined) {
return undefined
}
const parsed = Date.parse(value)
if (!Number.isFinite(parsed)) {
throw new RuntimeClientError('invalid_argument', '--since must be an ISO 8601 timestamp.')
}
return new Date(parsed).toISOString()
}
export const SEARCH_HANDLERS: Record<string, CommandHandler> = {
search: async ({ client, flags, json, cwd }) => {
const query = getOptionalStringFlag(flags, 'agent-session')
if (!query) {
throw new RuntimeClientError(
'invalid_argument',
'Missing --agent-session <query>. Example: orca search --agent-session "strict mode violation"'
)
}
const host = parseHostFlag(flags)
if (host?.kind === 'ssh') {
throw new RuntimeClientError(
'invalid_argument',
'Agent session search runs on a runtime host. Use --host runtime:<environment> or omit --host.'
)
}
const scopePaths = getRepeatedStringFlag(flags, 'path').map((value) =>
value.startsWith('~') ? value.replace(/^~/, process.env.HOME ?? '~') : value
)
const result = await client.call<AiVaultSearchResult>('aiVault.searchSessions', {
query,
limit: getOptionalPositiveIntegerFlag(flags, 'limit'),
agents: parseAgents(flags),
scopePaths: scopePaths.length > 0 ? scopePaths : undefined,
since: parseSince(flags),
sort: flags.get('newest') === true ? 'newest' : 'relevance',
...(host?.kind === 'runtime' ? { executionHostId: host.id } : {})
})
printResult(result, json, (value) => formatAgentSessionSearch(value, { query, cwd }))
}
}
export type { AiVaultSearchHit }
+5
View File
@@ -1,3 +1,4 @@
import { formatSearchFlagHelp } from './agent-session-search-help'
import type { CommandSpec } from './args'
import { findCommandSpec, isCommandGroup, supportsBrowserPageFlag } from './args'
import { unknownCommandData } from './command-suggestion'
@@ -72,6 +73,10 @@ export function formatGroupHelp(specs: CommandSpec[], group: string): string {
function formatCommandFlagHelp(flag: string, commandPath: string[]): string {
const command = commandPath.join(' ')
const searchHelp = formatSearchFlagHelp(command, flag)
if (searchHelp) {
return searchHelp
}
if (command === 'skills install' && flag === 'agent') {
return '--agent <names> Comma-separated install targets; default is detected agents'
}
+1
View File
@@ -42,6 +42,7 @@ export const ROOT_HELP_TEXT_SECONDARY = [
' orca agent-context [--json]',
' orca account add [--agent claude|codex] [--json]',
' orca account list [--json]',
' orca search --agent-session "<query>" [--limit <n>] [--agent <agent>] [--newest] [--json]',
' orca host list [--json]',
' orca environment add --name <name> --pairing-code <code> [--json]',
' orca environment list [--json]',
+3 -1
View File
@@ -17,6 +17,7 @@ import { LINEAR_COMMAND_SPECS } from './linear'
import { VM_COMMAND_SPECS } from './vm'
import { SKILL_COMMAND_SPECS } from './skills'
import { ARTIFACT_COMMAND_SPECS } from './artifacts'
import { SEARCH_COMMAND_SPECS } from './search'
export const COMMAND_SPECS: CommandSpec[] = [
...CORE_COMMAND_SPECS,
@@ -36,5 +37,6 @@ export const COMMAND_SPECS: CommandSpec[] = [
...LINEAR_COMMAND_SPECS,
...VM_COMMAND_SPECS,
...EMULATOR_COMMAND_SPECS,
...SKILL_COMMAND_SPECS
...SKILL_COMMAND_SPECS,
...SEARCH_COMMAND_SPECS
]
+32
View File
@@ -0,0 +1,32 @@
import { GLOBAL_FLAGS, type CommandSpec } from '../args'
export const SEARCH_COMMAND_SPECS: CommandSpec[] = [
{
path: ['search'],
summary: 'Search the full text of every local coding-agent session',
usage:
'orca search --agent-session "<query>" [--limit <n>] [--agent <agent>] [--path <dir>] [--since <iso>] [--newest] [--host <host>] [--json]',
allowedFlags: [
...GLOBAL_FLAGS,
'agent-session',
'limit',
'agent',
'path',
'since',
'newest',
'host'
],
notes: [
'Searches what you typed, what the agent said, the commands it ran, and the first 3 KB of each tool output across Claude Code, Codex, Cursor, Gemini, OpenCode, and the other agents Orca scans.',
'Quote paths, identifiers, or error text to match them exactly; plain words match anywhere. A misspelled word is repaired from the index vocabulary when nothing matches.',
'--agent-session takes the query. --agent and --path may repeat. --since takes an ISO timestamp. --newest sorts by session time instead of relevance.',
'The index builds in the background on first use; the result reports how many sessions are covered so far.',
'--host runtime:<environment> searches that host; the index always lives with the transcripts.'
],
examples: [
'orca search --agent-session "strict mode violation getByRole"',
'orca search --agent-session resolveTerminalPath --agent claude --newest',
'orca search --agent-session "kernel panic" --path ~/orca --since 2026-08-01T00:00:00Z --json'
]
}
]
@@ -0,0 +1,54 @@
// Identifier shadow terms: `resolveTerminalPath` → `resolve terminal path`,
// `src/main/foo-bar.ts` → `src main foo bar ts`. Stored in a separate FTS5
// column so a partial identifier still matches; the largest single accuracy
// win measured in the retrieval shoot-out (MRR 0.50 → 0.55).
const RAW_TOKEN = /[A-Za-z0-9_./-]+/g
const CAMEL_PIECE = /[A-Z]+(?![a-z])|[A-Z][a-z0-9]*|[a-z0-9]+/g
const SEPARATOR = /[_./-]+/
// Worth shadowing: has a separator, a camel boundary, or is SCREAMING_CASE.
const INTERESTING = /[_./-]|[a-z0-9][A-Z]|^[A-Z]{2,}[0-9_]*$/
const MIN_TOKEN = 3
const MAX_TOKEN = 120
const MIN_PIECE = 2
function hasMixedCase(piece: string): boolean {
return /[a-z]/.test(piece) && /[A-Z]/.test(piece)
}
export function identifierShadowTerms(text: string, limit = 4000): string[] {
const out: string[] = []
const seen = new Set<string>()
for (const match of text.matchAll(RAW_TOKEN)) {
const token = match[0]
if (token.length < MIN_TOKEN || token.length > MAX_TOKEN || !INTERESTING.test(token)) {
continue
}
const parts: string[] = []
for (const piece of token.split(SEPARATOR)) {
if (!piece) {
continue
}
parts.push(piece)
if (hasMixedCase(piece)) {
parts.push(...(piece.match(CAMEL_PIECE) ?? []))
}
}
for (const part of parts) {
const lowered = part.toLowerCase()
if (lowered.length < MIN_PIECE || seen.has(lowered)) {
continue
}
seen.add(lowered)
out.push(lowered)
if (out.length >= limit) {
return out
}
}
}
return out
}
export function identifierShadowText(text: string, limit?: number): string {
return identifierShadowTerms(text, limit).join(' ')
}
@@ -0,0 +1,214 @@
import type SyncDatabase from '../sqlite/sync-database'
import type {
SessionSearchFileIdentity,
SessionSearchIndexedFile,
SessionSearchIndexUpdate
} from '../ai-vault/session-search-capture'
import { identifierShadowText } from './session-search-identifier-split'
// Why: FTS5's length normalization buries a 100 KB tool log even when it holds
// the query many times; chunks at line boundaries keep every row rankable.
const CHUNK_TARGET_CHARS = 8000
type FileRow = {
dev: number | null
ino: number | null
byte_offset: number
mtime_ms: number
size_bytes: number | null
session_row_id: number | null
}
export function chunkMessageText(text: string): string[] {
if (text.length <= CHUNK_TARGET_CHARS) {
return [text]
}
const chunks: string[] = []
let start = 0
while (start < text.length) {
let end = Math.min(text.length, start + CHUNK_TARGET_CHARS)
if (end < text.length) {
const newline = text.lastIndexOf('\n', end)
if (newline > start + CHUNK_TARGET_CHARS / 2) {
end = newline + 1
}
}
chunks.push(text.slice(start, end))
start = end
}
return chunks
}
export class SessionSearchIndexWriter {
constructor(private readonly db: SyncDatabase) {}
indexedFile(path: string, identity: SessionSearchFileIdentity): SessionSearchIndexedFile | null {
const row = this.db
.prepare(
'SELECT dev, ino, byte_offset, mtime_ms, size_bytes, session_row_id FROM files WHERE path = ?'
)
.get(path) as FileRow | undefined
if (!row) {
return null
}
if (identity && row.dev !== null && row.ino !== null) {
if (row.dev !== identity.dev || row.ino !== identity.ino) {
return null
}
}
return { byteOffset: row.byte_offset, mtimeMs: row.mtime_ms, sizeBytes: row.size_bytes }
}
apply(update: SessionSearchIndexUpdate): void {
const path = update.candidate.file.path
this.db.exec('BEGIN IMMEDIATE')
try {
const existing = this.db
.prepare('SELECT byte_offset, session_row_id FROM files WHERE path = ?')
.get(path) as Pick<FileRow, 'byte_offset' | 'session_row_id'> | undefined
const appendable =
update.mode === 'append' &&
existing !== undefined &&
existing.byte_offset === update.previousByteOffset &&
existing.session_row_id !== null
if (!appendable) {
this.deleteFile(path, existing?.session_row_id ?? null)
}
// Why: an append that does not continue from the stored offset (a racing
// parse advanced it) would leave a hole; drop the file so the next parse
// is whole instead of storing a partial session.
if (update.mode === 'append' && !appendable) {
this.db.exec('COMMIT')
return
}
if (update.session === null) {
this.upsertFile(update, null)
this.db.exec('COMMIT')
return
}
const sessionRowId = this.upsertSession(update, appendable ? existing.session_row_id : null)
this.insertMessages(sessionRowId, update)
this.upsertFile(update, sessionRowId)
this.db.exec('COMMIT')
} catch (error) {
this.db.exec('ROLLBACK')
throw error
}
}
removeFile(path: string): void {
const existing = this.db
.prepare('SELECT session_row_id FROM files WHERE path = ?')
.get(path) as Pick<FileRow, 'session_row_id'> | undefined
if (!existing) {
return
}
this.db.exec('BEGIN IMMEDIATE')
try {
this.deleteFile(path, existing.session_row_id)
this.db.exec('COMMIT')
} catch (error) {
this.db.exec('ROLLBACK')
throw error
}
}
private deleteFile(path: string, sessionRowId: number | null): void {
if (sessionRowId !== null) {
const ids = this.db
.prepare('SELECT id FROM messages WHERE session_row_id = ?')
.all(sessionRowId) as { id: number }[]
const deleteFts = this.db.prepare('DELETE FROM messages_fts WHERE rowid = ?')
const deleteConversation = this.db.prepare('DELETE FROM conversation_fts WHERE rowid = ?')
for (const { id } of ids) {
deleteFts.run(id)
deleteConversation.run(id)
}
this.db.prepare('DELETE FROM messages WHERE session_row_id = ?').run(sessionRowId)
this.db.prepare('DELETE FROM sessions WHERE id = ?').run(sessionRowId)
}
this.db.prepare('DELETE FROM files WHERE path = ?').run(path)
}
private upsertSession(update: SessionSearchIndexUpdate, rowId: number | null): number {
const session = update.session!
const values = [
session.agent,
session.sessionId,
session.filePath,
session.codexHome,
session.title,
session.cwd,
session.branch,
session.createdAt,
session.updatedAt,
session.messageCount,
session.resumeCommand
]
if (rowId !== null) {
this.db
.prepare(
`UPDATE sessions SET agent = ?, session_id = ?, file_path = ?, codex_home = ?, title = ?,
cwd = ?, branch = ?, created_at = ?, updated_at = ?, message_count = ?, resume_command = ?
WHERE id = ?`
)
.run(...values, rowId)
return rowId
}
const result = this.db
.prepare(
`INSERT INTO sessions(agent, session_id, file_path, codex_home, title, cwd, branch,
created_at, updated_at, message_count, resume_command)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
)
.run(...values)
return Number(result.lastInsertRowid)
}
private insertMessages(sessionRowId: number, update: SessionSearchIndexUpdate): void {
const insertMessage = this.db.prepare(
'INSERT INTO messages(session_row_id, role, ts) VALUES (?, ?, ?)'
)
const insertFts = this.db.prepare(
'INSERT INTO messages_fts(rowid, user_text, assistant_text, tool_text, identifiers) VALUES (?, ?, ?, ?, ?)'
)
const insertConversation = this.db.prepare(
'INSERT INTO conversation_fts(rowid, user_text, assistant_text) VALUES (?, ?, ?)'
)
for (const message of update.messages) {
for (const chunk of chunkMessageText(message.text)) {
const id = Number(
insertMessage.run(sessionRowId, message.role, message.timestamp).lastInsertRowid
)
const user = message.role === 'user' ? chunk : ''
const assistant = message.role === 'assistant' ? chunk : ''
const tool = message.role === 'tool' ? chunk : ''
insertFts.run(id, user, assistant, tool, identifierShadowText(chunk))
if (message.role !== 'tool') {
insertConversation.run(id, user, assistant)
}
}
}
}
private upsertFile(update: SessionSearchIndexUpdate, sessionRowId: number | null): void {
const { file } = update.candidate
this.db
.prepare(
`INSERT INTO files(path, dev, ino, byte_offset, mtime_ms, size_bytes, session_row_id)
VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(path) DO UPDATE SET dev = excluded.dev, ino = excluded.ino,
byte_offset = excluded.byte_offset, mtime_ms = excluded.mtime_ms,
size_bytes = excluded.size_bytes, session_row_id = excluded.session_row_id`
)
.run(
file.path,
file.dev ?? null,
file.ino ?? null,
update.byteOffset,
file.mtimeMs,
file.sizeBytes ?? null,
sessionRowId
)
}
}
@@ -0,0 +1,18 @@
import { join } from 'node:path'
import type { AiVaultSessionSearchInit } from '../ai-vault/session-scanner-service-protocol'
// Why: like the parse cache, the index path is captured once at the composition
// root from the canonical userData dir; every export is a no-op until then.
let init: AiVaultSessionSearchInit | null = null
export function initSessionSearchPaths(userDataPath: string): void {
init = { databasePath: join(userDataPath, 'ai-vault-search', 'index.sqlite') }
}
export function getSessionSearchInitOptions(): AiVaultSessionSearchInit | null {
return init ? { ...init } : null
}
export function resetSessionSearchPathsForTests(): void {
init = null
}
@@ -0,0 +1,86 @@
import { identifierShadowTerms } from './session-search-identifier-split'
// Tokens exactly as the unicode61 tokenizer with `_ . - /` tokenchars emits them.
const INDEX_TOKEN = /[A-Za-z0-9_./-]+/g
const STOP_WORDS = new Set(
(
'a an and are as at be but by for from how i if in into is it its of on or that the this to ' +
'was were what when where which who why with you your we my me do does did not no can could ' +
'should would about our us they them there their has have had been being so such then than ' +
"these those there's im ive dont"
).split(' ')
)
const MAX_BODY_TERMS = 48
const MAX_TERMS = 64
// A query that quotes something from a transcript: camelCase, SCREAMING_SNAKE,
// a dotted or snake_case name, a path, a filename, a PR number, a ticket, code
// punctuation, or an error word.
const LITERAL_SHAPE =
/[A-Za-z0-9_]*[a-z][A-Z][A-Za-z0-9_]*|\b[A-Z][A-Z0-9]{2,}(_[A-Z0-9]+)+\b|\b\w{2,}[._]\w{2,}\b|\b[\w.-]+\/[\w/.-]+\b|\b\w+\.(ts|tsx|js|jsx|py|rs|go|json|md|sh|yml|yaml|toml|c|cc|h|java|sql)\b|#\d{3,}|\b[A-Z]{2,6}-\d{2,}\b|[(){};=]|::|->|--\w|\b(Error|Exception|Traceback|error:|warning:)\b/
const QUOTED = /"[^"]{3,}"|'[^']{3,}'/
export type SessionSearchQueryPlan = {
literal: boolean
/** Deduplicated index-faithful terms for the OR fallback, incl. identifier pieces. */
terms: string[]
/** Query-order tokens minus stop words: the phrase / AND candidate. */
body: string[]
}
export function isLiteralQuery(query: string): boolean {
return QUOTED.test(query) || LITERAL_SHAPE.test(query)
}
function indexTokens(query: string): string[] {
const out: string[] = []
for (const match of query.matchAll(INDEX_TOKEN)) {
const token = match[0]
if (token.length > 1 && /[A-Za-z0-9]/.test(token)) {
out.push(token)
if (out.length >= MAX_BODY_TERMS) {
break
}
}
}
return out
}
export function planSessionSearchQuery(query: string): SessionSearchQueryPlan {
const raw = indexTokens(query)
let body = raw.filter((token) => !STOP_WORDS.has(token.toLowerCase()))
if (body.length < 2) {
body = raw
}
const terms = [...new Set(body)]
const extra: string[] = []
for (const term of terms) {
for (const piece of identifierShadowTerms(term, 12)) {
if (!terms.includes(piece) && !STOP_WORDS.has(piece) && !extra.includes(piece)) {
extra.push(piece)
}
}
}
return {
literal: isLiteralQuery(query),
terms: [...terms, ...extra].slice(0, MAX_TERMS),
body: body.slice(0, MAX_BODY_TERMS)
}
}
// Why: `cli.mjs`, `foo-bar`, and `C++` are all FTS5 syntax errors unquoted.
export function quoteFtsTerm(term: string): string {
return `"${term.replaceAll('"', '""')}"`
}
export function phraseExpression(terms: readonly string[]): string {
return terms.map(quoteFtsTerm).join(' ')
}
export function andExpression(terms: readonly string[]): string {
return terms.map(quoteFtsTerm).join(' AND ')
}
export function orExpression(terms: readonly string[]): string {
return terms.map(quoteFtsTerm).join(' OR ')
}
@@ -0,0 +1,236 @@
import type SyncDatabase from '../sqlite/sync-database'
import type {
AiVaultSearchArgs,
AiVaultSearchHit,
AiVaultSearchRoute
} from '../../shared/ai-vault-search-types'
import {
AI_VAULT_SEARCH_LIMIT_DEFAULT,
AI_VAULT_SEARCH_LIMIT_MAX
} from '../../shared/ai-vault-search-types'
import type { AiVaultAgent } from '../../shared/ai-vault-types'
import {
andExpression,
orExpression,
phraseExpression,
planSessionSearchQuery,
type SessionSearchQueryPlan
} from './session-search-query-planner'
import { SessionSearchTypoRepair } from './session-search-typo-repair'
// Measured: user 3 / assistant 2 / tool 1 / identifiers 1 (MRR 0.503 vs 0.475 flat).
const FULL_WEIGHTS = '3.0, 2.0, 1.0, 1.0'
const CONVERSATION_WEIGHTS = '3.0, 2.0'
// Candidate messages fetched before rolling up to sessions; more does not help.
const MESSAGE_CANDIDATE_LIMIT = 600
// Subtracted per session: `0.02 · ln(1 + messages)`; slightly positive on both eval sets.
const LENGTH_PRIOR = 0.02
const SNIPPET_TOKENS = 12
type MessageRow = {
rowid: number
score: number
session_row_id: number
role: string
ts: string | null
}
type SessionRow = {
id: number
agent: AiVaultAgent
session_id: string
file_path: string
codex_home: string | null
title: string
cwd: string | null
branch: string | null
updated_at: string | null
message_count: number
resume_command: string
}
export type SessionSearchExecution = {
hits: AiVaultSearchHit[]
route: AiVaultSearchRoute
repairedTerms?: string[]
}
export class SessionSearchQuery {
private readonly typoRepair: SessionSearchTypoRepair
constructor(private readonly db: SyncDatabase) {
this.typoRepair = new SessionSearchTypoRepair(db)
}
execute(args: AiVaultSearchArgs): SessionSearchExecution {
const plan = planSessionSearchQuery(args.query)
if (plan.terms.length === 0) {
return { hits: [], route: 'or' }
}
const tier = args.tier ?? 'full'
const exact = this.retrieveLiteral(plan, tier)
if (exact) {
return { hits: this.rollUp(exact.rows, args, tier), route: exact.route }
}
// Why: repair runs before the OR fallback, not after it fails; a typo next
// to a common word would otherwise be masked by the common word's hits.
const repaired = this.repair(plan)
const effective = repaired ?? plan
const literal = repaired ? this.retrieveLiteral(repaired, tier) : null
const result = literal ?? {
rows: this.match(orExpression(effective.terms), tier),
route: 'or' as const
}
return {
hits: this.rollUp(result.rows, args, tier),
route: repaired ? (`typo+${result.route}` as AiVaultSearchRoute) : result.route,
...(repaired ? { repairedTerms: repaired.body } : {})
}
}
private repair(plan: SessionSearchQueryPlan): SessionSearchQueryPlan | null {
let changed = false
const body = plan.body.map((term) => {
const fix = this.typoRepair.correct(term)
if (fix && fix !== term.toLowerCase()) {
changed = true
return fix
}
return term
})
return changed ? planSessionSearchQuery(body.join(' ')) : null
}
/** Phrase, then AND, for literal-looking queries; null when neither matches. */
private retrieveLiteral(
plan: SessionSearchQueryPlan,
tier: 'full' | 'conversation'
): { rows: MessageRow[]; route: 'phrase' | 'and' } | null {
if (!plan.literal || plan.body.length === 0) {
return null
}
// A one-token literal (`resolveTerminalPath`, `src/a/b.ts`) is its own
// phrase: the tokenizer keeps it whole, so the exact token is the cheap,
// precise first try before the identifier pieces fan out over OR.
const phrase = this.match(phraseExpression(plan.body), tier)
if (phrase.length > 0) {
return { rows: phrase, route: 'phrase' }
}
if (plan.body.length < 2) {
return null
}
const and = this.match(andExpression(plan.body), tier)
return and.length > 0 ? { rows: and, route: 'and' } : null
}
private match(expression: string, tier: 'full' | 'conversation'): MessageRow[] {
const table = tier === 'full' ? 'messages_fts' : 'conversation_fts'
const weights = tier === 'full' ? FULL_WEIGHTS : CONVERSATION_WEIGHTS
try {
// Why: FTS5 aux functions (bm25, snippet) take the table name, never an alias.
return this.db
.prepare(
`SELECT ${table}.rowid AS rowid, -bm25(${table}, ${weights}) AS score,
m.session_row_id, m.role, m.ts
FROM ${table} JOIN messages m ON m.id = ${table}.rowid
WHERE ${table} MATCH ? ORDER BY score DESC LIMIT ${MESSAGE_CANDIDATE_LIMIT}`
)
.all(expression) as MessageRow[]
} catch {
// A term the tokenizer rejects outright (e.g. only punctuation) is a miss, not a fault.
return []
}
}
private rollUp(
rows: MessageRow[],
args: AiVaultSearchArgs,
tier: 'full' | 'conversation'
): AiVaultSearchHit[] {
const best = new Map<number, MessageRow>()
for (const row of rows) {
const current = best.get(row.session_row_id)
if (!current || row.score > current.score) {
best.set(row.session_row_id, row)
}
}
if (best.size === 0) {
return []
}
const sessions = this.loadSessions([...best.keys()], args)
const limit = Math.min(args.limit ?? AI_VAULT_SEARCH_LIMIT_DEFAULT, AI_VAULT_SEARCH_LIMIT_MAX)
const scored = sessions.map((session) => {
const message = best.get(session.id)!
return {
session,
message,
score: message.score - LENGTH_PRIOR * Math.log(1 + session.message_count)
}
})
scored.sort((left, right) =>
args.sort === 'newest'
? (right.session.updated_at ?? '').localeCompare(left.session.updated_at ?? '')
: right.score - left.score
)
const table = tier === 'full' ? 'messages_fts' : 'conversation_fts'
return scored.slice(0, limit).map(({ session, message, score }) => ({
agent: session.agent,
sessionId: session.session_id,
filePath: session.file_path,
codexHome: session.codex_home,
title: session.title,
cwd: session.cwd,
branch: session.branch,
updatedAt: session.updated_at,
messageCount: session.message_count,
resumeCommand: session.resume_command,
score,
evidence: {
role: message.role as AiVaultSearchHit['evidence']['role'],
timestamp: message.ts,
snippet: this.snippet(table, message.rowid, args.query)
}
}))
}
private snippet(table: string, rowid: number, query: string): string {
const expression = orExpression(planSessionSearchQuery(query).terms)
try {
// Why: a bound `rowid = ?` or `rowid IN (?)` next to MATCH is silently
// ignored by the FTS5 planner (it returns the first match); only the
// subselect form is honoured. Column -1 picks whichever column matched.
const row = this.db
.prepare(
`SELECT snippet(${table}, -1, '[', ']', '…', ${SNIPPET_TOKENS}) AS s
FROM ${table} WHERE ${table} MATCH ? AND rowid IN (SELECT ?)`
)
.get(expression, rowid) as { s: string } | undefined
return row?.s ?? ''
} catch {
return ''
}
}
private loadSessions(ids: number[], args: AiVaultSearchArgs): SessionRow[] {
const conditions = [`id IN (${ids.map(() => '?').join(',')})`]
const values: (string | number)[] = [...ids]
if (args.agents && args.agents.length > 0) {
conditions.push(`agent IN (${args.agents.map(() => '?').join(',')})`)
values.push(...args.agents)
}
if (args.since) {
conditions.push('updated_at >= ?')
values.push(args.since)
}
if (args.scopePaths && args.scopePaths.length > 0) {
conditions.push(`(${args.scopePaths.map(() => '(cwd = ? OR cwd LIKE ?)').join(' OR ')})`)
for (const scope of args.scopePaths) {
const trimmed = scope.replace(/[\\/]+$/, '')
values.push(trimmed, `${trimmed}${trimmed.includes('\\') ? '\\' : '/'}%`)
}
}
return this.db
.prepare(`SELECT * FROM sessions WHERE ${conditions.join(' AND ')}`)
.all(...values) as SessionRow[]
}
}
@@ -0,0 +1,108 @@
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 = 2
// unicode61 keeps `_ . - /` inside tokens so paths and identifiers match exactly;
// the `identifiers` column carries the split form (see session-search-identifier-split).
const TOKENIZER = `tokenize="unicode61 tokenchars '_.-/'"`
const SCHEMA_SQL = `
CREATE TABLE IF NOT EXISTS meta(key TEXT PRIMARY KEY, value TEXT NOT NULL);
CREATE TABLE IF NOT EXISTS sessions(
id INTEGER PRIMARY KEY,
agent TEXT NOT NULL,
session_id TEXT NOT NULL,
-- Not unique: OpenCode/Cursor SQLite sessions share one store path; files.path is the key.
file_path TEXT NOT NULL,
codex_home TEXT,
title TEXT NOT NULL,
cwd TEXT,
branch TEXT,
created_at TEXT,
updated_at TEXT,
message_count INTEGER NOT NULL DEFAULT 0,
resume_command TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS sessions_agent ON sessions(agent);
CREATE INDEX IF NOT EXISTS sessions_updated_at ON sessions(updated_at);
CREATE TABLE IF NOT EXISTS files(
path TEXT PRIMARY KEY,
dev INTEGER,
ino INTEGER,
byte_offset INTEGER NOT NULL,
mtime_ms REAL NOT NULL,
size_bytes INTEGER,
session_row_id INTEGER
);
CREATE TABLE IF NOT EXISTS messages(
id INTEGER PRIMARY KEY,
session_row_id INTEGER NOT NULL,
role TEXT NOT NULL,
ts TEXT
);
CREATE INDEX IF NOT EXISTS messages_session ON messages(session_row_id);
CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5(
user_text, assistant_text, tool_text, identifiers, ${TOKENIZER}, detail=full
);
CREATE VIRTUAL TABLE IF NOT EXISTS conversation_fts USING fts5(
user_text, assistant_text, ${TOKENIZER}, detail=full
);
CREATE VIRTUAL TABLE IF NOT EXISTS messages_vocab USING fts5vocab(messages_fts, 'row');
CREATE TABLE IF NOT EXISTS search_log(
id INTEGER PRIMARY KEY,
ts TEXT NOT NULL,
query TEXT NOT NULL,
route TEXT NOT NULL,
hits INTEGER NOT NULL,
duration_ms REAL NOT NULL
);
`
const DROP_SQL = `
DROP TABLE IF EXISTS messages_vocab;
DROP TABLE IF EXISTS conversation_fts;
DROP TABLE IF EXISTS messages_fts;
DROP TABLE IF EXISTS messages;
DROP TABLE IF EXISTS files;
DROP TABLE IF EXISTS sessions;
DROP TABLE IF EXISTS search_log;
DROP TABLE IF EXISTS meta;
`
export function openSessionSearchDatabase(path: string): SyncDatabase {
const db = new SyncDatabase(path)
db.pragma('journal_mode = WAL')
db.pragma('synchronous = NORMAL')
db.pragma('busy_timeout = 5000')
const version = readSchemaVersion(db)
if (version !== null && version !== SESSION_SEARCH_SCHEMA_VERSION) {
db.exec(DROP_SQL)
}
db.exec(SCHEMA_SQL)
db.prepare('INSERT OR REPLACE INTO meta(key, value) VALUES (?, ?)').run(
'schema_version',
String(SESSION_SEARCH_SCHEMA_VERSION)
)
return db
}
export function openSessionSearchDatabaseReadOnly(path: string): SyncDatabase {
const db = new SyncDatabase(path, { readonly: true, fileMustExist: true })
db.pragma('busy_timeout = 1500')
return db
}
function readSchemaVersion(db: SyncDatabase): number | null {
const table = db
.prepare("SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'meta'")
.get()
if (!table) {
return null
}
const row = db.prepare("SELECT value FROM meta WHERE key = 'schema_version'").get() as
| { value: string }
| undefined
const parsed = row ? Number(row.value) : Number.NaN
return Number.isFinite(parsed) ? parsed : null
}
@@ -0,0 +1,159 @@
import { mkdirSync } from 'node:fs'
import { dirname } from 'node:path'
import type {
AiVaultSearchArgs,
AiVaultSearchCoverage,
AiVaultSearchResult
} from '../../shared/ai-vault-search-types'
import { throwIfAiVaultScanCancelled } from '../ai-vault/ai-vault-scan-cancellation'
import { ensureSessionParseCacheLoaded } from '../ai-vault/session-parse-cache-persistence'
import { sessionCandidatesFromDiscoveries } from '../ai-vault/session-scanner-candidates'
import {
createSessionParseStats,
parseAgentSessionFileCached
} from '../ai-vault/session-scanner-parse-cache'
import { discoverAiVaultSessionSources } from '../ai-vault/session-scanner-source-discovery'
import type { AiVaultScanIssue } from '../../shared/ai-vault-types'
import type { AiVaultScanOptions, SessionFileCandidate } from '../ai-vault/session-scanner-types'
import {
registerSessionSearchIndexSink,
withSessionSearchIndexRequired
} from '../ai-vault/session-search-capture'
import { SessionSearchStore } from './session-search-store'
// Why: the backfill shares the scanner process's cache lane with list scans,
// so it yields between files and never holds the lane for long.
const BACKFILL_YIELD_EVERY_FILES = 8
const BACKFILL_YIELD_MS = 5
// Why: a search must see a session that is being written right now even when
// no list scan has run; re-reading the newest few files per provider is a
// readdir + stat plus the appended bytes, well under the query budget.
const REFRESH_RECENT_PER_AGENT = 12
export type SessionSearchServiceOptions = {
databasePath: string
}
/** Scan roots the backfill enumerates; the parent resolves them so they match list scans. */
export type SessionSearchScanRoots = Omit<AiVaultScanOptions, 'signal' | 'limit' | 'unlimited'>
/**
* Runs inside the ai-vault scanner process. Owns the index, feeds it from every
* parse the list scan performs, and fills in the long tail in the background.
*/
export class SessionSearchService {
private readonly store: SessionSearchStore
private backfillRun: Promise<void> | null = null
private backfillController: AbortController | null = null
constructor(options: SessionSearchServiceOptions) {
mkdirSync(dirname(options.databasePath), { recursive: true })
this.store = new SessionSearchStore(options.databasePath)
registerSessionSearchIndexSink(this.store)
}
/** Starts the backfill if needed, folds any appends list scans noticed, then queries. */
async search(
args: AiVaultSearchArgs,
roots: SessionSearchScanRoots,
signal?: AbortSignal
): Promise<AiVaultSearchResult> {
const backfill = this.ensureBackfill(roots)
if (args.refresh !== false) {
await this.refreshRecent(roots, signal)
await this.reindexStale(signal)
}
void backfill
return this.store.search(args)
}
coverage(roots: SessionSearchScanRoots): AiVaultSearchCoverage {
this.ensureBackfill(roots)
return this.store.coverage()
}
/** Idempotent: a running backfill is reused, a finished one is not restarted. */
ensureBackfill(roots: SessionSearchScanRoots): Promise<void> {
if (!this.backfillRun) {
this.backfillController = new AbortController()
this.backfillRun = this.runBackfill(roots, this.backfillController.signal)
.catch((error) => console.warn('[ai-vault-search] backfill stopped:', error))
.finally(() => {
this.backfillController = null
})
}
return this.backfillRun
}
invalidate(paths: readonly string[]): void {
for (const path of paths) {
this.store.removeFile(path)
}
}
dispose(): void {
this.backfillController?.abort()
registerSessionSearchIndexSink(null)
this.store.close()
}
private async refreshRecent(roots: SessionSearchScanRoots, signal?: AbortSignal): Promise<void> {
const issues: AiVaultScanIssue[] = []
const options: AiVaultScanOptions = { ...roots, signal }
const discoveries = await discoverAiVaultSessionSources({
options,
limitPerAgent: REFRESH_RECENT_PER_AGENT,
issues
})
const candidates = await sessionCandidatesFromDiscoveries(discoveries, options)
await this.parseAll(candidates, signal)
}
private async reindexStale(signal?: AbortSignal): Promise<void> {
const stale = this.store.takeStale()
if (stale.length === 0) {
return
}
await this.parseAll(stale, signal)
}
private async runBackfill(roots: SessionSearchScanRoots, signal: AbortSignal): Promise<void> {
this.store.setBackfillState('running')
try {
await ensureSessionParseCacheLoaded()
const issues: AiVaultScanIssue[] = []
const options: AiVaultScanOptions = { ...roots, signal }
const discoveries = await discoverAiVaultSessionSources({
options,
limitPerAgent: Number.POSITIVE_INFINITY,
issues
})
const candidates = await sessionCandidatesFromDiscoveries(discoveries, options)
await this.parseAll(candidates, signal)
this.store.setBackfillState('complete')
} catch (error) {
this.store.setBackfillState('idle')
throw error
}
}
private async parseAll(candidates: SessionFileCandidate[], signal?: AbortSignal): Promise<void> {
const stats = createSessionParseStats()
let sinceYield = 0
await withSessionSearchIndexRequired(async () => {
for (const candidate of candidates) {
throwIfAiVaultScanCancelled(signal)
try {
await parseAgentSessionFileCached(candidate, process.platform, stats)
} catch (error) {
console.warn('[ai-vault-search] backfill skipped', candidate.file.path, error)
}
sinceYield += 1
if (sinceYield >= BACKFILL_YIELD_EVERY_FILES) {
sinceYield = 0
await new Promise((resolve) => setTimeout(resolve, BACKFILL_YIELD_MS))
}
}
})
}
}
@@ -0,0 +1,250 @@
import { appendFile, mkdtemp, rename, rm, stat, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
import {
createSessionParseStats,
parseAgentSessionFileCached,
resetSessionParseCacheForTests
} from '../ai-vault/session-scanner-parse-cache'
import {
registerSessionSearchIndexSink,
withSessionSearchIndexRequired
} from '../ai-vault/session-search-capture'
import type { SessionFileCandidate } from '../ai-vault/session-scanner-types'
import { SessionSearchStore } from './session-search-store'
const SESSION_ID = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee'
let tempRoots: string[] = []
let store: SessionSearchStore
beforeEach(async () => {
resetSessionParseCacheForTests()
const root = await makeTempDir()
store = new SessionSearchStore(join(root, 'index.sqlite'), (error) => {
throw error
})
registerSessionSearchIndexSink(store)
})
afterEach(async () => {
registerSessionSearchIndexSink(null)
store.close()
await Promise.all(tempRoots.map((root) => rm(root, { recursive: true, force: true })))
tempRoots = []
})
async function makeTempDir(): Promise<string> {
const root = await mkdtemp(join(tmpdir(), 'orca-session-search-'))
tempRoots.push(root)
return root
}
async function claudeCandidate(path: string): Promise<SessionFileCandidate> {
const fileStat = await stat(path)
return {
agent: 'claude',
codexHome: null,
file: {
path,
mtimeMs: fileStat.mtimeMs,
modifiedAt: fileStat.mtime.toISOString(),
sizeBytes: fileStat.size,
dev: fileStat.dev,
ino: fileStat.ino
}
}
}
function userRecord(index: number, content: unknown, sessionId = SESSION_ID): string {
return JSON.stringify({
type: 'user',
sessionId,
timestamp: new Date(1740000000000 + index * 60_000).toISOString(),
cwd: '/repo/app',
gitBranch: 'main',
message: { role: 'user', content }
})
}
function assistantRecord(index: number, content: unknown, sessionId = SESSION_ID): string {
return JSON.stringify({
type: 'assistant',
sessionId,
timestamp: new Date(1740000000000 + index * 60_000).toISOString(),
message: { role: 'assistant', model: 'claude-fable-5', content }
})
}
async function parse(path: string) {
const stats = createSessionParseStats()
const session = await parseAgentSessionFileCached(
await claudeCandidate(path),
process.platform,
stats
)
return { session, stats }
}
describe('SessionSearchStore', () => {
it('indexes a transcript through the parse cache and finds it by mid-session text', async () => {
const root = await makeTempDir()
const path = join(root, `${SESSION_ID}.jsonl`)
await writeFile(
path,
`${[
userRecord(0, 'first question about the tab switcher'),
assistantRecord(1, [{ type: 'text', text: 'Looking at resolveTerminalPath now.' }]),
userRecord(2, 'the locator is flaky again'),
assistantRecord(3, [
{
type: 'tool_use',
id: 'toolu_1',
name: 'Bash',
input: { command: 'pnpm test src/tabs' }
}
]),
userRecord(4, [
{
type: 'tool_result',
tool_use_id: 'toolu_1',
content: 'Error: strict mode violation: getByRole(button) resolved to 2 elements'
}
]),
JSON.stringify({ type: 'ai-title', aiTitle: 'Fix flaky locator' })
].join('\n')}\n`
)
await parse(path)
const literal = store.search({ query: 'strict mode violation getByRole' })
expect(literal.hits).toHaveLength(1)
expect(literal.hits[0]).toMatchObject({
agent: 'claude',
sessionId: SESSION_ID,
title: 'Fix flaky locator',
cwd: '/repo/app',
evidence: { role: 'tool' }
})
expect(literal.hits[0]?.evidence.snippet).toContain('[strict]')
expect(literal.route).toBe('phrase')
// Identifier split: a partial camelCase name still matches.
expect(store.search({ query: 'TerminalPath' }).hits).toHaveLength(1)
// Tool command indexed from the tool_use block.
expect(store.search({ query: 'pnpm test src/tabs' }).hits).toHaveLength(1)
// Conversation tier excludes tool rows but keeps prompts.
expect(store.search({ query: 'getByRole', tier: 'conversation' }).hits).toHaveLength(0)
expect(store.search({ query: 'locator flaky', tier: 'conversation' }).hits).toHaveLength(1)
// user, assistant, user, tool_use, tool_result
expect(store.coverage()).toMatchObject({ sessionsIndexed: 1, messagesIndexed: 5 })
})
it('appends only the new lines on an incremental parse and keeps the cursor in step', async () => {
const root = await makeTempDir()
const path = join(root, `${SESSION_ID}.jsonl`)
await writeFile(path, `${userRecord(0, 'alpha question')}\n`)
await parse(path)
expect(store.search({ query: 'omega' }).hits).toHaveLength(0)
await appendFile(path, `${assistantRecord(1, 'omega answer')}\n`)
const { stats } = await parse(path)
expect(stats.incremental).toBe(1)
expect(store.search({ query: 'omega' }).hits).toHaveLength(1)
expect(store.search({ query: 'alpha' }).hits).toHaveLength(1)
expect(store.coverage().messagesIndexed).toBe(2)
})
it('re-indexes a rename-replaced file instead of appending onto stale rows', async () => {
const root = await makeTempDir()
const path = join(root, `${SESSION_ID}.jsonl`)
await writeFile(path, `${userRecord(0, 'stale question')}\n`)
await parse(path)
const staging = join(root, 'staging.jsonl')
await writeFile(
staging,
`${userRecord(0, 'fresh question')}\n${assistantRecord(1, 'fresh answer')}\n`
)
await rename(staging, path)
await parse(path)
expect(store.search({ query: 'stale' }).hits).toHaveLength(0)
expect(store.search({ query: 'fresh' }).hits).toHaveLength(1)
expect(store.coverage().messagesIndexed).toBe(2)
})
it('marks a cache-known file stale in opportunistic mode and re-parses it in required mode', async () => {
const root = await makeTempDir()
const path = join(root, `${SESSION_ID}.jsonl`)
await writeFile(path, `${userRecord(0, 'before the index existed')}\n`)
registerSessionSearchIndexSink(null)
await parse(path)
expect(store.coverage().sessionsIndexed).toBe(0)
// A list scan must not pay for the index: reuse the cache, queue the file.
registerSessionSearchIndexSink(store)
const opportunistic = await parse(path)
expect(opportunistic.stats.reused).toBe(1)
expect(store.coverage()).toMatchObject({ sessionsIndexed: 0, filesPending: 1 })
// The backfill lane drains the queue with a whole-file parse.
const stale = store.takeStale()
expect(stale.map((candidate) => candidate.file.path)).toEqual([path])
const required = await withSessionSearchIndexRequired(() => parse(path))
expect(required.stats.reused).toBe(0)
expect(store.search({ query: 'before the index existed' }).hits).toHaveLength(1)
expect(store.coverage().filesPending).toBe(0)
})
it('repairs a typo from the index vocabulary', async () => {
const root = await makeTempDir()
const path = join(root, `${SESSION_ID}.jsonl`)
await writeFile(
path,
`${userRecord(0, 'the watcher coalesces events')}\n${assistantRecord(1, 'watcher coalesces them')}\n`
)
await parse(path)
const result = store.search({ query: 'watcher coalesces' })
expect(result.hits).toHaveLength(1)
const typo = store.search({ query: 'watcher coalesces'.replace('coalesces', 'coalesecs') })
expect(typo.hits).toHaveLength(1)
expect(typo.route).toMatch(/^typo\+/)
expect(typo.repairedTerms).toContain('coalesces')
})
it('ranks by newest when asked and filters by agent and scope', async () => {
const root = await makeTempDir()
const older = join(root, 'aaaaaaaa-0000-4000-8000-000000000001.jsonl')
const newer = join(root, 'aaaaaaaa-0000-4000-8000-000000000002.jsonl')
await writeFile(
older,
`${userRecord(0, 'shared phrase one', 'aaaaaaaa-0000-4000-8000-000000000001')}\n`
)
await writeFile(
newer,
`${userRecord(500, 'shared phrase two', 'aaaaaaaa-0000-4000-8000-000000000002')}\n`
)
await parse(older)
await parse(newer)
const newest = store.search({ query: 'shared phrase', sort: 'newest' })
expect(newest.hits.map((hit) => hit.sessionId)).toEqual([
'aaaaaaaa-0000-4000-8000-000000000002',
'aaaaaaaa-0000-4000-8000-000000000001'
])
expect(store.search({ query: 'shared phrase', agents: ['codex'] }).hits).toHaveLength(0)
expect(store.search({ query: 'shared phrase', scopePaths: ['/repo'] }).hits).toHaveLength(2)
expect(store.search({ query: 'shared phrase', scopePaths: ['/other'] }).hits).toHaveLength(0)
})
it('does not choke on FTS5 syntax in user text', async () => {
const root = await makeTempDir()
const path = join(root, `${SESSION_ID}.jsonl`)
await writeFile(path, `${userRecord(0, 'run cli.mjs with foo-bar and C++')}\n`)
await parse(path)
for (const query of ['cli.mjs', 'foo-bar', 'C++', '"quoted phrase"', 'AND OR NOT']) {
expect(() => store.search({ query })).not.toThrow()
}
expect(store.search({ query: 'cli.mjs' }).hits).toHaveLength(1)
expect(store.search({ query: 'foo-bar' }).hits).toHaveLength(1)
})
})
@@ -0,0 +1,154 @@
import type SyncDatabase from '../sqlite/sync-database'
import type {
AiVaultSearchArgs,
AiVaultSearchCoverage,
AiVaultSearchProviderCoverage,
AiVaultSearchResult
} from '../../shared/ai-vault-search-types'
import type { AiVaultAgent } from '../../shared/ai-vault-types'
import type {
SessionSearchFileIdentity,
SessionSearchIndexedFile,
SessionSearchIndexSink,
SessionSearchIndexUpdate
} from '../ai-vault/session-search-capture'
import type { SessionFileCandidate } from '../ai-vault/session-scanner-types'
import { SessionSearchIndexWriter } from './session-search-index-writer'
import { SessionSearchQuery } from './session-search-query'
import { openSessionSearchDatabase } from './session-search-schema'
// Why: the query log keeps only the surface form so the eval set can be rebuilt
// from real usage (the shoot-out's highest-value follow-up); bounded ring.
const SEARCH_LOG_LIMIT = 5000
export type SessionSearchBackfillState = 'idle' | 'running' | 'complete'
/** Owns the index database: the scanner writes through it, search reads from it. */
export class SessionSearchStore implements SessionSearchIndexSink {
private readonly db: SyncDatabase
private readonly writer: SessionSearchIndexWriter
private readonly query: SessionSearchQuery
private backfill: SessionSearchBackfillState = 'idle'
private lastIndexedAt: string | null = null
private applyFailures = 0
private readonly stale = new Map<string, SessionFileCandidate>()
constructor(
path: string,
private readonly onError: (error: unknown) => void = (error) =>
console.warn('[ai-vault-search] index write failed:', error)
) {
this.db = openSessionSearchDatabase(path)
this.writer = new SessionSearchIndexWriter(this.db)
this.query = new SessionSearchQuery(this.db)
}
indexedFile(path: string, identity: SessionSearchFileIdentity): SessionSearchIndexedFile | null {
try {
return this.writer.indexedFile(path, identity)
} catch (error) {
this.onError(error)
return null
}
}
apply(update: SessionSearchIndexUpdate): void {
try {
this.writer.apply(update)
this.lastIndexedAt = new Date().toISOString()
} catch (error) {
this.applyFailures += 1
this.onError(error)
}
}
markStale(candidate: SessionFileCandidate): void {
this.stale.set(candidate.file.path, candidate)
}
/** Hands the stale set to the backfill lane and clears it. */
takeStale(): SessionFileCandidate[] {
const candidates = [...this.stale.values()]
this.stale.clear()
return candidates
}
get staleCount(): number {
return this.stale.size
}
removeFile(path: string): void {
try {
this.writer.removeFile(path)
} catch (error) {
this.onError(error)
}
}
setBackfillState(state: SessionSearchBackfillState): void {
this.backfill = state
}
search(args: AiVaultSearchArgs): AiVaultSearchResult {
const startedAt = performance.now()
const execution = this.query.execute(args)
const durationMs = performance.now() - startedAt
this.logQuery(args.query, execution.route, execution.hits.length, durationMs)
return {
hits: execution.hits,
route: execution.route,
...(execution.repairedTerms ? { repairedTerms: execution.repairedTerms } : {}),
durationMs,
coverage: this.coverage()
}
}
coverage(): AiVaultSearchCoverage {
const providers = this.db
.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
GROUP BY s.agent ORDER BY s.agent`
)
.all() as { agent: AiVaultAgent; sessions: number; messages: number }[]
const byProvider: AiVaultSearchProviderCoverage[] = providers.map((row) => ({
agent: row.agent,
sessionsIndexed: row.sessions,
messagesIndexed: row.messages
}))
return {
sessionsIndexed: byProvider.reduce((sum, row) => sum + row.sessionsIndexed, 0),
messagesIndexed: byProvider.reduce((sum, row) => sum + row.messagesIndexed, 0),
providers: byProvider,
backfill: this.backfill,
filesPending: this.stale.size,
lastIndexedAt: this.lastIndexedAt
}
}
get failures(): number {
return this.applyFailures
}
close(): void {
this.db.close()
}
private logQuery(query: string, route: string, hits: number, durationMs: number): void {
try {
this.db
.prepare(
'INSERT INTO search_log(ts, query, route, hits, duration_ms) VALUES (?, ?, ?, ?, ?)'
)
.run(new Date().toISOString(), query, route, hits, durationMs)
this.db
.prepare(
`DELETE FROM search_log WHERE id <= (
SELECT id FROM search_log ORDER BY id DESC LIMIT 1 OFFSET ?)`
)
.run(SEARCH_LOG_LIMIT)
} catch (error) {
this.onError(error)
}
}
}
@@ -0,0 +1,99 @@
import type SyncDatabase from '../sqlite/sync-database'
// Why: a query term with zero postings is usually a typo. The index's own
// vocabulary (fts5vocab) is the dictionary, so repair needs no model and can
// never suggest a word the index does not contain. Measured MRR 0.553 → 0.566.
const MIN_TERM_LENGTH = 4
const MAX_TERM_LENGTH = 40
const LENGTH_SLACK = 2
const MIN_DOC_FREQUENCY = 2
const MIN_SIMILARITY = 0.82
const MAX_CANDIDATES = 4000
type VocabRow = { term: string; doc: number }
// Longest common subsequence length; the indel distance is len(a)+len(b)-2·LCS.
function commonSubsequenceLength(a: string, b: string): number {
let previous = Array.from<number>({ length: b.length + 1 }).fill(0)
let current = Array.from<number>({ length: b.length + 1 }).fill(0)
for (let i = 1; i <= a.length; i += 1) {
for (let j = 1; j <= b.length; j += 1) {
current[j] =
a.charCodeAt(i - 1) === b.charCodeAt(j - 1)
? previous[j - 1] + 1
: Math.max(previous[j], current[j - 1])
}
;[previous, current] = [current, previous]
}
return previous[b.length]
}
/** Normalized indel similarity in [0, 1], the scale rapidfuzz's `fuzz.ratio` uses. */
function similarity(a: string, b: string): number {
const total = a.length + b.length
return total === 0 ? 1 : (2 * commonSubsequenceLength(a, b)) / total
}
export class SessionSearchTypoRepair {
private readonly documentFrequency: ReturnType<SyncDatabase['prepare']>
private readonly candidatesByPrefix: ReturnType<SyncDatabase['prepare']>
constructor(db: SyncDatabase) {
this.documentFrequency = db.prepare('SELECT doc FROM messages_vocab WHERE term = ?')
// fts5vocab is ordered by term, so a prefix range plus a length band is a
// bounded scan; the most frequent terms are kept when the band overflows.
this.candidatesByPrefix = db.prepare(
`SELECT term, doc FROM messages_vocab
WHERE term >= ? AND term < ? AND length(term) BETWEEN ? AND ? AND doc >= ?
ORDER BY doc DESC LIMIT ?`
)
}
hasPostings(term: string): boolean {
const row = this.documentFrequency.get(term.toLowerCase()) as VocabRow | undefined
return row !== undefined && row.doc > 0
}
/** Returns the closest indexed term, or null when `term` exists or nothing is close enough. */
correct(term: string): string | null {
const lowered = term.toLowerCase()
if (lowered.length < MIN_TERM_LENGTH || lowered.length > MAX_TERM_LENGTH) {
return null
}
if (this.hasPostings(lowered)) {
return null
}
// Two-letter prefix first (a typo rarely hits both), then the transposed
// pair, then the bare first letter as the wide fallback.
const prefixes = [lowered.slice(0, 2), lowered[1] + lowered[0], lowered[0]]
let best: { term: string; score: number; doc: number } | null = null
for (const prefix of prefixes) {
for (const row of this.candidates(prefix, lowered.length)) {
const score = similarity(lowered, row.term)
if (score < MIN_SIMILARITY) {
continue
}
if (!best || score > best.score || (score === best.score && row.doc > best.doc)) {
best = { term: row.term, score, doc: row.doc }
}
}
if (best) {
return best.term
}
}
return null
}
private candidates(prefix: string, length: number): VocabRow[] {
const last = prefix.charCodeAt(prefix.length - 1)
const upper = prefix.slice(0, -1) + String.fromCharCode(last + 1)
return this.candidatesByPrefix.all(
prefix,
upper,
Math.max(MIN_TERM_LENGTH - 1, length - LENGTH_SLACK),
length + LENGTH_SLACK,
MIN_DOC_FREQUENCY,
MAX_CANDIDATES
) as VocabRow[]
}
}
+43 -14
View File
@@ -1,14 +1,22 @@
import { join } from 'node:path'
import {
clearAiVaultBackgroundRestartCircuit,
readAiVaultSearchCoverageInBackground,
resetAiVaultScannerBackgroundForTests,
scanAiVaultSessionsInBackground
scanAiVaultSessionsInBackground,
searchAiVaultSessionsInBackground
} from './session-scanner-background'
import { listRunningWslHomeDirsAsync } from '../wsl'
import { filterPathsToRunningWslDistrosAsync } from '../wsl-running-path-filter'
import type { AiVaultListArgs, AiVaultListResult } from '../../shared/ai-vault-types'
import type {
AiVaultSearchArgs,
AiVaultSearchCoverage,
AiVaultSearchResult
} from '../../shared/ai-vault-search-types'
import { LOCAL_EXECUTION_HOST_ID } from '../../shared/execution-host'
import { AiVaultScanCoordinator } from './ai-vault-scan-coordinator'
import type { AiVaultWorkerScanOptions } from './session-scanner-worker-protocol'
import {
aiVaultSessionDepthCovers,
requestedAiVaultSessionDepth,
@@ -49,6 +57,39 @@ export function configureAiVaultSessionSources(next: AiVaultSessionSources): voi
sources = next
}
/** Host-local source roots every scan of this machine shares (managed Codex homes, WSL homes). */
export async function resolveAiVaultHostScanSources(): Promise<
Pick<AiVaultWorkerScanOptions, 'additionalCodexSessionsDirs' | 'wslHomeDirs' | 'executionHostId'>
> {
const configuredCodexHomes = sources.getAdditionalCodexHomePaths?.() ?? []
const [additionalCodexHomes, wslHomeDirs] = await Promise.all([
filterPathsToRunningWslDistrosAsync(configuredCodexHomes),
getAiVaultWslHomeDirs()
])
return {
additionalCodexSessionsDirs: additionalCodexHomes.map((homePath) => join(homePath, 'sessions')),
wslHomeDirs,
// Why: this scan is always host-local; callers addressing this host by a
// runtime id get the result restamped at the RPC edge, never rescanned.
executionHostId: LOCAL_EXECUTION_HOST_ID
}
}
export async function searchAiVaultSessions(
args: AiVaultSearchArgs,
options: { signal?: AbortSignal } = {}
): Promise<AiVaultSearchResult> {
const roots = await resolveAiVaultHostScanSources()
return searchAiVaultSessionsInBackground({ args, roots }, options.signal)
}
export async function readAiVaultSearchCoverage(
options: { signal?: AbortSignal } = {}
): Promise<AiVaultSearchCoverage> {
const roots = await resolveAiVaultHostScanSources()
return readAiVaultSearchCoverageInBackground({ roots }, options.signal)
}
export async function listAiVaultSessions(
args?: AiVaultListArgs,
options: { signal?: AbortSignal } = {}
@@ -80,24 +121,12 @@ export async function listAiVaultSessions(
force: args?.force,
signal: options.signal,
start: async (scanSignal) => {
const configuredCodexHomes = sources.getAdditionalCodexHomePaths?.() ?? []
const [additionalCodexHomes, wslHomeDirs] = await Promise.all([
filterPathsToRunningWslDistrosAsync(configuredCodexHomes),
getAiVaultWslHomeDirs()
])
const additionalCodexSessionsDirs = additionalCodexHomes.map((homePath) =>
join(homePath, 'sessions')
)
const result = await scanAiVaultSessionsInBackground(
{
limit: args?.limit,
unlimited: args?.unlimited,
scopePaths: args?.scopePaths,
additionalCodexSessionsDirs,
wslHomeDirs,
// Why: this scan is always host-local; callers addressing this host by a
// runtime id get the result restamped at the RPC edge, never rescanned.
executionHostId: LOCAL_EXECUTION_HOST_ID
...(await resolveAiVaultHostScanSources())
},
scanSignal
)
@@ -0,0 +1,78 @@
import type { SessionParseCacheEntry } from './session-scanner-parse-cache'
// Sized past the default recency cap (1000) plus the in-scope cap (2000) so a
// full steady-state result set stays resident between forced rescans.
const MAX_CACHE_ENTRIES = 4096
const cache = new Map<string, SessionParseCacheEntry>()
export function resetSessionParseCacheForTests(): void {
cache.clear()
}
// Drops one entry after its file is deleted. Cleanliness, not correctness:
// discovery walks disk first, so a trashed file is never rediscovered anyway.
export function invalidateSessionParseCacheEntry(path: string): void {
cache.delete(path)
}
// Persisted subset of a cache entry: the non-serializable `resume` parser
// state is dropped (see session-parse-cache-persistence.ts).
export type PersistedSessionParseCacheEntry = Omit<SessionParseCacheEntry, 'resume'>
export function snapshotSessionParseCacheForPersistence(): [
string,
PersistedSessionParseCacheEntry
][] {
return [...cache].map(([path, entry]): [string, PersistedSessionParseCacheEntry] => [
path,
{
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session
}
])
}
// Seeded entries carry `resume: null`: after a restart an unchanged file is a
// cache hit; a file that changed while the app was closed pays one full
// (not incremental) re-parse.
export function seedSessionParseCache(
entries: Iterable<[string, PersistedSessionParseCacheEntry]>
): void {
const list = [...entries]
// Snapshot order is oldest→newest (LRU); an over-cap list keeps the newest
// tail rather than seeding the oldest entries and dropping the tail.
for (const [path, entry] of list.slice(Math.max(0, list.length - MAX_CACHE_ENTRIES))) {
if (cache.size >= MAX_CACHE_ENTRIES) {
return
}
// In-process entries are always fresher than persisted ones; never clobber.
if (cache.has(path)) {
continue
}
cache.set(path, {
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session,
resume: null
})
}
}
export function storeSessionParseCacheEntry(path: string, entry: SessionParseCacheEntry): void {
cache.delete(path)
cache.set(path, entry)
if (cache.size > MAX_CACHE_ENTRIES) {
const oldest = cache.keys().next()
if (!oldest.done) {
cache.delete(oldest.value)
}
}
}
export function getSessionParseCacheEntry(path: string): SessionParseCacheEntry | undefined {
return cache.get(path)
}
@@ -23,6 +23,7 @@ import {
normalizePreviewText,
timestampMs
} from './session-scanner-values'
import { captureIndexableContent, captureIndexableText } from './session-search-content'
const SESSION_PREVIEW_MESSAGE_LIMIT = 5
@@ -161,6 +162,8 @@ export function addPreviewMessage(
// Why: Claude meta/injected turns still preview, but must not seed the
// copyable first-prompt row.
seedFirstUserPrompt?: boolean
// Set by addPreviewContent, which already captured the search rows.
indexedByContent?: boolean
}
): void {
// Seeded before the preview-empty return so the copy body never depends on
@@ -171,6 +174,11 @@ export function addPreviewMessage(
() => (args.text ? normalizeFullFirstUserPromptText(args.text) : null),
args.seedFirstUserPrompt
)
// Search rows ride the same funnel every parser already feeds; the content
// path below captures its own richer (tool-aware) rows before reaching here.
if (!args.indexedByContent) {
captureIndexableText(args.role, args.text, args.timestamp)
}
const text = normalizePreviewText(args.text ?? '')
if (!text) {
return
@@ -199,12 +207,14 @@ export function addPreviewContent(
() => extractFullFirstUserPromptText(content),
options?.seedFirstUserPrompt
)
captureIndexableContent(role, content, timestamp)
addPreviewMessage(accumulator, {
role,
text: extractPreviewContentText(content),
timestamp,
// Content path already seeded above when capture is enabled.
seedFirstUserPrompt: false
seedFirstUserPrompt: false,
indexedByContent: true
})
}
@@ -13,16 +13,24 @@ import {
invalidateAiVaultServiceCache,
listAiVaultSubagentSessionsInService,
readAiVaultFirstUserPromptInService,
readAiVaultSearchCoverageInService,
resetAiVaultScannerServiceForTests,
resolveAiVaultSessionTitlesInService,
scanAiVaultSessionsInService
scanAiVaultSessionsInService,
searchAiVaultSessionsInService
} from './session-scanner-service-spawn'
import type { AiVaultServiceSubagentRequest } from './session-scanner-service-protocol'
import type {
AiVaultServiceSearchRequest,
AiVaultServiceSubagentRequest
} from './session-scanner-service-protocol'
import {
readAiVaultSearchCoverageInWorker,
resetAiVaultScannerWorkerForTests,
resolveAiVaultSessionTitlesInWorker,
scanAiVaultSessionsInWorker
scanAiVaultSessionsInWorker,
searchAiVaultSessionsInWorker
} from './session-scanner-worker-spawn'
import type { AiVaultSearchCoverage, AiVaultSearchResult } from '../../shared/ai-vault-search-types'
import type { AiVaultWorkerScanOptions } from './session-scanner-worker-protocol'
import { listLocalAiVaultSubagentSessions } from './session-subagent-reader'
@@ -78,6 +86,24 @@ export function readAiVaultFirstUserPromptInBackground(
: readAiVaultFirstUserPrompt(request)
}
export function searchAiVaultSessionsInBackground(
request: AiVaultServiceSearchRequest,
signal?: AbortSignal
): Promise<AiVaultSearchResult> {
return shouldUseAiVaultServiceProcess()
? searchAiVaultSessionsInService(request, signal)
: searchAiVaultSessionsInWorker(request, signal)
}
export function readAiVaultSearchCoverageInBackground(
request: Pick<AiVaultServiceSearchRequest, 'roots'>,
signal?: AbortSignal
): Promise<AiVaultSearchCoverage> {
return shouldUseAiVaultServiceProcess()
? readAiVaultSearchCoverageInService(request, signal)
: readAiVaultSearchCoverageInWorker(request)
}
export function invalidateAiVaultBackgroundCache(paths: string[]): Promise<void> {
return shouldUseAiVaultServiceProcess() ? invalidateAiVaultServiceCache(paths) : Promise.resolve()
}
@@ -0,0 +1,46 @@
import { readCodexRolloutSessionMetaId } from '../codex/codex-rollout-session-meta'
import { codexRolloutHardlinkIdentity, dedupeCodexRolloutAliases } from './codex-session-root-dedup'
import { antigravityHistoryPathForBrainDir } from './session-scanner-antigravity-paths'
import { codexHomeForSessionsDir } from './session-scanner-codex-paths'
import { DEFAULT_CODEX_HOME_DIR } from './session-scanner-source-discovery'
import type {
AiVaultScanOptions,
SessionFileCandidate,
SessionFileDiscovery
} from './session-scanner-types'
/** Newest-first parse candidates for a discovery set, with Codex hardlink aliases collapsed. */
export async function sessionCandidatesFromDiscoveries(
discoveries: SessionFileDiscovery[],
options: AiVaultScanOptions
): Promise<SessionFileCandidate[]> {
return dedupeCodexRolloutAliases(
discoveries
.flatMap((discovery) =>
discovery.files.map((file): SessionFileCandidate => ({
agent: discovery.agent,
file,
codexHome:
discovery.agent === 'codex'
? codexHomeForSessionsDir(
discovery.rootDir,
options.defaultCodexHomeDir ?? DEFAULT_CODEX_HOME_DIR
)
: null,
antigravityHistoryPath:
discovery.agent === 'antigravity'
? antigravityHistoryPathForBrainDir(discovery.rootDir)
: undefined
}))
)
.sort((left, right) => right.file.mtimeMs - left.file.mtimeMs),
{
isCodex: (candidate) => candidate.agent === 'codex',
getFilePath: (candidate) => candidate.file.path,
getCodexHome: (candidate) => candidate.codexHome,
getHardlinkIdentity: (candidate) => codexRolloutHardlinkIdentity(candidate.file)
},
(filePath) => readCodexRolloutSessionMetaId(filePath, options.signal, 'scan'),
options.signal
)
}
@@ -29,12 +29,17 @@ import {
extractModel,
extractString,
normalizeCodexUsage,
normalizeTitleText,
parseJsonObject,
subtractCodexUsage
} from './session-scanner-values'
import { remoteSessionContentLines } from './remote-session-content-lines'
import { readCodexTimelineOnlyRecord } from './session-scanner-codex-record-fast-path'
import { captureCodexToolRecord } from './session-search-codex-tool-records'
import { isSessionSearchCaptureActive } from './session-search-capture'
import {
extractCodexSessionMetadataTitle,
isCodexWorkerSession
} from './session-scanner-codex-session-meta'
export async function parseCodexSessionFile(
file: FileWithMtime,
@@ -167,6 +172,8 @@ function consumeCodexRecordLine(state: CodexSessionParseState, line: string): vo
return
}
captureCodexToolRecord(record.type, payload, record.timestamp, state.historyMode)
if (record.type === 'response_item' && payload.type === 'message') {
if (state.historyMode === 'paginated') {
return
@@ -272,7 +279,10 @@ function codexResumeStateFromParseState(
return {
consumeLine: (line) => consumeCodexRecordLine(state, line),
consumeLineBytes: (line) => {
const timelineOnlyRecord = readCodexTimelineOnlyRecord(line)
// The prefix fast path skips exactly the tool records the search index wants.
const timelineOnlyRecord = isSessionSearchCaptureActive()
? null
: readCodexTimelineOnlyRecord(line)
if (timelineOnlyRecord) {
updateTimeline(state.accumulator, timelineOnlyRecord.timestamp)
} else {
@@ -314,21 +324,3 @@ async function parseCodexSessionLines(args: {
executionHostPlatform: args.executionHostPlatform
})
}
function isCodexWorkerSession(payload: Record<string, unknown>): boolean {
const threadSource = extractString(payload.thread_source) ?? extractString(payload.threadSource)
if (threadSource) {
return threadSource.toLowerCase() !== 'user'
}
const source = asRecord(payload.source)
return Boolean(asRecord(source?.subagent))
}
function extractCodexSessionMetadataTitle(payload: Record<string, unknown>): string | null {
return (
normalizeTitleText(extractString(payload.title) ?? '') ??
normalizeTitleText(extractString(payload.thread_name) ?? '') ??
normalizeTitleText(extractString(payload.threadName) ?? '')
)
}
@@ -0,0 +1,20 @@
import { asRecord, extractString, normalizeTitleText } from './session-scanner-values'
/** Codex writes worker/sub-agent transcripts into the same tree; only user-started threads list. */
export function isCodexWorkerSession(payload: Record<string, unknown>): boolean {
const threadSource = extractString(payload.thread_source) ?? extractString(payload.threadSource)
if (threadSource) {
return threadSource.toLowerCase() !== 'user'
}
const source = asRecord(payload.source)
return Boolean(asRecord(source?.subagent))
}
export function extractCodexSessionMetadataTitle(payload: Record<string, unknown>): string | null {
return (
normalizeTitleText(extractString(payload.title) ?? '') ??
normalizeTitleText(extractString(payload.thread_name) ?? '') ??
normalizeTitleText(extractString(payload.threadName) ?? '')
)
}
@@ -3,7 +3,7 @@ import { openTranscriptReadStream } from '../native-chat/wsl-transcript-fs-acces
const NEWLINE_BYTE = 0x0a
const CARRIAGE_RETURN_BYTE = 0x0d
type JsonlReadResult = {
export type JsonlReadResult = {
consumedThrough: number
trailingPartialLine: string | null
bytesRead: number
@@ -1,4 +1,4 @@
import { appendFile, mkdir, mkdtemp, rm, stat, writeFile } from 'node:fs/promises'
import { appendFile, mkdir, mkdtemp, rename, rm, stat, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
@@ -38,6 +38,14 @@ async function claudeCandidate(path: string): Promise<SessionFileCandidate> {
return { agent: 'claude', file, codexHome: null }
}
// Mirrors discovery, which stats identity; the plain helper above models
// synthetic candidates that carry no inode.
async function identifiedClaudeCandidate(path: string): Promise<SessionFileCandidate> {
const candidate = await claudeCandidate(path)
const fileStat = await stat(path)
return { ...candidate, file: { ...candidate.file, dev: fileStat.dev, ino: fileStat.ino } }
}
function userRecord(index: number, text: string): string {
return JSON.stringify({
type: 'user',
@@ -231,6 +239,59 @@ describe('parseAgentSessionFileCached', () => {
expect(reparsed?.messageCount).toBe(2)
})
it('refuses to resume across a rename-replace even when the old offset lands on a newline', async () => {
const root = await makeTempDir()
const path = join(root, 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee.jsonl')
const original = `${userRecord(0, 'aa')}\n`
await writeFile(path, original)
const stats = createSessionParseStats()
await parseAgentSessionFileCached(
await identifiedClaudeCandidate(path),
process.platform,
stats
)
// A larger replacement whose byte at the old offset is exactly '\n', so the
// newline guard alone would resume mid-file and skip the first record.
const padded = userRecord(0, 'aa'.padEnd(original.length - 1 - userRecord(0, '').length, 'b'))
const replacement = `${padded}\n${assistantRecord(1, 'answer')}\n${userRecord(2, 'more')}\n`
expect(replacement[original.length - 1]).toBe('\n')
expect(replacement.length).toBeGreaterThan(original.length)
const staging = join(root, 'staging.jsonl')
await writeFile(staging, replacement)
await rename(staging, path)
const reparsed = await parseAgentSessionFileCached(
await identifiedClaudeCandidate(path),
process.platform,
stats
)
expect(stats.fullParses).toBe(2)
expect(stats.incremental).toBe(0)
expect(reparsed).toEqual(await freshParse(path))
expect(reparsed?.messageCount).toBe(3)
})
it('still resumes an append when file identity is unchanged', async () => {
const root = await makeTempDir()
const path = join(root, 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee.jsonl')
await writeFile(path, `${userRecord(0, 'question')}\n`)
const stats = createSessionParseStats()
await parseAgentSessionFileCached(
await identifiedClaudeCandidate(path),
process.platform,
stats
)
await appendFile(path, `${assistantRecord(1, 'answer')}\n`)
const resumed = await parseAgentSessionFileCached(
await identifiedClaudeCandidate(path),
process.platform,
stats
)
expect(stats.incremental).toBe(1)
expect(resumed).toEqual(await freshParse(path))
})
it('parses CRLF transcripts identically to the streaming parser', async () => {
const root = await makeTempDir()
const path = join(root, 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee.jsonl')
+121 -115
View File
@@ -1,4 +1,3 @@
import { readTranscriptSlice } from '../native-chat/wsl-transcript-fs-access'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import { createAntigravitySessionResumeState } from './session-scanner-antigravity-parser'
import { parseAgentSessionFile } from './session-scanner-agent-parser'
@@ -13,21 +12,23 @@ import { countSubagentTranscripts } from './session-scanner-subagent-transcripts
import { countOmpSubagentTranscripts } from './session-scanner-omp-subagent-transcripts'
import type { ResumableSessionParseState, SessionFileCandidate } from './session-scanner-types'
import { refreshCachedCodexTitle } from './session-scanner-codex-cached-title'
import { consumeCompleteJsonlLines } from './session-scanner-jsonl-reader'
import { consumeCompleteJsonlLines, type JsonlReadResult } from './session-scanner-jsonl-reader'
import { getSessionParseCacheEntry, storeSessionParseCacheEntry } from './session-parse-cache-store'
import {
endsWithNewlineAt,
fileIdentity,
sameFileIdentity,
type ResumePoint
} from './session-scanner-resume-point'
import {
getSessionSearchIndexMode,
getSessionSearchIndexSink,
withoutSessionSearchCapture,
withSessionSearchCapture,
type SessionSearchIndexSink
} from './session-search-capture'
// Sized past the default recency cap (1000) plus the in-scope cap (2000) so a
// full steady-state result set stays resident between forced rescans.
const MAX_CACHE_ENTRIES = 4096
const NEWLINE_BYTE = 0x0a
type ResumePoint = {
state: ResumableSessionParseState
// Byte offset just past the last complete ('\n'-terminated) line consumed;
// a trailing unterminated line is deliberately left before this point.
byteOffset: number
}
type SessionParseCacheEntry = {
export type SessionParseCacheEntry = {
mtimeMs: number
sizeBytes: number | null
platform: NodeJS.Platform
@@ -80,6 +81,14 @@ function resumableStateFactoryFor(
}
}
export {
invalidateSessionParseCacheEntry,
resetSessionParseCacheForTests,
seedSessionParseCache,
snapshotSessionParseCacheForPersistence,
type PersistedSessionParseCacheEntry
} from './session-parse-cache-store'
export type SessionParseStats = {
reused: number
incremental: number
@@ -95,75 +104,6 @@ export function createSessionParseStats(): SessionParseStats {
return { reused: 0, incremental: 0, fullParses: 0, earlyStopped: 0, bytesRead: 0 }
}
const cache = new Map<string, SessionParseCacheEntry>()
export function resetSessionParseCacheForTests(): void {
cache.clear()
}
// Drops one entry after its file is deleted. Cleanliness, not correctness:
// discovery walks disk first, so a trashed file is never rediscovered anyway.
export function invalidateSessionParseCacheEntry(path: string): void {
cache.delete(path)
}
// Persisted subset of a cache entry: the non-serializable `resume` parser
// state is dropped (see session-parse-cache-persistence.ts).
export type PersistedSessionParseCacheEntry = Omit<SessionParseCacheEntry, 'resume'>
export function snapshotSessionParseCacheForPersistence(): [
string,
PersistedSessionParseCacheEntry
][] {
return [...cache].map(([path, entry]): [string, PersistedSessionParseCacheEntry] => [
path,
{
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session
}
])
}
// Seeded entries carry `resume: null`: after a restart an unchanged file is a
// cache hit; a file that changed while the app was closed pays one full
// (not incremental) re-parse.
export function seedSessionParseCache(
entries: Iterable<[string, PersistedSessionParseCacheEntry]>
): void {
const list = [...entries]
// Snapshot order is oldest→newest (LRU); an over-cap list keeps the newest
// tail rather than seeding the oldest entries and dropping the tail.
for (const [path, entry] of list.slice(Math.max(0, list.length - MAX_CACHE_ENTRIES))) {
if (cache.size >= MAX_CACHE_ENTRIES) {
return
}
// In-process entries are always fresher than persisted ones; never clobber.
if (cache.has(path)) {
continue
}
cache.set(path, {
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session,
resume: null
})
}
}
function storeEntry(path: string, entry: SessionParseCacheEntry): void {
cache.delete(path)
cache.set(path, entry)
if (cache.size > MAX_CACHE_ENTRIES) {
const oldest = cache.keys().next()
if (!oldest.done) {
cache.delete(oldest.value)
}
}
}
/**
* Parse a session file, reusing prior work where the file is provably
* unchanged (mtime+size) and, for append-only JSONL transcripts (Claude,
@@ -179,14 +119,32 @@ export async function parseAgentSessionFileCached(
stats?: SessionParseStats
): Promise<AiVaultSession | null> {
const { file } = candidate
const entry = cache.get(file.path)
const entry = getSessionParseCacheEntry(file.path)
const sink = getSessionSearchIndexSink()
const indexed = sink ? sink.indexedFile(file.path, fileIdentity(file)) : null
const indexCurrent =
sink === null ||
(indexed !== null &&
indexed.mtimeMs === file.mtimeMs &&
(indexed.sizeBytes === null ||
file.sizeBytes === undefined ||
indexed.sizeBytes === file.sizeBytes))
// In `required` mode a file the index has not caught up on is never
// "unchanged"; in `opportunistic` mode it is reused and handed to the backfill.
const indexRequired = sink !== null && getSessionSearchIndexMode() === 'required'
const unchanged =
entry !== undefined &&
entry.platform === platform &&
entry.mtimeMs === file.mtimeMs &&
(entry.sizeBytes === null || file.sizeBytes === undefined || entry.sizeBytes === file.sizeBytes)
(entry.sizeBytes === null ||
file.sizeBytes === undefined ||
entry.sizeBytes === file.sizeBytes) &&
(indexCurrent || !indexRequired)
if (unchanged) {
if (sink && !indexCurrent) {
sink.markStale(candidate)
}
if (stats) {
stats.reused++
}
@@ -214,7 +172,7 @@ export async function parseAgentSessionFileCached(
if (entry.session && candidate.agent === 'codex') {
entry.session = await refreshCachedCodexTitle(candidate, entry.session)
}
storeEntry(file.path, entry)
storeSessionParseCacheEntry(file.path, entry)
return entry.session
}
@@ -225,9 +183,12 @@ export async function parseAgentSessionFileCached(
platform,
entry,
stats,
stateFactory
stateFactory,
sink,
indexRequired,
indexedOffset: indexed?.byteOffset ?? null
})
storeEntry(file.path, parsed)
storeSessionParseCacheEntry(file.path, parsed)
return parsed.session
}
@@ -235,8 +196,10 @@ export async function parseAgentSessionFileCached(
stats.fullParses++
stats.bytesRead += file.sizeBytes ?? 0
}
const session = await parseAgentSessionFile(candidate, platform)
storeEntry(file.path, {
const session = sink
? await indexWholeFileParse(sink, candidate, () => parseAgentSessionFile(candidate, platform))
: await parseAgentSessionFile(candidate, platform)
storeSessionParseCacheEntry(file.path, {
mtimeMs: file.mtimeMs,
sizeBytes: file.sizeBytes ?? null,
platform,
@@ -246,21 +209,52 @@ export async function parseAgentSessionFileCached(
return session
}
async function indexWholeFileParse(
sink: SessionSearchIndexSink,
candidate: SessionFileCandidate,
parse: () => Promise<AiVaultSession | null>
): Promise<AiVaultSession | null> {
const captured = await withSessionSearchCapture(parse)
sink.apply({
candidate,
session: captured.value,
mode: 'replace',
messages: captured.messages,
previousByteOffset: 0,
byteOffset: candidate.file.sizeBytes ?? 0
})
return captured.value
}
async function parseResumableCandidate(args: {
candidate: SessionFileCandidate
platform: NodeJS.Platform
entry: SessionParseCacheEntry | undefined
stats?: SessionParseStats
stateFactory: () => ResumableSessionParseState
sink: SessionSearchIndexSink | null
indexRequired: boolean
indexedOffset: number | null
}): Promise<SessionParseCacheEntry> {
const { file } = args.candidate
const resume = args.entry?.platform === args.platform ? args.entry.resume : null
const canResume =
const parserCanResume =
resume !== null &&
resume !== undefined &&
typeof file.sizeBytes === 'number' &&
file.sizeBytes >= resume.byteOffset &&
sameFileIdentity(resume.identity, fileIdentity(file)) &&
(resume.byteOffset === 0 || (await endsWithNewlineAt(file.path, resume.byteOffset)))
// The index can only take an append that continues from its own offset.
const indexInStep =
args.sink !== null && parserCanResume && args.indexedOffset === resume.byteOffset
const canResume = parserCanResume && (indexInStep || !args.indexRequired)
// Whole-file parses always feed the index (the bytes are read anyway); an
// append feeds it only when in step, otherwise the backfill re-parses later.
const feedIndex = args.sink !== null && (!canResume || indexInStep)
if (args.sink && canResume && !indexInStep) {
args.sink.markStale(args.candidate)
}
// Clone before consuming: a failed read must not corrupt the cached state,
// or the next resume would double-count the lines applied before the error.
@@ -279,15 +273,22 @@ async function parseResumableCandidate(args: {
}
}
const readResult = await consumeCompleteJsonlLines({
path: file.path,
start: startOffset,
onLine: (line) => state.consumeLine(line),
// Bound: the optional hooks are declared as methods, so a parser written
// with method syntax must not lose `this` on the way into the reader.
onLineBytes: state.consumeLineBytes?.bind(state),
shouldStop: state.shouldStop?.bind(state)
})
const read = (): Promise<JsonlReadResult> =>
consumeCompleteJsonlLines({
path: file.path,
start: startOffset,
onLine: (line) => state.consumeLine(line),
// Bound: the optional hooks are declared as methods, so a parser written
// with method syntax must not lose `this` on the way into the reader.
onLineBytes: state.consumeLineBytes?.bind(state),
shouldStop: state.shouldStop?.bind(state)
})
// Capture only when the index takes these rows: it disables the Codex
// byte-prefix fast path, which is exactly the cost a list-only scan must not pay.
const captured = feedIndex
? await withSessionSearchCapture(read)
: { value: await read(), messages: [] }
const readResult = captured.value
if (args.stats) {
args.stats.bytesRead += readResult.bytesRead
}
@@ -297,28 +298,33 @@ async function parseResumableCandidate(args: {
// Keep parity with the one-shot parser: a final unterminated line is shown,
// but stays out of the resumable state so the (possibly still-growing) line
// is re-read once complete instead of being half-counted.
// is re-read once complete instead of being half-counted. The index only
// stores complete lines, so the display-only consume must not emit rows.
let displayState = state
if (readResult.trailingPartialLine !== null) {
displayState = state.clone()
displayState.consumeLine(readResult.trailingPartialLine)
withoutSessionSearchCapture(() => displayState.consumeLine(readResult.trailingPartialLine!))
}
const session = await displayState.finalize(args.platform)
if (feedIndex && args.sink) {
args.sink.apply({
candidate: args.candidate,
// The index stores what the complete lines say; the trailing partial line
// only changes the displayed session.
session: displayState === state ? session : await state.finalize(args.platform),
mode: canResume ? 'append' : 'replace',
messages: captured.messages,
previousByteOffset: startOffset,
byteOffset: readResult.consumedThrough
})
}
return {
mtimeMs: file.mtimeMs,
sizeBytes: file.sizeBytes ?? null,
platform: args.platform,
session: await displayState.finalize(args.platform),
resume: { state, byteOffset: readResult.consumedThrough }
session,
resume: { state, byteOffset: readResult.consumedThrough, identity: fileIdentity(file) }
}
}
// A resume point is only valid if it still sits just past a line break;
// anything else means the file was rewritten, not appended. Heuristic: a
// grown rewrite keeping '\n' at exactly this byte would slip through, but
// agent transcripts are append-only so that trade is accepted (worst case is
// a stale vault row until the file is next truncated or the app restarts).
async function endsWithNewlineAt(path: string, offset: number): Promise<boolean> {
const slice = await readTranscriptSlice(path, offset - 1, 1, 'scan')
return slice.length === 1 && slice[0] === NEWLINE_BYTE
}
@@ -0,0 +1,45 @@
import { readTranscriptSlice } from '../native-chat/wsl-transcript-fs-access'
import type { FileWithMtime, ResumableSessionParseState } from './session-scanner-types'
const NEWLINE_BYTE = 0x0a
export type FileIdentity = { dev: number; ino: number }
export type ResumePoint = {
state: ResumableSessionParseState
// Byte offset just past the last complete ('\n'-terminated) line consumed;
// a trailing unterminated line is deliberately left before this point.
byteOffset: number
// Inode the offset belongs to; null when discovery could not stat identity.
identity: FileIdentity | null
}
export function fileIdentity(file: FileWithMtime): FileIdentity | null {
return typeof file.dev === 'number' && typeof file.ino === 'number'
? { dev: file.dev, ino: file.ino }
: null
}
// A rename-replace (atomic rewrite) keeps the path and can keep '\n' at the
// old offset, so the newline guard alone would resume into a different file.
// Identity is only compared when both sides have it; a candidate without a
// stat (synthetic files) keeps the newline-only heuristic.
export function sameFileIdentity(
previous: FileIdentity | null,
current: FileIdentity | null
): boolean {
if (previous === null || current === null) {
return true
}
return previous.dev === current.dev && previous.ino === current.ino
}
// A resume point is only valid if it still sits just past a line break;
// anything else means the file was rewritten in place, not appended. An
// in-place grown rewrite keeping '\n' at exactly this byte still slips
// through (rename-replace is caught by the inode check above); agent
// transcripts are append-only so that residual trade is accepted.
export async function endsWithNewlineAt(path: string, offset: number): Promise<boolean> {
const slice = await readTranscriptSlice(path, offset - 1, 1, 'scan')
return slice.length === 1 && slice[0] === NEWLINE_BYTE
}
@@ -19,6 +19,7 @@ import {
import { readAiVaultSessionTitlesFromFiles } from './session-title-file-reader'
import { resolveHostReadableAiVaultTitleRequests } from './session-title-request-paths'
import { listLocalAiVaultSubagentSessions } from './session-subagent-reader'
import { SessionSearchService } from '../ai-vault-search/session-search-service'
if (!process.send) {
throw new Error('AI Vault service requires a parent IPC channel.')
@@ -29,6 +30,7 @@ const cancelled = new Set<number>()
const pending = new Set<number>()
const titleIndex = new Map<string, AiVaultSessionTitle>()
const invalidatedPaths = new Set<string>()
let sessionSearch: SessionSearchService | null = null
let initialized = false
let shuttingDown = false
let cacheLane = Promise.resolve()
@@ -42,6 +44,13 @@ function titleKey(request: { agent: string; sessionId: string }): string {
return `${request.agent}\0${request.sessionId}`
}
function requireSessionSearch(): SessionSearchService {
if (!sessionSearch) {
throw new Error('Agent session search is not enabled on this host.')
}
return sessionSearch
}
async function executeRequest(request: AiVaultServiceRequest): Promise<AiVaultServiceResultValue> {
const controller = new AbortController()
controllers.set(request.id, controller)
@@ -77,6 +86,22 @@ async function executeRequest(request: AiVaultServiceRequest): Promise<AiVaultSe
value: await readAiVaultFirstUserPrompt(request.request)
}
}
if (request.operation === 'search') {
return {
operation: 'search',
value: await requireSessionSearch().search(
request.request.args,
request.request.roots,
controller.signal
)
}
}
if (request.operation === 'searchCoverage') {
return {
operation: 'searchCoverage',
value: requireSessionSearch().coverage(request.request.roots)
}
}
const startedAt = performance.now()
const result = await scanAiVaultSessions({ ...request.options, signal: controller.signal })
for (const session of result.sessions) {
@@ -151,6 +176,7 @@ async function shutdown(): Promise<void> {
}
await Promise.allSettled([cacheLane, interactiveLane])
await flushSessionParseCachePersist()
sessionSearch?.dispose()
process.disconnect?.()
}
@@ -164,6 +190,13 @@ process.on('message', (raw: AiVaultServiceParentMessage) => {
if (raw.sessionParseCache) {
initSessionParseCachePersistence(raw.sessionParseCache)
}
if (raw.sessionSearch) {
try {
sessionSearch = new SessionSearchService(raw.sessionSearch)
} catch (error) {
console.error('[ai-vault] session search index unavailable:', error)
}
}
send({ type: 'ready', protocol: AI_VAULT_SERVICE_PROTOCOL_VERSION, pid: process.pid })
return
}
@@ -176,6 +209,7 @@ process.on('message', (raw: AiVaultServiceParentMessage) => {
return
}
if (raw?.type === 'invalidate') {
sessionSearch?.invalidate(raw.paths)
for (const path of raw.paths) {
invalidatedPaths.delete(path)
invalidatedPaths.add(path)
@@ -5,23 +5,43 @@ import type {
AiVaultSessionTitlesResult
} from '../../shared/ai-vault-session-title'
import type { ReadAiVaultFirstUserPromptArgs } from './session-first-user-prompt-read'
import type {
AiVaultSearchArgs,
AiVaultSearchCoverage,
AiVaultSearchResult
} from '../../shared/ai-vault-search-types'
import type { SessionParseCachePersistenceOptions } from './session-parse-cache-persistence'
import type { AiVaultWorkerScanOptions } from './session-scanner-worker-protocol'
export const AI_VAULT_SERVICE_PROTOCOL_VERSION = 1
export type AiVaultServiceLane = 'cache' | 'interactive'
export type AiVaultServiceOperation = 'scan' | 'titles' | 'subagents' | 'firstPrompt'
export type AiVaultServiceOperation =
| 'scan'
| 'titles'
| 'subagents'
| 'firstPrompt'
| 'search'
| 'searchCoverage'
export type AiVaultServiceSubagentRequest = {
agent: 'claude' | 'omp'
parentFilePath: string
}
export type AiVaultSessionSearchInit = { databasePath: string }
export type AiVaultServiceInit = {
type: 'init'
protocol: typeof AI_VAULT_SERVICE_PROTOCOL_VERSION
sessionParseCache: SessionParseCachePersistenceOptions | null
sessionSearch?: AiVaultSessionSearchInit | null
}
/** Search requests carry the scan roots so the child's backfill sees what list scans see. */
export type AiVaultServiceSearchRequest = {
args: AiVaultSearchArgs
roots: Omit<AiVaultWorkerScanOptions, 'limit' | 'unlimited' | 'scopePaths'>
}
export type AiVaultServiceRequestBody =
@@ -41,6 +61,12 @@ export type AiVaultServiceRequestBody =
operation: 'firstPrompt'
request: ReadAiVaultFirstUserPromptArgs
}
| { type: 'request'; operation: 'search'; request: AiVaultServiceSearchRequest }
| {
type: 'request'
operation: 'searchCoverage'
request: Pick<AiVaultServiceSearchRequest, 'roots'>
}
export type AiVaultServiceRequest = AiVaultServiceRequestBody & { id: number }
@@ -56,6 +82,8 @@ export type AiVaultServiceResultValue =
| { operation: 'titles'; value: AiVaultSessionTitlesResult }
| { operation: 'subagents'; value: AiVaultSubagentListResult }
| { operation: 'firstPrompt'; value: { prompt: string | null } }
| { operation: 'search'; value: AiVaultSearchResult }
| { operation: 'searchCoverage'; value: AiVaultSearchCoverage }
export type AiVaultServiceChildMessage =
| {
@@ -68,7 +96,7 @@ export type AiVaultServiceChildMessage =
| { type: 'invalidated'; generation: number }
export function aiVaultServiceLane(operation: AiVaultServiceOperation): AiVaultServiceLane {
return operation === 'subagents' || operation === 'firstPrompt' ? 'interactive' : 'cache'
return operation === 'scan' || operation === 'titles' ? 'cache' : 'interactive'
}
export function isAiVaultServiceRequest(value: unknown): value is AiVaultServiceRequest {
@@ -82,7 +110,9 @@ export function isAiVaultServiceRequest(value: unknown): value is AiVaultService
(message.operation === 'scan' ||
message.operation === 'titles' ||
message.operation === 'subagents' ||
message.operation === 'firstPrompt')
message.operation === 'firstPrompt' ||
message.operation === 'search' ||
message.operation === 'searchCoverage')
)
}
@@ -15,8 +15,13 @@ import { buildAiVaultServiceEnv } from './session-scanner-service-env'
import { AiVaultScannerServiceClient } from './session-scanner-service-client'
import { getAiVaultServiceEntryPath } from './session-scanner-service-entry-path'
import { lowerAiVaultServicePriority } from './session-scanner-service-priority'
import type { AiVaultServiceSubagentRequest } from './session-scanner-service-protocol'
import type {
AiVaultServiceSearchRequest,
AiVaultServiceSubagentRequest
} from './session-scanner-service-protocol'
import type { AiVaultWorkerScanOptions } from './session-scanner-worker-protocol'
import { getSessionSearchInitOptions } from '../ai-vault-search/session-search-paths'
import type { AiVaultSearchCoverage, AiVaultSearchResult } from '../../shared/ai-vault-search-types'
export function spawnAiVaultServiceProcess(): ChildProcess {
const entryPath = getAiVaultServiceEntryPath()
@@ -39,7 +44,10 @@ let sharedClient: AiVaultScannerServiceClient | null = null
function getSharedClient(): AiVaultScannerServiceClient {
sharedClient ??= new AiVaultScannerServiceClient({
processFactory: spawnAiVaultServiceProcess,
init: { sessionParseCache: getSessionParseCachePersistenceOptions() },
init: {
sessionParseCache: getSessionParseCachePersistenceOptions(),
sessionSearch: getSessionSearchInitOptions()
},
onStderr: (text) => console.error('[ai-vault-service]', text.trimEnd())
})
return sharedClient
@@ -81,6 +89,23 @@ export function readAiVaultFirstUserPromptInService(
return getSharedClient().request({ type: 'request', operation: 'firstPrompt', request }, signal)
}
export function searchAiVaultSessionsInService(
request: AiVaultServiceSearchRequest,
signal?: AbortSignal
): Promise<AiVaultSearchResult> {
return getSharedClient().request({ type: 'request', operation: 'search', request }, signal)
}
export function readAiVaultSearchCoverageInService(
request: Pick<AiVaultServiceSearchRequest, 'roots'>,
signal?: AbortSignal
): Promise<AiVaultSearchCoverage> {
return getSharedClient().request(
{ type: 'request', operation: 'searchCoverage', request },
signal
)
}
export function invalidateAiVaultServiceCache(paths: string[]): Promise<void> {
return sharedClient?.invalidate(paths) ?? Promise.resolve()
}
@@ -0,0 +1,38 @@
import { describe, expect, it } from 'vitest'
import { extractPreviewContentText } from './session-scanner-text-normalization'
describe('extractPreviewContentText', () => {
it('renders a tool_use block as the tool name and its command', () => {
expect(
extractPreviewContentText([
{ type: 'tool_use', id: 'toolu_1', name: 'Bash', input: { command: 'pnpm test src/foo' } }
])
).toBe('Bash: pnpm test src/foo')
})
it('falls back through file_path, pattern, and description inputs', () => {
expect(
extractPreviewContentText([
{ type: 'tool_use', name: 'Read', input: { file_path: '/a/b.ts' } }
])
).toBe('Read: /a/b.ts')
expect(
extractPreviewContentText([{ type: 'tool_use', name: 'Grep', input: { pattern: 'foo\\(' } }])
).toBe('Grep: foo\\(')
expect(extractPreviewContentText([{ type: 'tool_use', name: 'Task', input: {} }])).toBe('Task')
})
it('interleaves tool calls with surrounding text and bounds the argument', () => {
const long = 'x'.repeat(2000)
const text = extractPreviewContentText([
{ type: 'text', text: 'Running tests.' },
{ type: 'tool_use', name: 'Bash', input: { command: long } }
])
expect(text?.startsWith('Running tests. Bash: xxx')).toBe(true)
expect(text?.length).toBeLessThan(long.length)
})
it('ignores a tool_use block with no name and no usable input', () => {
expect(extractPreviewContentText([{ type: 'tool_use', input: { unrelated: 1 } }])).toBeNull()
})
})
@@ -95,9 +95,35 @@ function contentItemText(item: unknown): string | null {
return null
}
if (record.type === 'tool_use') {
return toolUseText(record)
}
return nonBlankString(record.text) ?? nonBlankString(record.content)
}
// Why: a tool call names the command the agent ran; without it an
// assistant turn that only invokes tools reads as empty, and "which session
// ran this" is unanswerable. Bounded so a huge Write payload cannot dominate.
const TOOL_USE_INPUT_SCAN_LIMIT = 512
function toolUseText(record: Record<string, unknown>): string | null {
const name = nonBlankString(record.name)
const input = objectRecord(record.input)
const argument = input
? (nonBlankString(input.command) ??
nonBlankString(input.file_path) ??
nonBlankString(input.path) ??
nonBlankString(input.pattern) ??
nonBlankString(input.query) ??
nonBlankString(input.description))
: null
if (!name && !argument) {
return null
}
const argumentText = argument ? sliceAtCodeUnitLimit(argument, TOOL_USE_INPUT_SCAN_LIMIT) : null
return name && argumentText ? `${name}: ${argumentText}` : (name ?? argumentText)
}
function nonBlankString(value: unknown): string | null {
if (typeof value !== 'string') {
return null
@@ -10,16 +10,22 @@ import type {
AiVaultWorkerResponse,
AiVaultWorkerScanOptions
} from './session-scanner-worker-protocol'
import type { AiVaultServiceSearchRequest } from './session-scanner-service-protocol'
import type { AiVaultSearchCoverage, AiVaultSearchResult } from '../../shared/ai-vault-search-types'
const SCAN_TIMEOUT_MS = 130_000
const TITLE_TIMEOUT_MS = 15_000
// Why: the first search may fold a burst of stale files before answering.
const SEARCH_TIMEOUT_MS = 60_000
const MAX_QUEUED_CALLS = 16
export type AiVaultWorkerFactory = () => Worker
type RequestBody =
| Omit<Extract<AiVaultWorkerRequest, { kind: 'scan' }>, 'id'>
| Omit<Extract<AiVaultWorkerRequest, { kind: 'titles' }>, 'id'>
type RequestBody = AiVaultWorkerRequest extends infer R
? R extends { id: number }
? Omit<R, 'id'>
: never
: never
type PendingCall = {
request: AiVaultWorkerRequest
@@ -64,6 +70,23 @@ export class AiVaultScannerWorkerClient {
) as Promise<AiVaultSessionTitlesResult>
}
search(request: AiVaultServiceSearchRequest, signal?: AbortSignal): Promise<AiVaultSearchResult> {
return this.dispatch(
{ kind: 'search', request },
SEARCH_TIMEOUT_MS,
signal
) as Promise<AiVaultSearchResult>
}
searchCoverage(
request: Pick<AiVaultServiceSearchRequest, 'roots'>
): Promise<AiVaultSearchCoverage> {
return this.dispatch(
{ kind: 'searchCoverage', request },
TITLE_TIMEOUT_MS
) as Promise<AiVaultSearchCoverage>
}
dispose(): void {
this.destroyWorker()
const pending = this.queue
@@ -4,6 +4,7 @@ import type {
AiVaultSessionTitleRequest
} from '../../shared/ai-vault-session-title'
import { scanAiVaultSessions } from './session-scanner'
import { SessionSearchService } from '../ai-vault-search/session-search-service'
import { initSessionParseCachePersistence } from './session-parse-cache-persistence'
import { readAiVaultSessionTitlesFromFiles } from './session-title-file-reader'
import { resolveHostReadableAiVaultTitleRequests } from './session-title-request-paths'
@@ -24,6 +25,7 @@ const data = workerData as AiVaultWorkerData | undefined
if (data?.sessionParseCache) {
initSessionParseCachePersistence(data.sessionParseCache)
}
const sessionSearch = data?.sessionSearch ? new SessionSearchService(data.sessionSearch) : null
const controllers = new Map<number, AbortController>()
const titleIndex = new Map<string, AiVaultSessionTitle>()
@@ -66,6 +68,28 @@ async function handleRequest(request: AiVaultWorkerRequest): Promise<AiVaultWork
})
}
}
if (request.kind === 'search' || request.kind === 'searchCoverage') {
if (!sessionSearch) {
throw new Error('Agent session search is not enabled on this host.')
}
return request.kind === 'search'
? {
id: request.id,
ok: true,
kind: 'search',
value: await sessionSearch.search(
request.request.args,
request.request.roots,
controller.signal
)
}
: {
id: request.id,
ok: true,
kind: 'searchCoverage',
value: sessionSearch.coverage(request.request.roots)
}
}
const startedAt = performance.now()
const result = await scanAiVaultSessions({ ...request.options, signal: controller.signal })
for (const session of result.sessions) {
@@ -5,16 +5,24 @@ import type {
} from '../../shared/ai-vault-session-title'
import type { AiVaultScanOptions } from './session-scanner-types'
import type { SessionParseCachePersistenceOptions } from './session-parse-cache-persistence'
import type { AiVaultSearchCoverage, AiVaultSearchResult } from '../../shared/ai-vault-search-types'
import type {
AiVaultServiceSearchRequest,
AiVaultSessionSearchInit
} from './session-scanner-service-protocol'
export type AiVaultWorkerScanOptions = Omit<AiVaultScanOptions, 'signal'>
export type AiVaultWorkerData = {
sessionParseCache: SessionParseCachePersistenceOptions | null
sessionSearch?: AiVaultSessionSearchInit | null
}
export type AiVaultWorkerRequest =
| { id: number; kind: 'scan'; options: AiVaultWorkerScanOptions }
| { id: number; kind: 'titles'; requests: AiVaultSessionTitleRequest[] }
| { id: number; kind: 'search'; request: AiVaultServiceSearchRequest }
| { id: number; kind: 'searchCoverage'; request: Pick<AiVaultServiceSearchRequest, 'roots'> }
export type AiVaultWorkerControl = { id: number; kind: 'cancel' }
@@ -26,4 +34,6 @@ export type AiVaultWorkerResponse =
value: { result: AiVaultListResult; durationMs: number }
}
| { id: number; ok: true; kind: 'titles'; value: AiVaultSessionTitlesResult }
| { id: number; ok: true; kind: 'search'; value: AiVaultSearchResult }
| { id: number; ok: true; kind: 'searchCoverage'; value: AiVaultSearchCoverage }
| { id: number; ok: false; error: string }
@@ -10,6 +10,9 @@ import { withSpan } from '../observability/tracer'
import { getSessionParseCachePersistenceOptions } from './session-parse-cache-persistence'
import { AiVaultScannerWorkerClient } from './session-scanner-worker-client'
import type { AiVaultWorkerData, AiVaultWorkerScanOptions } from './session-scanner-worker-protocol'
import type { AiVaultServiceSearchRequest } from './session-scanner-service-protocol'
import { getSessionSearchInitOptions } from '../ai-vault-search/session-search-paths'
import type { AiVaultSearchCoverage, AiVaultSearchResult } from '../../shared/ai-vault-search-types'
const WORKER_ENTRY_FILENAME = 'session-scanner-worker-entry.js'
@@ -20,7 +23,8 @@ function defaultWorkerFactory(): Worker {
}
return new Worker(workerPath, {
workerData: {
sessionParseCache: getSessionParseCachePersistenceOptions()
sessionParseCache: getSessionParseCachePersistenceOptions(),
sessionSearch: getSessionSearchInitOptions()
} satisfies AiVaultWorkerData
})
}
@@ -51,6 +55,19 @@ export function resolveAiVaultSessionTitlesInWorker(
return getSharedClient().resolveTitles(requests, signal)
}
export function searchAiVaultSessionsInWorker(
request: AiVaultServiceSearchRequest,
signal?: AbortSignal
): Promise<AiVaultSearchResult> {
return getSharedClient().search(request, signal)
}
export function readAiVaultSearchCoverageInWorker(
request: Pick<AiVaultServiceSearchRequest, 'roots'>
): Promise<AiVaultSearchCoverage> {
return getSharedClient().searchCoverage(request)
}
export function resetAiVaultScannerWorkerForTests(): void {
sharedClient?.dispose()
sharedClient = null
+4 -41
View File
@@ -6,19 +6,13 @@ import type {
import { LOCAL_EXECUTION_HOST_ID, type ExecutionHostId } from '../../shared/execution-host'
import { withSpan } from '../observability/tracer'
import { sessionSortTime } from './session-scanner-accumulator'
import {
codexRolloutHardlinkIdentity,
dedupeCodexRolloutAliases,
dedupeCodexSessionsBySessionId
} from './codex-session-root-dedup'
import { readCodexRolloutSessionMetaId } from '../codex/codex-rollout-session-meta'
import { dedupeCodexSessionsBySessionId } from './codex-session-root-dedup'
import { sessionCandidatesFromDiscoveries } from './session-scanner-candidates'
import {
createAntigravityWorkspaceResolver,
readLocalAntigravityHistory,
type AntigravityWorkspaceResolver
} from './session-scanner-antigravity-history'
import { antigravityHistoryPathForBrainDir } from './session-scanner-antigravity-paths'
import { codexHomeForSessionsDir } from './session-scanner-codex-paths'
import {
ensureSessionParseCacheLoaded,
scheduleSessionParseCachePersist
@@ -30,10 +24,7 @@ import {
} from './session-scanner-parse-cache'
import { recordSessionScanIssue } from './session-scan-issues'
import { discoverInScopeClaudeFiles } from './session-scanner-scope-discovery'
import {
DEFAULT_CODEX_HOME_DIR,
discoverAiVaultSessionSources
} from './session-scanner-source-discovery'
import { discoverAiVaultSessionSources } from './session-scanner-source-discovery'
import type {
AiVaultScanOptions,
SessionFileCandidate,
@@ -83,35 +74,7 @@ export async function scanAiVaultSessions(
const discoveries = await discoverAiVaultSessionSources({ options, limitPerAgent, issues })
throwIfAiVaultScanCancelled(options.signal)
const candidates = await dedupeCodexRolloutAliases(
discoveries
.flatMap((discovery) =>
discovery.files.map((file): SessionFileCandidate => ({
agent: discovery.agent,
file,
codexHome:
discovery.agent === 'codex'
? codexHomeForSessionsDir(
discovery.rootDir,
options.defaultCodexHomeDir ?? DEFAULT_CODEX_HOME_DIR
)
: null,
antigravityHistoryPath:
discovery.agent === 'antigravity'
? antigravityHistoryPathForBrainDir(discovery.rootDir)
: undefined
}))
)
.sort((left, right) => right.file.mtimeMs - left.file.mtimeMs),
{
isCodex: (candidate) => candidate.agent === 'codex',
getFilePath: (candidate) => candidate.file.path,
getCodexHome: (candidate) => candidate.codexHome,
getHardlinkIdentity: (candidate) => codexRolloutHardlinkIdentity(candidate.file)
},
(filePath) => readCodexRolloutSessionMetaId(filePath, options.signal, 'scan'),
options.signal
)
const candidates = await sessionCandidatesFromDiscoveries(discoveries, options)
const parsedSessions = await parseSessionCandidates({
candidates: candidates.slice(0, limit * SESSION_PARSE_CANDIDATE_MULTIPLIER),
@@ -0,0 +1,99 @@
import { AsyncLocalStorage } from 'node:async_hooks'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import type { SessionFileCandidate } from './session-scanner-types'
// Why: the parsers already fold every provider's transcript into one
// accumulator. Instead of a second reader per format, a parse runs inside a
// capture scope and the preview funnel also emits full-text rows; incremental
// resumes emit only the newly consumed lines, which is exactly what the search
// index needs to append.
export type SessionSearchCapturedRole = 'user' | 'assistant' | 'tool'
export type SessionSearchCapturedMessage = {
role: SessionSearchCapturedRole
text: string
timestamp: string | null
}
export type SessionSearchIndexUpdate = {
candidate: SessionFileCandidate
/** Null when the parser rejected the file (e.g. a Codex worker transcript): drop its rows. */
session: AiVaultSession | null
/** `replace`: whole-file parse, rows supersede the session; `append`: resumed parse. */
mode: 'replace' | 'append'
messages: SessionSearchCapturedMessage[]
/** Byte offset the appended rows continue from; the sink refuses a mismatch. */
previousByteOffset: number
byteOffset: number
}
export type SessionSearchFileIdentity = { dev: number; ino: number } | null
export type SessionSearchIndexedFile = {
byteOffset: number
mtimeMs: number
sizeBytes: number | null
}
export type SessionSearchIndexSink = {
/**
* What the index holds for this file, or null when it is not indexed or its
* identity changed. In `required` mode the parse cache may only reuse an
* entry the index also has and may only resume when the parser's resume
* offset equals `byteOffset`; anything else forces a whole-file parse.
*/
indexedFile(path: string, identity: SessionSearchFileIdentity): SessionSearchIndexedFile | null
/** Never throws: an index failure must not break the session list. */
apply(update: SessionSearchIndexUpdate): void
/** `opportunistic` mode saw a file the index is behind on; the backfill lane re-parses it. */
markStale(candidate: SessionFileCandidate): void
}
// Why: list scans have a latency budget and must never pay for the index; they
// feed it only when a whole-file parse happens anyway. The backfill lane runs
// in `required` mode, where index consistency wins over parse reuse.
export type SessionSearchIndexMode = 'opportunistic' | 'required'
type CaptureScope = { messages: SessionSearchCapturedMessage[] } | null
const captureStorage = new AsyncLocalStorage<CaptureScope>()
const indexModeStorage = new AsyncLocalStorage<SessionSearchIndexMode>()
let sink: SessionSearchIndexSink | null = null
export function getSessionSearchIndexMode(): SessionSearchIndexMode {
return indexModeStorage.getStore() ?? 'opportunistic'
}
export function withSessionSearchIndexRequired<T>(fn: () => Promise<T>): Promise<T> {
return indexModeStorage.run('required', fn)
}
export function registerSessionSearchIndexSink(next: SessionSearchIndexSink | null): void {
sink = next
}
export function getSessionSearchIndexSink(): SessionSearchIndexSink | null {
return sink
}
export function captureSessionSearchMessage(message: SessionSearchCapturedMessage): void {
captureStorage.getStore()?.messages.push(message)
}
export function isSessionSearchCaptureActive(): boolean {
return captureStorage.getStore() != null
}
/** Runs `fn` with capture suppressed: display-only re-reads must not emit rows. */
export function withoutSessionSearchCapture<T>(fn: () => T): T {
return captureStorage.run(null, fn)
}
export async function withSessionSearchCapture<T>(
fn: () => Promise<T>
): Promise<{ value: T; messages: SessionSearchCapturedMessage[] }> {
const scope = { messages: [] as SessionSearchCapturedMessage[] }
const value = await captureStorage.run(scope, fn)
return { value, messages: scope.messages }
}
@@ -0,0 +1,73 @@
import { asRecord, extractString } from './session-scanner-values'
import { captureIndexableText, toolCallText } from './session-search-content'
import { isSessionSearchCaptureActive } from './session-search-capture'
// Codex writes tool traffic twice: the raw model call/output as response_item
// records, and a rendered CommandExecution/FileChange item once the turn
// completes. Index the rendered item when the history is paginated (it has the
// real argv and output) and the raw pair otherwise, so nothing is stored twice.
const RAW_CALL_TYPES = new Set(['function_call', 'custom_tool_call', 'local_shell_call'])
const RAW_OUTPUT_TYPES = new Set(['function_call_output', 'custom_tool_call_output'])
function argvText(command: unknown): string | null {
if (typeof command === 'string') {
return command
}
return Array.isArray(command)
? command.filter((part): part is string => typeof part === 'string').join(' ')
: null
}
export function captureCodexToolRecord(
recordType: unknown,
payload: Record<string, unknown>,
timestamp: unknown,
historyMode: string | null
): void {
if (!isSessionSearchCaptureActive()) {
return
}
const payloadType = extractString(payload.type)
if (recordType === 'response_item' && payloadType) {
if (historyMode === 'paginated') {
return
}
if (RAW_CALL_TYPES.has(payloadType)) {
const input = payload.arguments ?? payload.input ?? asRecord(payload.action)?.command
captureIndexableText('tool', toolCallText(payload.name ?? 'shell', input), timestamp)
} else if (RAW_OUTPUT_TYPES.has(payloadType)) {
captureIndexableText('tool', extractString(payload.output), timestamp)
}
return
}
if (recordType !== 'event_msg' || payloadType !== 'item_completed') {
return
}
const item = asRecord(payload.item)
const itemType = item ? extractString(item.type) : null
if (!item || !itemType) {
return
}
if (itemType === 'CommandExecution') {
captureIndexableText('tool', argvText(item.command), timestamp)
captureIndexableText(
'tool',
extractString(item.aggregated_output) ?? extractString(item.stdout),
timestamp
)
return
}
if (itemType === 'FileChange') {
const changes = asRecord(item.changes)
const paths = changes ? Object.keys(changes).join(' ') : null
captureIndexableText(
'tool',
paths ? `FileChange: ${paths}` : extractString(item.stdout),
timestamp
)
return
}
if (itemType === 'Reasoning') {
captureIndexableText('assistant', extractString(item.text), timestamp)
}
}
+181
View File
@@ -0,0 +1,181 @@
import {
captureSessionSearchMessage,
isSessionSearchCaptureActive,
type SessionSearchCapturedRole
} from './session-search-capture'
import { timestampMs } from './session-scanner-values'
// Why: literal accuracy tracks this cap almost linearly (MRR 0.36 → 0.62 from
// none to 3 KB) while latency and disk scale the other way; 3 KB is the knee.
export const SESSION_SEARCH_TOOL_OUTPUT_CAP = 3000
// Sanity bound on a single conversational message.
const SESSION_SEARCH_TEXT_CAP = 200_000
const TOOL_INPUT_CAP = 2000
const HIDDEN_BLOCK_PATTERN =
/<(system-reminder|codex_internal_context|goal_context)\b[^>]*>[\s\S]*?<\/\1>/gi
const TEXT_BLOCK_TYPES = new Set(['text', 'input_text', 'output_text', 'thinking', 'reasoning'])
const TOOL_INPUT_KEYS = ['command', 'cmd', 'file_path', 'path', 'pattern', 'query', 'description']
type PreviewRole = 'user' | 'assistant' | 'system' | 'tool' | 'unknown'
function record(value: unknown): Record<string, unknown> | null {
return value && typeof value === 'object' && !Array.isArray(value)
? (value as Record<string, unknown>)
: null
}
function stripHiddenBlocks(text: string): string {
return text.includes('<') ? text.replace(HIDDEN_BLOCK_PATTERN, ' ') : text
}
function cap(text: string, limit: number): string {
return text.length > limit ? text.slice(0, limit) : text
}
function indexRole(role: PreviewRole): SessionSearchCapturedRole | null {
return role === 'user' || role === 'assistant' || role === 'tool' ? role : null
}
/** Flattens a tool_result body (string, or array of text blocks) to one string. */
function toolResultText(content: unknown): string {
if (typeof content === 'string') {
return content
}
if (!Array.isArray(content)) {
return ''
}
const parts: string[] = []
let length = 0
for (const item of content) {
const text = typeof item === 'string' ? item : record(item)?.text
if (typeof text === 'string' && text) {
parts.push(text)
length += text.length
if (length >= SESSION_SEARCH_TOOL_OUTPUT_CAP) {
break
}
}
}
return parts.join('\n')
}
export function toolCallText(name: unknown, input: unknown): string | null {
const toolName = typeof name === 'string' && name ? name : null
const inputRecord = record(input)
let argument: string | null = null
if (inputRecord) {
for (const key of TOOL_INPUT_KEYS) {
const value = inputRecord[key]
if (typeof value === 'string' && value.trim()) {
argument = value
break
}
}
} else if (typeof input === 'string' && input.trim()) {
argument = input
}
if (!toolName && !argument) {
return null
}
const bounded = argument ? cap(argument, TOOL_INPUT_CAP) : null
return toolName && bounded ? `${toolName}: ${bounded}` : (toolName ?? bounded)
}
type IndexableMessage = { role: SessionSearchCapturedRole; text: string }
/**
* Splits a provider content value into indexable messages. Text blocks keep
* the record's role; tool_use and tool_result blocks become `tool` rows no
* matter which record carried them (Claude stores tool results on user records).
*/
export function indexableMessagesFromContent(
role: PreviewRole,
content: unknown
): IndexableMessage[] {
const out: IndexableMessage[] = []
const textRole = indexRole(role)
if (typeof content === 'string') {
if (textRole) {
out.push({ role: textRole, text: content })
}
return out
}
const blocks = Array.isArray(content) ? content : content != null ? [content] : []
const textParts: string[] = []
for (const block of blocks) {
if (typeof block === 'string') {
textParts.push(block)
continue
}
const item = record(block)
if (!item) {
continue
}
const type = typeof item.type === 'string' ? item.type : null
if (type === 'tool_use') {
const text = toolCallText(item.name, item.input)
if (text) {
out.push({ role: 'tool', text })
}
continue
}
if (type === 'tool_result') {
const text = toolResultText(item.content)
if (text.trim()) {
out.push({ role: 'tool', text })
}
continue
}
if (type !== null && !TEXT_BLOCK_TYPES.has(type)) {
continue
}
const text = typeof item.text === 'string' ? item.text : item.content
if (typeof text === 'string' && text) {
textParts.push(text)
}
}
if (textRole && textParts.length > 0) {
out.unshift({ role: textRole, text: textParts.join('\n') })
}
return out
}
/** Emits index rows for a message when a capture scope is active; no-op otherwise. */
export function captureIndexableContent(
role: PreviewRole,
content: unknown,
timestamp: unknown
): void {
if (!isSessionSearchCaptureActive()) {
return
}
for (const message of indexableMessagesFromContent(role, content)) {
captureIndexableText(message.role, message.text, timestamp)
}
}
export function captureIndexableText(
role: PreviewRole,
text: string | null,
timestamp: unknown
): void {
if (!text || !isSessionSearchCaptureActive()) {
return
}
const indexed = indexRole(role)
if (!indexed) {
return
}
const limit = indexed === 'tool' ? SESSION_SEARCH_TOOL_OUTPUT_CAP : SESSION_SEARCH_TEXT_CAP
const cleaned = cap(stripHiddenBlocks(cap(text, limit * 4)), limit).trim()
if (!cleaned) {
return
}
const parsed = timestampMs(timestamp)
captureSessionSearchMessage({
role: indexed,
text: cleaned,
timestamp: Number.isFinite(parsed) ? new Date(parsed).toISOString() : null
})
}
+32
View File
@@ -3,6 +3,10 @@ import { defineMethod, type RpcMethod } from '../core'
import { OptionalBoolean } from '../schemas'
import { restampAiVaultListResult } from '../../../ai-vault/session-list-results'
import { AI_VAULT_AGENTS, AI_VAULT_SCOPE_PATHS_MAX_COUNT } from '../../../../shared/ai-vault-types'
import {
AI_VAULT_SEARCH_LIMIT_MAX,
AI_VAULT_SEARCH_QUERY_MAX_LENGTH
} from '../../../../shared/ai-vault-search-types'
import { AI_VAULT_SESSION_TITLE_REQUEST_MAX_COUNT } from '../../../../shared/ai-vault-session-title'
import type { AiVaultPrepareSessionResumeArgs } from '../../../../shared/ai-vault-resume-preparation'
import { LOCAL_EXECUTION_HOST_ID, parseExecutionHostId } from '../../../../shared/execution-host'
@@ -80,7 +84,35 @@ export const AiVaultSessionTitlesParams = z.object({
.max(AI_VAULT_SESSION_TITLE_REQUEST_MAX_COUNT)
})
export const AiVaultSearchSessionsParams = z.object({
query: z.string().trim().min(1).max(AI_VAULT_SEARCH_QUERY_MAX_LENGTH),
limit: z.number().int().min(1).max(AI_VAULT_SEARCH_LIMIT_MAX).optional(),
agents: z.array(z.enum(AI_VAULT_AGENTS)).max(AI_VAULT_AGENTS.length).optional(),
scopePaths: z
.array(z.string().min(1).max(AI_VAULT_SCOPE_PATH_MAX_LENGTH))
.transform((paths) => paths.slice(0, AI_VAULT_SCOPE_PATHS_MAX_COUNT))
.optional(),
since: z.string().datetime({ offset: true }).optional(),
sort: z.enum(['relevance', 'newest']).optional(),
tier: z.enum(['full', 'conversation']).optional(),
refresh: OptionalBoolean,
executionHostId: executionHostIdSchema.optional()
})
export const AI_VAULT_METHODS: RpcMethod[] = [
defineMethod({
name: 'aiVault.searchSessions',
params: AiVaultSearchSessionsParams,
// Why: the index lives with the transcripts, so this runs on the host the
// client addressed; the id only names that host, it never redirects the search.
handler: ({ executionHostId: _host, ...params }, { runtime, signal }) =>
runtime.searchAiVaultSessions(params, signal)
}),
defineMethod({
name: 'aiVault.searchCoverage',
params: z.object({ executionHostId: executionHostIdSchema.optional() }),
handler: (_params, { runtime, signal }) => runtime.readAiVaultSearchCoverage(signal)
}),
defineMethod({
name: 'aiVault.resolveSessionTitles',
params: AiVaultSessionTitlesParams,
+18 -1
View File
@@ -7,7 +7,16 @@ import type {
AiVaultSessionTitlesResult
} from '../../shared/ai-vault-session-title'
import type { AiVaultListArgs, AiVaultListResult } from '../../shared/ai-vault-types'
import { listAiVaultSessions } from '../ai-vault/cached-session-list'
import {
listAiVaultSessions,
readAiVaultSearchCoverage,
searchAiVaultSessions
} from '../ai-vault/cached-session-list'
import type {
AiVaultSearchArgs,
AiVaultSearchCoverage,
AiVaultSearchResult
} from '../../shared/ai-vault-search-types'
import { resolveLocalAiVaultSessionTitles } from '../ai-vault/session-title-resolver'
export class RuntimeAiVaultCommands {
@@ -21,6 +30,14 @@ export class RuntimeAiVaultCommands {
return listAiVaultSessions(args)
}
search(args: AiVaultSearchArgs, signal?: AbortSignal): Promise<AiVaultSearchResult> {
return searchAiVaultSessions(args, { signal })
}
searchCoverage(signal?: AbortSignal): Promise<AiVaultSearchCoverage> {
return readAiVaultSearchCoverage({ signal })
}
resolveTitles(
requests: AiVaultSessionTitleRequest[],
signal?: AbortSignal
@@ -11,6 +11,8 @@ import type { RuntimeSubscriptionRegistry } from './runtime-subscription-registr
export type RuntimeServiceCommandSurface = {
listAiVaultSessions: RuntimeAiVaultCommands['list']
searchAiVaultSessions: RuntimeAiVaultCommands['search']
readAiVaultSearchCoverage: RuntimeAiVaultCommands['searchCoverage']
resolveAiVaultSessionTitles: RuntimeAiVaultCommands['resolveTitles']
prepareAiVaultSessionResume: RuntimeAiVaultCommands['prepare']
onClientEvent: RuntimeClientEventBus['on']
@@ -90,6 +92,8 @@ export function installRuntimeServiceCommandSurface(
const waiters = owners.messageWaiters
Object.assign(target, {
listAiVaultSessions: vault.list.bind(vault),
searchAiVaultSessions: vault.search.bind(vault),
readAiVaultSearchCoverage: vault.searchCoverage.bind(vault),
resolveAiVaultSessionTitles: vault.resolveTitles.bind(vault),
prepareAiVaultSessionResume: vault.prepare.bind(vault),
onClientEvent: events.on.bind(events),
@@ -66,6 +66,7 @@ import { setDefaultProxySessionResolver } from '../network/proxy-settings'
import { initDataPath, getCanonicalUserDataPath } from '../persistence'
import { applyMacPressAndHoldDefaultAtStartup } from '../macos-press-and-hold-default'
import { initSessionParseCachePersistence } from '../ai-vault/session-parse-cache-persistence'
import { initSessionSearchPaths } from '../ai-vault-search/session-search-paths'
import { initOrcaProfilePaths } from '../orca-profiles/profile-index-store'
import { initStatsPath } from '../stats/collector'
import { initClaudeUsagePath } from '../claude-usage/store'
@@ -269,6 +270,7 @@ export function runMainProcessPreflight(options: MainProcessPreflightOptions): b
filePath: join(getCanonicalUserDataPath(), 'ai-vault', 'session-parse-cache.json'),
appVersion: app.getVersion()
})
initSessionSearchPaths(getCanonicalUserDataPath())
initOrcaProfilePaths()
// Why: same timing as initDataPath — capture userData before app.setName changes it. See persistence.ts:20-28.
initStatsPath()
+73
View File
@@ -0,0 +1,73 @@
import type { AiVaultAgent, AiVaultSessionPreviewMessage } from './ai-vault-types'
export const AI_VAULT_SEARCH_QUERY_MAX_LENGTH = 512
export const AI_VAULT_SEARCH_LIMIT_MAX = 100
export const AI_VAULT_SEARCH_LIMIT_DEFAULT = 20
export type AiVaultSearchSort = 'relevance' | 'newest'
export type AiVaultSearchArgs = {
query: string
limit?: number
agents?: readonly AiVaultAgent[]
/** Restrict to sessions whose cwd is inside one of these paths. */
scopePaths?: readonly string[]
/** ISO timestamp; only sessions updated at or after it. */
since?: string
sort?: AiVaultSearchSort
/** As-you-type tier: conversation-only index (no tool output), ~10x faster. */
tier?: 'full' | 'conversation'
/** Fold in transcript appends before searching (default true; the as-you-type tier passes false). */
refresh?: boolean
}
export type AiVaultSearchEvidence = {
role: AiVaultSessionPreviewMessage['role']
timestamp: string | null
/** FTS5 snippet with the matched terms wrapped in `[` `]`. */
snippet: string
}
export type AiVaultSearchHit = {
agent: AiVaultAgent
sessionId: string
filePath: string
codexHome: string | null
title: string
cwd: string | null
branch: string | null
updatedAt: string | null
messageCount: number
resumeCommand: string
score: number
evidence: AiVaultSearchEvidence
}
/** How the query was executed; logged locally so the eval set can be rebuilt from real usage. */
export type AiVaultSearchRoute = 'phrase' | 'and' | 'or' | 'typo+phrase' | 'typo+and' | 'typo+or'
export type AiVaultSearchResult = {
hits: AiVaultSearchHit[]
route: AiVaultSearchRoute
/** Query terms after typo repair, when any were changed. */
repairedTerms?: string[]
durationMs: number
coverage: AiVaultSearchCoverage
}
export type AiVaultSearchProviderCoverage = {
agent: AiVaultAgent
sessionsIndexed: number
messagesIndexed: number
}
export type AiVaultSearchCoverage = {
sessionsIndexed: number
messagesIndexed: number
providers: AiVaultSearchProviderCoverage[]
/** `running` means older sessions are still being added; results are partial until `complete`. */
backfill: 'idle' | 'running' | 'complete'
/** Files a list scan saw change that the index has not re-read yet. */
filesPending: number
lastIndexedAt: string | null
}
+1
View File
@@ -24,6 +24,7 @@ export const CLI_BOOLEAN_FLAGS = new Set([
'me',
'mobile',
'mobile-pairing',
'newest',
'no-pairing',
'screen',
'parent-current',
+2
View File
@@ -69,6 +69,7 @@ export const JIRA_USER_FIELDS_UPDATE_REQUIRED_MESSAGE =
// conditional like browser.headless.v1.
export const AI_VAULT_RUNTIME_CAPABILITY = 'aiVault.v1' as const
export const AI_VAULT_SESSION_TITLES_RUNTIME_CAPABILITY = 'aiVault.session-titles.v1' as const
export const AI_VAULT_SESSION_SEARCH_RUNTIME_CAPABILITY = 'aiVault.session-search.v1' as const
// Why: signals a host owns browser pages with no renderer (headless serve via the
// offscreen backend). Advertised only when that backend is actually available, so
// clients never fall back to a local desktop browser tab for a remote-owned page.
@@ -225,6 +226,7 @@ export const RUNTIME_CAPABILITIES = [
JIRA_USER_FIELDS_RUNTIME_CAPABILITY,
AI_VAULT_RUNTIME_CAPABILITY,
AI_VAULT_SESSION_TITLES_RUNTIME_CAPABILITY,
AI_VAULT_SESSION_SEARCH_RUNTIME_CAPABILITY,
TERMINAL_QUERY_REPLY_INPUT_RUNTIME_CAPABILITY,
TERMINAL_PAIRED_PARKING_RUNTIME_CAPABILITY,
TERMINAL_QUICK_COMMANDS_RUNTIME_CAPABILITY,