diff --git a/src/cli/agent-session-search-format.ts b/src/cli/agent-session-search-format.ts new file mode 100644 index 00000000000..d7b1535d862 --- /dev/null +++ b/src/cli/agent-session-search-format.ts @@ -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 = { + 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') +} diff --git a/src/cli/agent-session-search-help.ts b/src/cli/agent-session-search-help.ts new file mode 100644 index 00000000000..edf851bf218 --- /dev/null +++ b/src/cli/agent-session-search-help.ts @@ -0,0 +1,12 @@ +const SEARCH_FLAG_HELP: Record = { + 'agent-session': '--agent-session Text to find in any coding-agent session (required)', + agent: '--agent Only this agent (claude, codex, cursor, …); repeatable', + path: '--path Only sessions whose working directory is under this path; repeatable', + since: '--since Only sessions updated at or after this ISO 8601 timestamp', + newest: '--newest Sort by session time instead of relevance', + host: '--host Search a paired runtime host (runtime:)' +} + +export function formatSearchFlagHelp(command: string, flag: string): string | null { + return command === 'search' ? (SEARCH_FLAG_HELP[flag] ?? null) : null +} diff --git a/src/cli/handler-group-manifest.ts b/src/cli/handler-group-manifest.ts index d995b4d8687..8e5bf409048 100644 --- a/src/cli/handler-group-manifest.ts +++ b/src/cli/handler-group-manifest.ts @@ -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 } ] diff --git a/src/cli/handlers/search.ts b/src/cli/handlers/search.ts new file mode 100644 index 00000000000..aedfa93ada2 --- /dev/null +++ b/src/cli/handlers/search.ts @@ -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): 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 | 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 = { + search: async ({ client, flags, json, cwd }) => { + const query = getOptionalStringFlag(flags, 'agent-session') + if (!query) { + throw new RuntimeClientError( + 'invalid_argument', + 'Missing --agent-session . 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: or omit --host.' + ) + } + const scopePaths = getRepeatedStringFlag(flags, 'path').map((value) => + value.startsWith('~') ? value.replace(/^~/, process.env.HOME ?? '~') : value + ) + const result = await client.call('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 } diff --git a/src/cli/help.ts b/src/cli/help.ts index c7beec4018a..0613610c67b 100644 --- a/src/cli/help.ts +++ b/src/cli/help.ts @@ -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 Comma-separated install targets; default is detected agents' } diff --git a/src/cli/root-help-text-secondary.ts b/src/cli/root-help-text-secondary.ts index 870c4f50836..d9364c38d61 100644 --- a/src/cli/root-help-text-secondary.ts +++ b/src/cli/root-help-text-secondary.ts @@ -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 "" [--limit ] [--agent ] [--newest] [--json]', ' orca host list [--json]', ' orca environment add --name --pairing-code [--json]', ' orca environment list [--json]', diff --git a/src/cli/specs/index.ts b/src/cli/specs/index.ts index 92829793119..ee9db22dc1a 100644 --- a/src/cli/specs/index.ts +++ b/src/cli/specs/index.ts @@ -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 ] diff --git a/src/cli/specs/search.ts b/src/cli/specs/search.ts new file mode 100644 index 00000000000..641bcb2b4e1 --- /dev/null +++ b/src/cli/specs/search.ts @@ -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 "" [--limit ] [--agent ] [--path ] [--since ] [--newest] [--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: 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' + ] + } +] diff --git a/src/main/ai-vault-search/session-search-identifier-split.ts b/src/main/ai-vault-search/session-search-identifier-split.ts new file mode 100644 index 00000000000..e2df1822cfa --- /dev/null +++ b/src/main/ai-vault-search/session-search-identifier-split.ts @@ -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() + 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(' ') +} diff --git a/src/main/ai-vault-search/session-search-index-writer.ts b/src/main/ai-vault-search/session-search-index-writer.ts new file mode 100644 index 00000000000..359c30ce764 --- /dev/null +++ b/src/main/ai-vault-search/session-search-index-writer.ts @@ -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 | 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 | 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 + ) + } +} diff --git a/src/main/ai-vault-search/session-search-paths.ts b/src/main/ai-vault-search/session-search-paths.ts new file mode 100644 index 00000000000..d2b896e7ba5 --- /dev/null +++ b/src/main/ai-vault-search/session-search-paths.ts @@ -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 +} diff --git a/src/main/ai-vault-search/session-search-query-planner.ts b/src/main/ai-vault-search/session-search-query-planner.ts new file mode 100644 index 00000000000..163b6395320 --- /dev/null +++ b/src/main/ai-vault-search/session-search-query-planner.ts @@ -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 ') +} diff --git a/src/main/ai-vault-search/session-search-query.ts b/src/main/ai-vault-search/session-search-query.ts new file mode 100644 index 00000000000..558c213d565 --- /dev/null +++ b/src/main/ai-vault-search/session-search-query.ts @@ -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() + 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[] + } +} diff --git a/src/main/ai-vault-search/session-search-schema.ts b/src/main/ai-vault-search/session-search-schema.ts new file mode 100644 index 00000000000..631f60e1d4f --- /dev/null +++ b/src/main/ai-vault-search/session-search-schema.ts @@ -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 +} diff --git a/src/main/ai-vault-search/session-search-service.ts b/src/main/ai-vault-search/session-search-service.ts new file mode 100644 index 00000000000..696e71f3b93 --- /dev/null +++ b/src/main/ai-vault-search/session-search-service.ts @@ -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 + +/** + * 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 | 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 { + 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 { + 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 { + 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 { + const stale = this.store.takeStale() + if (stale.length === 0) { + return + } + await this.parseAll(stale, signal) + } + + private async runBackfill(roots: SessionSearchScanRoots, signal: AbortSignal): Promise { + 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 { + 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)) + } + } + }) + } +} diff --git a/src/main/ai-vault-search/session-search-store.test.ts b/src/main/ai-vault-search/session-search-store.test.ts new file mode 100644 index 00000000000..ddfc4025070 --- /dev/null +++ b/src/main/ai-vault-search/session-search-store.test.ts @@ -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 { + const root = await mkdtemp(join(tmpdir(), 'orca-session-search-')) + tempRoots.push(root) + return root +} + +async function claudeCandidate(path: string): Promise { + 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) + }) +}) diff --git a/src/main/ai-vault-search/session-search-store.ts b/src/main/ai-vault-search/session-search-store.ts new file mode 100644 index 00000000000..c04ca049979 --- /dev/null +++ b/src/main/ai-vault-search/session-search-store.ts @@ -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() + + 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) + } + } +} diff --git a/src/main/ai-vault-search/session-search-typo-repair.ts b/src/main/ai-vault-search/session-search-typo-repair.ts new file mode 100644 index 00000000000..f768717229b --- /dev/null +++ b/src/main/ai-vault-search/session-search-typo-repair.ts @@ -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({ length: b.length + 1 }).fill(0) + let current = Array.from({ 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 + private readonly candidatesByPrefix: ReturnType + + 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[] + } +} diff --git a/src/main/ai-vault/cached-session-list.ts b/src/main/ai-vault/cached-session-list.ts index c8d04205091..5a0d187f7ca 100644 --- a/src/main/ai-vault/cached-session-list.ts +++ b/src/main/ai-vault/cached-session-list.ts @@ -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 +> { + 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 { + const roots = await resolveAiVaultHostScanSources() + return searchAiVaultSessionsInBackground({ args, roots }, options.signal) +} + +export async function readAiVaultSearchCoverage( + options: { signal?: AbortSignal } = {} +): Promise { + 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 ) diff --git a/src/main/ai-vault/session-parse-cache-store.ts b/src/main/ai-vault/session-parse-cache-store.ts new file mode 100644 index 00000000000..31beed5d520 --- /dev/null +++ b/src/main/ai-vault/session-parse-cache-store.ts @@ -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() + +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 + +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) +} diff --git a/src/main/ai-vault/session-scanner-accumulator.ts b/src/main/ai-vault/session-scanner-accumulator.ts index caf08f4259d..bb1aef36f34 100644 --- a/src/main/ai-vault/session-scanner-accumulator.ts +++ b/src/main/ai-vault/session-scanner-accumulator.ts @@ -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 }) } diff --git a/src/main/ai-vault/session-scanner-background.ts b/src/main/ai-vault/session-scanner-background.ts index 6903e689a64..eff94842379 100644 --- a/src/main/ai-vault/session-scanner-background.ts +++ b/src/main/ai-vault/session-scanner-background.ts @@ -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 { + return shouldUseAiVaultServiceProcess() + ? searchAiVaultSessionsInService(request, signal) + : searchAiVaultSessionsInWorker(request, signal) +} + +export function readAiVaultSearchCoverageInBackground( + request: Pick, + signal?: AbortSignal +): Promise { + return shouldUseAiVaultServiceProcess() + ? readAiVaultSearchCoverageInService(request, signal) + : readAiVaultSearchCoverageInWorker(request) +} + export function invalidateAiVaultBackgroundCache(paths: string[]): Promise { return shouldUseAiVaultServiceProcess() ? invalidateAiVaultServiceCache(paths) : Promise.resolve() } diff --git a/src/main/ai-vault/session-scanner-candidates.ts b/src/main/ai-vault/session-scanner-candidates.ts new file mode 100644 index 00000000000..80f146367e5 --- /dev/null +++ b/src/main/ai-vault/session-scanner-candidates.ts @@ -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 { + 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 + ) +} diff --git a/src/main/ai-vault/session-scanner-codex-parser.ts b/src/main/ai-vault/session-scanner-codex-parser.ts index 80a76ce0a92..ab2992527b0 100644 --- a/src/main/ai-vault/session-scanner-codex-parser.ts +++ b/src/main/ai-vault/session-scanner-codex-parser.ts @@ -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): 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 | null { - return ( - normalizeTitleText(extractString(payload.title) ?? '') ?? - normalizeTitleText(extractString(payload.thread_name) ?? '') ?? - normalizeTitleText(extractString(payload.threadName) ?? '') - ) -} diff --git a/src/main/ai-vault/session-scanner-codex-session-meta.ts b/src/main/ai-vault/session-scanner-codex-session-meta.ts new file mode 100644 index 00000000000..9b939b6b514 --- /dev/null +++ b/src/main/ai-vault/session-scanner-codex-session-meta.ts @@ -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): 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 | null { + return ( + normalizeTitleText(extractString(payload.title) ?? '') ?? + normalizeTitleText(extractString(payload.thread_name) ?? '') ?? + normalizeTitleText(extractString(payload.threadName) ?? '') + ) +} diff --git a/src/main/ai-vault/session-scanner-jsonl-reader.ts b/src/main/ai-vault/session-scanner-jsonl-reader.ts index 613c50ef0ba..1c80598fedf 100644 --- a/src/main/ai-vault/session-scanner-jsonl-reader.ts +++ b/src/main/ai-vault/session-scanner-jsonl-reader.ts @@ -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 diff --git a/src/main/ai-vault/session-scanner-parse-cache.test.ts b/src/main/ai-vault/session-scanner-parse-cache.test.ts index 8434f26bac1..e5a470d9a69 100644 --- a/src/main/ai-vault/session-scanner-parse-cache.test.ts +++ b/src/main/ai-vault/session-scanner-parse-cache.test.ts @@ -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 { 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 { + 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') diff --git a/src/main/ai-vault/session-scanner-parse-cache.ts b/src/main/ai-vault/session-scanner-parse-cache.ts index 139940ba005..8b20458c2ab 100644 --- a/src/main/ai-vault/session-scanner-parse-cache.ts +++ b/src/main/ai-vault/session-scanner-parse-cache.ts @@ -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() - -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 - -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 { 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 +): Promise { + 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 { 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 => + 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 { - const slice = await readTranscriptSlice(path, offset - 1, 1, 'scan') - return slice.length === 1 && slice[0] === NEWLINE_BYTE -} diff --git a/src/main/ai-vault/session-scanner-resume-point.ts b/src/main/ai-vault/session-scanner-resume-point.ts new file mode 100644 index 00000000000..8760b08569b --- /dev/null +++ b/src/main/ai-vault/session-scanner-resume-point.ts @@ -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 { + const slice = await readTranscriptSlice(path, offset - 1, 1, 'scan') + return slice.length === 1 && slice[0] === NEWLINE_BYTE +} diff --git a/src/main/ai-vault/session-scanner-service-entry.ts b/src/main/ai-vault/session-scanner-service-entry.ts index 74a5ea9c355..e75fca4ebc3 100644 --- a/src/main/ai-vault/session-scanner-service-entry.ts +++ b/src/main/ai-vault/session-scanner-service-entry.ts @@ -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() const pending = new Set() const titleIndex = new Map() const invalidatedPaths = new Set() +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 { const controller = new AbortController() controllers.set(request.id, controller) @@ -77,6 +86,22 @@ async function executeRequest(request: AiVaultServiceRequest): Promise { } 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) diff --git a/src/main/ai-vault/session-scanner-service-protocol.ts b/src/main/ai-vault/session-scanner-service-protocol.ts index f2842751eb1..55eae963f85 100644 --- a/src/main/ai-vault/session-scanner-service-protocol.ts +++ b/src/main/ai-vault/session-scanner-service-protocol.ts @@ -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 } 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 + } 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') ) } diff --git a/src/main/ai-vault/session-scanner-service-spawn.ts b/src/main/ai-vault/session-scanner-service-spawn.ts index 3ba12322734..5b0df720b98 100644 --- a/src/main/ai-vault/session-scanner-service-spawn.ts +++ b/src/main/ai-vault/session-scanner-service-spawn.ts @@ -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 { + return getSharedClient().request({ type: 'request', operation: 'search', request }, signal) +} + +export function readAiVaultSearchCoverageInService( + request: Pick, + signal?: AbortSignal +): Promise { + return getSharedClient().request( + { type: 'request', operation: 'searchCoverage', request }, + signal + ) +} + export function invalidateAiVaultServiceCache(paths: string[]): Promise { return sharedClient?.invalidate(paths) ?? Promise.resolve() } diff --git a/src/main/ai-vault/session-scanner-text-normalization.test.ts b/src/main/ai-vault/session-scanner-text-normalization.test.ts new file mode 100644 index 00000000000..41bcc6c6d12 --- /dev/null +++ b/src/main/ai-vault/session-scanner-text-normalization.test.ts @@ -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() + }) +}) diff --git a/src/main/ai-vault/session-scanner-text-normalization.ts b/src/main/ai-vault/session-scanner-text-normalization.ts index 99f73a93d02..ad77ed8fa6b 100644 --- a/src/main/ai-vault/session-scanner-text-normalization.ts +++ b/src/main/ai-vault/session-scanner-text-normalization.ts @@ -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 | 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 diff --git a/src/main/ai-vault/session-scanner-worker-client.ts b/src/main/ai-vault/session-scanner-worker-client.ts index 70e75c4f982..c98cfd0e71c 100644 --- a/src/main/ai-vault/session-scanner-worker-client.ts +++ b/src/main/ai-vault/session-scanner-worker-client.ts @@ -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, 'id'> - | Omit, 'id'> +type RequestBody = AiVaultWorkerRequest extends infer R + ? R extends { id: number } + ? Omit + : never + : never type PendingCall = { request: AiVaultWorkerRequest @@ -64,6 +70,23 @@ export class AiVaultScannerWorkerClient { ) as Promise } + search(request: AiVaultServiceSearchRequest, signal?: AbortSignal): Promise { + return this.dispatch( + { kind: 'search', request }, + SEARCH_TIMEOUT_MS, + signal + ) as Promise + } + + searchCoverage( + request: Pick + ): Promise { + return this.dispatch( + { kind: 'searchCoverage', request }, + TITLE_TIMEOUT_MS + ) as Promise + } + dispose(): void { this.destroyWorker() const pending = this.queue diff --git a/src/main/ai-vault/session-scanner-worker-entry.ts b/src/main/ai-vault/session-scanner-worker-entry.ts index 33c06d30401..d40042ed7ae 100644 --- a/src/main/ai-vault/session-scanner-worker-entry.ts +++ b/src/main/ai-vault/session-scanner-worker-entry.ts @@ -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() const titleIndex = new Map() @@ -66,6 +68,28 @@ async function handleRequest(request: AiVaultWorkerRequest): Promise 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 } 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 } diff --git a/src/main/ai-vault/session-scanner-worker-spawn.ts b/src/main/ai-vault/session-scanner-worker-spawn.ts index 3c5abbd3551..2a9e9499866 100644 --- a/src/main/ai-vault/session-scanner-worker-spawn.ts +++ b/src/main/ai-vault/session-scanner-worker-spawn.ts @@ -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 { + return getSharedClient().search(request, signal) +} + +export function readAiVaultSearchCoverageInWorker( + request: Pick +): Promise { + return getSharedClient().searchCoverage(request) +} + export function resetAiVaultScannerWorkerForTests(): void { sharedClient?.dispose() sharedClient = null diff --git a/src/main/ai-vault/session-scanner.ts b/src/main/ai-vault/session-scanner.ts index 3d27dccd1ef..910ec42f2c9 100644 --- a/src/main/ai-vault/session-scanner.ts +++ b/src/main/ai-vault/session-scanner.ts @@ -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), diff --git a/src/main/ai-vault/session-search-capture.ts b/src/main/ai-vault/session-search-capture.ts new file mode 100644 index 00000000000..0eaec0f5b49 --- /dev/null +++ b/src/main/ai-vault/session-search-capture.ts @@ -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() +const indexModeStorage = new AsyncLocalStorage() +let sink: SessionSearchIndexSink | null = null + +export function getSessionSearchIndexMode(): SessionSearchIndexMode { + return indexModeStorage.getStore() ?? 'opportunistic' +} + +export function withSessionSearchIndexRequired(fn: () => Promise): Promise { + 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(fn: () => T): T { + return captureStorage.run(null, fn) +} + +export async function withSessionSearchCapture( + fn: () => Promise +): Promise<{ value: T; messages: SessionSearchCapturedMessage[] }> { + const scope = { messages: [] as SessionSearchCapturedMessage[] } + const value = await captureStorage.run(scope, fn) + return { value, messages: scope.messages } +} diff --git a/src/main/ai-vault/session-search-codex-tool-records.ts b/src/main/ai-vault/session-search-codex-tool-records.ts new file mode 100644 index 00000000000..7ecf0ba62bc --- /dev/null +++ b/src/main/ai-vault/session-search-codex-tool-records.ts @@ -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, + 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) + } +} diff --git a/src/main/ai-vault/session-search-content.ts b/src/main/ai-vault/session-search-content.ts new file mode 100644 index 00000000000..50cad8da623 --- /dev/null +++ b/src/main/ai-vault/session-search-content.ts @@ -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 | null { + return value && typeof value === 'object' && !Array.isArray(value) + ? (value as Record) + : 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 + }) +} diff --git a/src/main/runtime/rpc/methods/ai-vault.ts b/src/main/runtime/rpc/methods/ai-vault.ts index c689165924d..c21bb393a0c 100644 --- a/src/main/runtime/rpc/methods/ai-vault.ts +++ b/src/main/runtime/rpc/methods/ai-vault.ts @@ -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, diff --git a/src/main/runtime/runtime-ai-vault-commands.ts b/src/main/runtime/runtime-ai-vault-commands.ts index 2701d782b61..ce31e677773 100644 --- a/src/main/runtime/runtime-ai-vault-commands.ts +++ b/src/main/runtime/runtime-ai-vault-commands.ts @@ -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 { + return searchAiVaultSessions(args, { signal }) + } + + searchCoverage(signal?: AbortSignal): Promise { + return readAiVaultSearchCoverage({ signal }) + } + resolveTitles( requests: AiVaultSessionTitleRequest[], signal?: AbortSignal diff --git a/src/main/runtime/runtime-service-command-surface.ts b/src/main/runtime/runtime-service-command-surface.ts index 19545cc76e6..ddb9d00c543 100644 --- a/src/main/runtime/runtime-service-command-surface.ts +++ b/src/main/runtime/runtime-service-command-surface.ts @@ -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), diff --git a/src/main/startup/main-process-preflight.ts b/src/main/startup/main-process-preflight.ts index 177a2357441..79266d3661c 100644 --- a/src/main/startup/main-process-preflight.ts +++ b/src/main/startup/main-process-preflight.ts @@ -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() diff --git a/src/shared/ai-vault-search-types.ts b/src/shared/ai-vault-search-types.ts new file mode 100644 index 00000000000..21df08eab4e --- /dev/null +++ b/src/shared/ai-vault-search-types.ts @@ -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 +} diff --git a/src/shared/cli-argument-boundary.ts b/src/shared/cli-argument-boundary.ts index 7b29db49c45..7503ed3f092 100644 --- a/src/shared/cli-argument-boundary.ts +++ b/src/shared/cli-argument-boundary.ts @@ -24,6 +24,7 @@ export const CLI_BOOLEAN_FLAGS = new Set([ 'me', 'mobile', 'mobile-pairing', + 'newest', 'no-pairing', 'screen', 'parent-current', diff --git a/src/shared/protocol-version.ts b/src/shared/protocol-version.ts index e4121c95a67..72b4b826388 100644 --- a/src/shared/protocol-version.ts +++ b/src/shared/protocol-version.ts @@ -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,