From 4dfbc66e9769c60d3a6afda78e7958b4e1663921 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Wed, 9 Sep 2026 01:44:11 -0400 Subject: [PATCH] feat(ai-vault-search): add the session search index schema and row modules The FTS5 index that PR 2 folds the transcript reader's message stream into: two FTS tables behind visibility views, the identifier shadow column, the resumable fork digest, redaction at the row-insert choke point, and the bounded compaction, warm-up and retention lanes. Schema version starts at 1: this branch drops the query log and the fts5vocab table, which the query engine reintroduces with its own bump. --- .../session-search-content-hash.test.ts | 49 +++++ .../session-search-content-hash.ts | 45 +++++ .../session-search-file-cursor.ts | 36 ++++ .../session-search-identifier-split.test.ts | 29 +++ .../session-search-identifier-split.ts | 54 +++++ .../session-search-index-compaction.ts | 22 ++ .../session-search-message-rows.test.ts | 52 +++++ .../session-search-message-rows.ts | 71 +++++++ .../session-search-page-warmup.ts | 27 +++ .../session-search-path-key.ts | 12 ++ .../ai-vault-search/session-search-paths.ts | 17 ++ .../session-search-pending-deletes.ts | 35 ++++ .../session-search-redaction.test.ts | 27 +++ .../session-search-redaction.ts | 26 +++ .../session-search-retention-delete.ts | 97 +++++++++ .../session-search-schema.test.ts | 191 ++++++++++++++++++ .../ai-vault-search/session-search-schema.ts | 144 +++++++++++++ ...ession-search-staged-write-test-fixture.ts | 124 ++++++++++++ .../session-search-synthetic-corpus.ts | 132 ++++++++++++ .../session-search-transcript-fixtures.ts | 111 ++++++++++ .../session-search-wal-budget.ts | 24 +++ src/main/observability/redactor.ts | 2 +- 22 files changed, 1326 insertions(+), 1 deletion(-) create mode 100644 src/main/ai-vault-search/session-search-content-hash.test.ts create mode 100644 src/main/ai-vault-search/session-search-content-hash.ts create mode 100644 src/main/ai-vault-search/session-search-file-cursor.ts create mode 100644 src/main/ai-vault-search/session-search-identifier-split.test.ts create mode 100644 src/main/ai-vault-search/session-search-identifier-split.ts create mode 100644 src/main/ai-vault-search/session-search-index-compaction.ts create mode 100644 src/main/ai-vault-search/session-search-message-rows.test.ts create mode 100644 src/main/ai-vault-search/session-search-message-rows.ts create mode 100644 src/main/ai-vault-search/session-search-page-warmup.ts create mode 100644 src/main/ai-vault-search/session-search-path-key.ts create mode 100644 src/main/ai-vault-search/session-search-paths.ts create mode 100644 src/main/ai-vault-search/session-search-pending-deletes.ts create mode 100644 src/main/ai-vault-search/session-search-redaction.test.ts create mode 100644 src/main/ai-vault-search/session-search-redaction.ts create mode 100644 src/main/ai-vault-search/session-search-retention-delete.ts create mode 100644 src/main/ai-vault-search/session-search-schema.test.ts create mode 100644 src/main/ai-vault-search/session-search-schema.ts create mode 100644 src/main/ai-vault-search/session-search-staged-write-test-fixture.ts create mode 100644 src/main/ai-vault-search/session-search-synthetic-corpus.ts create mode 100644 src/main/ai-vault-search/session-search-transcript-fixtures.ts create mode 100644 src/main/ai-vault-search/session-search-wal-budget.ts diff --git a/src/main/ai-vault-search/session-search-content-hash.test.ts b/src/main/ai-vault-search/session-search-content-hash.test.ts new file mode 100644 index 00000000000..bff194e77d4 --- /dev/null +++ b/src/main/ai-vault-search/session-search-content-hash.test.ts @@ -0,0 +1,49 @@ +import { expect, it } from 'vitest' +import { + CONTENT_HASH_MESSAGE_LIMIT, + EMPTY_CONTENT_HASH, + foldContentHash, + isCollapsibleContentHash +} from './session-search-content-hash' +import { userMessages } from './session-search-staged-write-test-fixture' + +it('reaches the same digest whether the prefix arrives whole or in two appends', () => { + const messages = userMessages('turn', 5) + const whole = foldContentHash(EMPTY_CONTENT_HASH, messages) + const resumed = foldContentHash( + foldContentHash(EMPTY_CONTENT_HASH, messages.slice(0, 2)), + messages.slice(2) + ) + + expect(resumed).toEqual(whole) + expect(whole.count).toBe(5) +}) + +it('freezes once the prefix limit is reached so later appends cannot move it', () => { + const capped = foldContentHash( + EMPTY_CONTENT_HASH, + userMessages('turn', CONTENT_HASH_MESSAGE_LIMIT) + ) + expect(foldContentHash(capped, userMessages('later', 20))).toEqual(capped) +}) + +it('separates two conversations that share an opening prompt', () => { + const shared = userMessages('same opening', 1) + const first = foldContentHash(EMPTY_CONTENT_HASH, [ + ...shared, + { role: 'user', text: 'left', timestamp: null } + ]) + const second = foldContentHash(EMPTY_CONTENT_HASH, [ + ...shared, + { role: 'user', text: 'right', timestamp: null } + ]) + expect(first.hash).not.toBe(second.hash) +}) + +it('refuses to collapse on a prefix too short to mean anything', () => { + const one = foldContentHash(EMPTY_CONTENT_HASH, userMessages('only turn', 1)) + expect(isCollapsibleContentHash(one.hash, one.count)).toBe(false) + const two = foldContentHash(EMPTY_CONTENT_HASH, userMessages('two turns', 2)) + expect(isCollapsibleContentHash(two.hash, two.count)).toBe(true) + expect(isCollapsibleContentHash(null, 9)).toBe(false) +}) diff --git a/src/main/ai-vault-search/session-search-content-hash.ts b/src/main/ai-vault-search/session-search-content-hash.ts new file mode 100644 index 00000000000..acb5c54f512 --- /dev/null +++ b/src/main/ai-vault-search/session-search-content-hash.ts @@ -0,0 +1,45 @@ +import { createHash } from 'node:crypto' +import type { TranscriptMessage } from '../ai-vault/session-transcript-consumers' + +// Why: Claude `--resume` and Codex fork copy the parent transcript into a new +// file under a new session id, so one conversation lands N times in results. +// The shared opening prefix is what identifies the copy; the tail diverges. +export const CONTENT_HASH_MESSAGE_LIMIT = 8 +// One shared opening prompt is not evidence of a fork; two turns is. +export const CONTENT_HASH_MIN_MESSAGES = 2 + +export type SessionContentHash = { hash: string | null; count: number } + +export const EMPTY_CONTENT_HASH: SessionContentHash = { hash: null, count: 0 } + +/** + * Chained digest over the first `CONTENT_HASH_MESSAGE_LIMIT` messages. Chaining + * (rather than hashing one joined string) makes it resumable, so an `append` + * can finish a prefix a short `replace` started; once the limit is reached the + * value is frozen and later appends leave it untouched. + */ +export function foldContentHash( + previous: SessionContentHash, + messages: readonly TranscriptMessage[] +): SessionContentHash { + let { hash, count } = previous + for (const message of messages) { + if (count >= CONTENT_HASH_MESSAGE_LIMIT) { + break + } + hash = createHash('sha256') + .update(hash ?? '') + .update('\0') + .update(message.role) + .update('\0') + .update(message.text) + .digest('hex') + count += 1 + } + return { hash, count } +} + +/** Sessions collapse only on a hash that covers enough turns to mean anything. */ +export function isCollapsibleContentHash(hash: string | null, count: number): hash is string { + return hash !== null && count >= CONTENT_HASH_MIN_MESSAGES +} diff --git a/src/main/ai-vault-search/session-search-file-cursor.ts b/src/main/ai-vault-search/session-search-file-cursor.ts new file mode 100644 index 00000000000..7075b4f9155 --- /dev/null +++ b/src/main/ai-vault-search/session-search-file-cursor.ts @@ -0,0 +1,36 @@ +import type { FileWithMtime } from '../ai-vault/session-scanner-types' + +// Why the index keeps its own cursor: the parse cache's cursor answers "what +// does the session list already show", which is a different question from "what +// bytes of this file are already rows". They diverge the moment either side +// declines a read, so neither may consult the other. + +/** Filesystem identity, when discovery could prove it. */ +export type SessionSearchFileIdentity = { dev: number; ino: number } | null + +/** What the index holds for one transcript. */ +export type SessionSearchIndexedFile = { + byteOffset: number + mtimeMs: number + sizeBytes: number | null +} + +export function fileIdentity(file: FileWithMtime): SessionSearchFileIdentity { + return typeof file.dev === 'number' && typeof file.ino === 'number' + ? { dev: file.dev, ino: file.ino } + : null +} + +/** True when the index already covers this file at its current stat. */ +export function isSessionSearchFileCurrent( + indexed: SessionSearchIndexedFile | null, + file: FileWithMtime +): boolean { + return ( + indexed !== null && + indexed.mtimeMs === file.mtimeMs && + (indexed.sizeBytes === null || + file.sizeBytes === undefined || + indexed.sizeBytes === file.sizeBytes) + ) +} diff --git a/src/main/ai-vault-search/session-search-identifier-split.test.ts b/src/main/ai-vault-search/session-search-identifier-split.test.ts new file mode 100644 index 00000000000..24f8a67d5c3 --- /dev/null +++ b/src/main/ai-vault-search/session-search-identifier-split.test.ts @@ -0,0 +1,29 @@ +import { expect, it } from 'vitest' +import { identifierShadowTerms, identifierShadowText } from './session-search-identifier-split' + +it('splits a camel-case symbol into its pieces and keeps the whole', () => { + expect(identifierShadowTerms('call resolveTerminalPath here')).toEqual([ + 'resolveterminalpath', + 'resolve', + 'terminal', + 'path' + ]) +}) + +it('splits a path into its segments and extension', () => { + // The whole path already tokenizes on its own; only the pieces need shadowing. + expect(identifierShadowText('src/main/foo-bar.ts')).toBe('src main foo bar ts') +}) + +it('leaves ordinary prose alone', () => { + expect(identifierShadowTerms('the quick brown fox')).toEqual([]) +}) + +it('shadows a screaming-case constant', () => { + expect(identifierShadowTerms('MAX_RETRIES')).toEqual(['max', 'retries']) +}) + +it('stops at the term limit rather than growing with the message', () => { + const text = Array.from({ length: 50 }, (_unused, index) => `alpha_beta${index}`).join(' ') + expect(identifierShadowTerms(text, 10)).toHaveLength(10) +}) 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-compaction.ts b/src/main/ai-vault-search/session-search-index-compaction.ts new file mode 100644 index 00000000000..93d2b2649d0 --- /dev/null +++ b/src/main/ai-vault-search/session-search-index-compaction.ts @@ -0,0 +1,22 @@ +import { setImmediate as yieldToEventLoop } from 'node:timers/promises' +import type SyncDatabase from '../sqlite/sync-database' + +const COMPACT_PAGES_PER_STEP = 2000 + +/** Hands freed pages back to the filesystem in bounded steps, never one long stall. */ +export async function compactSessionSearchIndex( + db: SyncDatabase, + stopped: () => boolean +): Promise { + let freed = Number(db.pragma('freelist_count', { simple: true })) + while (!stopped() && freed > 0) { + db.pragma(`incremental_vacuum(${COMPACT_PAGES_PER_STEP})`) + const remaining = Number(db.pragma('freelist_count', { simple: true })) + // Why: without auto_vacuum the step is a no-op; never spin on it. + if (remaining >= freed) { + return + } + freed = remaining + await yieldToEventLoop() + } +} diff --git a/src/main/ai-vault-search/session-search-message-rows.test.ts b/src/main/ai-vault-search/session-search-message-rows.test.ts new file mode 100644 index 00000000000..9ff9cc130ea --- /dev/null +++ b/src/main/ai-vault-search/session-search-message-rows.test.ts @@ -0,0 +1,52 @@ +import { expect, it } from 'vitest' +import { chunkMessageText, insertSearchMessage } from './session-search-message-rows' +import { openSessionSearchIndexFile } from './session-search-staged-write-test-fixture' + +it('splits an oversized message on a line boundary and keeps every character', () => { + const line = `${'padding '.repeat(11)}word\n` + const text = line.repeat(400) + const chunks = chunkMessageText(text) + + expect(chunks.length).toBeGreaterThan(1) + expect(chunks.join('')).toBe(text) + for (const chunk of chunks) { + expect(chunk.length).toBeLessThanOrEqual(8000) + expect(chunk.endsWith('\n')).toBe(true) + } +}) + +it('leaves a message that fits as a single row', () => { + expect(chunkMessageText('short enough')).toEqual(['short enough']) +}) + +it('redacts before either FTS table sees the text', async () => { + const index = await openSessionSearchIndexFile('ss-message-rows') + try { + index.db.exec('INSERT INTO search_write_batches(id,session_row_id) VALUES (1,1)') + insertSearchMessage(index.db, 1, 1, { + role: 'assistant', + text: 'use AKIAIOSFODNN7EXAMPLE for the upload', + timestamp: null + }) + for (const table of ['messages_fts', 'conversation_fts']) { + const row = index.db.prepare(`SELECT assistant_text AS text FROM ${table}`).get() as { + text: string + } + expect(row.text).toBe('use [redacted:aws-access-key-id] for the upload') + } + } finally { + await index.close() + } +}) + +it('keeps a tool row out of the conversation half', async () => { + const index = await openSessionSearchIndexFile('ss-message-rows-tool') + try { + index.db.exec('INSERT INTO search_write_batches(id,session_row_id) VALUES (1,1)') + insertSearchMessage(index.db, 1, 1, { role: 'tool', text: 'rg pericardium', timestamp: null }) + expect(index.db.prepare('SELECT count(*) AS n FROM messages_fts').get()).toEqual({ n: 1 }) + expect(index.db.prepare('SELECT count(*) AS n FROM conversation_fts').get()).toEqual({ n: 0 }) + } finally { + await index.close() + } +}) diff --git a/src/main/ai-vault-search/session-search-message-rows.ts b/src/main/ai-vault-search/session-search-message-rows.ts new file mode 100644 index 00000000000..5f1ce68f066 --- /dev/null +++ b/src/main/ai-vault-search/session-search-message-rows.ts @@ -0,0 +1,71 @@ +import type SyncDatabase from '../sqlite/sync-database' +import type { TranscriptMessage } from '../ai-vault/session-transcript-consumers' +import { identifierShadowText } from './session-search-identifier-split' +import { redactSessionSearchText } from './session-search-redaction' + +const CHUNK_TARGET_CHARS = 8000 + +function* textChunks(text: string): Generator { + if (text.length <= CHUNK_TARGET_CHARS) { + yield text + return + } + 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 + } + } + yield text.slice(start, end) + start = end + } +} + +/** One message becomes N rows: FTS5 ranks a short row far better than a huge one. */ +export function* searchMessageRows( + messages: Iterable +): Generator { + for (const message of messages) { + for (const text of textChunks(message.text)) { + yield { ...message, text } + } + } +} + +/** + * Writes one row into `messages` and both FTS tables in the caller's + * transaction, so a published message is never present in one table and absent + * from the other. `tool` rows stay out of `conversation_fts`: that table is the + * conversation-only half of the split. + */ +export function insertSearchMessage( + db: SyncDatabase, + sessionId: number, + batchId: number, + message: TranscriptMessage +): void { + const text = redactSessionSearchText(message.text) + const id = db + .prepare('INSERT INTO messages(session_row_id, batch_id, role, ts) VALUES (?, ?, ?, ?)') + .run(sessionId, batchId, message.role, message.timestamp).lastInsertRowid + const user = message.role === 'user' ? text : '' + const assistant = message.role === 'assistant' ? text : '' + const tool = message.role === 'tool' ? text : '' + db.prepare( + 'INSERT INTO messages_fts(rowid,user_text,assistant_text,tool_text,identifiers) VALUES (?,?,?,?,?)' + ).run(id, user, assistant, tool, identifierShadowText(text)) + if (message.role !== 'tool') { + db.prepare('INSERT INTO conversation_fts(rowid,user_text,assistant_text) VALUES (?,?,?)').run( + id, + user, + assistant + ) + } +} + +export function chunkMessageText(text: string): string[] { + return [...textChunks(text)] +} diff --git a/src/main/ai-vault-search/session-search-page-warmup.ts b/src/main/ai-vault-search/session-search-page-warmup.ts new file mode 100644 index 00000000000..e4be5d63cfd --- /dev/null +++ b/src/main/ai-vault-search/session-search-page-warmup.ts @@ -0,0 +1,27 @@ +import { setImmediate as yieldToEventLoop } from 'node:timers/promises' +import type SyncDatabase from '../sqlite/sync-database' + +const WARM_ROWS_PER_STEP = 50_000 + +/** + * Reads the messages table through in slices so its pages sit in the OS cache + * before the first query joins against it. Measured on a 4 GB index: the first + * query after a cold start drops from ~1.3 s to ~0.45 s, and each slice holds + * the connection for under 50 ms. + */ +export async function warmSessionSearchPages( + db: SyncDatabase, + stopped: () => boolean +): Promise { + const max = (db.prepare('SELECT max(id) AS id FROM messages').get() as { id: number | null }).id + const touch = db.prepare( + 'SELECT count(*) FROM messages WHERE id BETWEEN ? AND ? AND role IS NOT NULL' + ) + for (let low = 1; max !== null && low <= max; low += WARM_ROWS_PER_STEP) { + if (stopped()) { + return + } + touch.get(low, low + WARM_ROWS_PER_STEP - 1) + await yieldToEventLoop() + } +} diff --git a/src/main/ai-vault-search/session-search-path-key.ts b/src/main/ai-vault-search/session-search-path-key.ts new file mode 100644 index 00000000000..11bd35dcf05 --- /dev/null +++ b/src/main/ai-vault-search/session-search-path-key.ts @@ -0,0 +1,12 @@ +import { normalizeRuntimePathForComparison } from '../../shared/cross-platform-path' +import { parseWslUncPath, toWindowsWslPath } from '../../shared/wsl-paths' + +export function sessionSearchPathKey(cwd: string, transcriptPath?: string): string { + if (cwd.startsWith('\\') && !cwd.startsWith('\\\\')) { + cwd = cwd.replaceAll('\\', '/') + } + const wsl = transcriptPath ? parseWslUncPath(transcriptPath) : null + return normalizeRuntimePathForComparison( + wsl && cwd.startsWith('/') && !cwd.startsWith('//') ? toWindowsWslPath(cwd, wsl.distro) : cwd + ).replace(/\/+$/, '') +} 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..bb7630383a9 --- /dev/null +++ b/src/main/ai-vault-search/session-search-paths.ts @@ -0,0 +1,17 @@ +import { join } from 'node:path' + +// 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 databasePath: string | null = null + +export function initSessionSearchPaths(userDataPath: string): void { + databasePath = join(userDataPath, 'ai-vault-search', 'index.sqlite') +} + +export function getSessionSearchDatabasePath(): string | null { + return databasePath +} + +export function resetSessionSearchPathsForTests(): void { + databasePath = null +} diff --git a/src/main/ai-vault-search/session-search-pending-deletes.ts b/src/main/ai-vault-search/session-search-pending-deletes.ts new file mode 100644 index 00000000000..69de9f9f603 --- /dev/null +++ b/src/main/ai-vault-search/session-search-pending-deletes.ts @@ -0,0 +1,35 @@ +import type SyncDatabase from '../sqlite/sync-database' + +// Why the NUL prefix: tombstones share the file-path key space with real files, and +// `\0` cannot occur in one, so a synthetic key still gets the per-path cleanup mutex. +export function retireSearchSession(db: SyncDatabase, sessionId: number): void { + db.prepare('INSERT OR IGNORE INTO search_pending_deletes(path,session_row_id) VALUES (?,?)').run( + `\0session:${sessionId}`, + sessionId + ) +} + +export function discardSearchBatch( + db: SyncDatabase, + sessionId: number, + batchId: number, + ownsSession: boolean +): void { + if (ownsSession) { + retireSearchSession(db, sessionId) + } else { + db.prepare( + 'INSERT OR IGNORE INTO search_pending_deletes(path,session_row_id,batch_id) VALUES (?,?,?)' + ).run(`\0batch:${batchId}`, sessionId, batchId) + } +} + +/** Only called on open, before this store can have active writers. */ +export function recoverSearchWrites(db: SyncDatabase): void { + // A staging session or a batch row that outlived its writer is by definition unfinished: + // publish clears both in the same transaction that makes the rows visible. + db.exec(`INSERT OR IGNORE INTO search_pending_deletes(path,session_row_id) + SELECT char(0)||'session:'||id,id FROM sessions WHERE index_ready=0; + INSERT OR IGNORE INTO search_pending_deletes(path,session_row_id,batch_id) + SELECT char(0)||'batch:'||id,session_row_id,id FROM search_write_batches`) +} diff --git a/src/main/ai-vault-search/session-search-redaction.test.ts b/src/main/ai-vault-search/session-search-redaction.test.ts new file mode 100644 index 00000000000..a35ce485a88 --- /dev/null +++ b/src/main/ai-vault-search/session-search-redaction.test.ts @@ -0,0 +1,27 @@ +import { expect, it } from 'vitest' +import { redactSessionSearchText } from './session-search-redaction' + +const AWS_KEY = 'AKIAIOSFODNN7EXAMPLE' +const JWT = + 'eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiIxMjM0NTY3ODkwIn0.dozjgNryP4J3jVmNHl0w5N_XgL0n3I9PlFUP0THsR8U' + +it('keeps the provider fingerprints and labels what it removed', () => { + expect(redactSessionSearchText(`key ${AWS_KEY} here`)).toBe( + 'key [redacted:aws-access-key-id] here' + ) + expect(redactSessionSearchText(`Authorization: Bearer ${JWT}`)).toBe( + 'Authorization: Bearer [redacted:jwt]' + ) + // Opaque (non-JWT) bearer values are covered too; the label survives. + expect(redactSessionSearchText('Bearer abcdefghijklmnopqrstuvwxyz012345')).toBe( + 'Bearer [redacted:bearer-token]' + ) +}) + +it('leaves ordinary transcript shapes searchable', () => { + // Why these and not `redactString`: env-shaped code lines and `token:` prose + // are ordinary transcript content. + for (const benign of ['MAX_RETRIES = 3', 'the auth token: refreshed on 401', 'API_KEY_HEADER']) { + expect(redactSessionSearchText(benign)).toBe(benign) + } +}) diff --git a/src/main/ai-vault-search/session-search-redaction.ts b/src/main/ai-vault-search/session-search-redaction.ts new file mode 100644 index 00000000000..13c800a12c7 --- /dev/null +++ b/src/main/ai-vault-search/session-search-redaction.ts @@ -0,0 +1,26 @@ +import { PROVIDER_PATTERNS } from '../observability/redactor' + +// Why not `redactString`: its labeled-kv and .env-line rules are tuned for +// stack traces and are far too eager over transcript text — every pasted +// `MAX_RETRIES = 3` diff line would lose its value, and prose like +// `token: the next token` would lose the word `token` itself, taking the +// searchable content with it. The provider fingerprints are shape-matched and +// safe over prose, so those are reused verbatim, plus the one shape they miss: +// an opaque (non-JWT) bearer token. +const BEARER_TOKEN = /\b(Bearer)\s+[A-Za-z0-9._~+/=-]{16,}/g + +// One scan decides whether the eight fingerprint passes run at all: every +// pattern above is anchored on one of these, so text without them cannot match. +const SECRET_ANCHOR = /sk-|gh[pousr]_|AKIA|eyJ|xox|-----|[Bb]earer|aws_secret_access_key/ + +/** Strips credential-shaped spans so the index (and every snippet) never holds one. */ +export function redactSessionSearchText(text: string): string { + if (!SECRET_ANCHOR.test(text)) { + return text + } + let out = text + for (const { tag, re } of PROVIDER_PATTERNS) { + out = out.replace(re, `[redacted:${tag}]`) + } + return out.replace(BEARER_TOKEN, '$1 [redacted:bearer-token]') +} diff --git a/src/main/ai-vault-search/session-search-retention-delete.ts b/src/main/ai-vault-search/session-search-retention-delete.ts new file mode 100644 index 00000000000..2695e916f3b --- /dev/null +++ b/src/main/ai-vault-search/session-search-retention-delete.ts @@ -0,0 +1,97 @@ +import { setImmediate as yieldToEventLoop } from 'node:timers/promises' +import type SyncDatabase from '../sqlite/sync-database' +import { inSessionParseFileLane } from '../ai-vault/session-parse-file-lane' + +export const RETENTION_DELETE_ROWS_PER_STEP = 256 + +/** A durable tombstone hides partial deletes and lets a reopened store finish them. */ +export async function deleteExpiredSearchFiles( + db: SyncDatabase, + cutoffMs: number | null, + closed: () => boolean, + changed: () => void, + yieldStep: () => Promise = yieldToEventLoop +): Promise { + const pending = db.prepare('SELECT path FROM search_pending_deletes').all() as { path: string }[] + const expired = + cutoffMs === null + ? [] + : (db + .prepare('SELECT path FROM files WHERE mtime_ms < ? ORDER BY mtime_ms') + .all(cutoffMs) as { path: string }[]) + for (const { path } of [...pending, ...expired]) { + if (closed()) { + return + } + await inSessionParseFileLane(path, async () => { + if (closed()) { + return + } + db.exec('BEGIN IMMEDIATE') + try { + const file = db + .prepare(`SELECT session_row_id FROM files WHERE path = ? AND mtime_ms < ? + AND path NOT IN (SELECT path FROM search_pending_deletes)`) + .get(path, cutoffMs ?? -Infinity) as { session_row_id: number | null } | undefined + if (file) { + if (file.session_row_id !== null) { + db.prepare( + 'INSERT OR IGNORE INTO search_pending_deletes(path, session_row_id) VALUES (?, ?)' + ).run(path, file.session_row_id) + } + db.prepare('DELETE FROM files WHERE path = ?').run(path) + } + db.exec('COMMIT') + changed() + } catch (error) { + db.exec('ROLLBACK') + throw error + } + while (!closed()) { + const pending = db + .prepare('SELECT session_row_id,batch_id FROM search_pending_deletes WHERE path = ?') + .get(path) as { session_row_id: number; batch_id: number | null } | undefined + if (!pending) { + return + } + db.exec('BEGIN IMMEDIATE') + try { + const ids = db + .prepare( + pending.batch_id === null + ? 'SELECT id FROM messages WHERE session_row_id = ? LIMIT ?' + : 'SELECT id FROM messages WHERE batch_id = ? LIMIT ?' + ) + .all(pending.batch_id ?? pending.session_row_id, RETENTION_DELETE_ROWS_PER_STEP) as { + id: number + }[] + const full = db.prepare('DELETE FROM messages_fts WHERE rowid = ?') + const conversation = db.prepare('DELETE FROM conversation_fts WHERE rowid = ?') + const message = db.prepare('DELETE FROM messages WHERE id = ?') + for (const { id } of ids) { + full.run(id) + conversation.run(id) + message.run(id) + } + if (ids.length < RETENTION_DELETE_ROWS_PER_STEP) { + if (pending.batch_id === null) { + db.prepare('DELETE FROM search_write_batches WHERE session_row_id=?').run( + pending.session_row_id + ) + db.prepare('DELETE FROM sessions WHERE id = ?').run(pending.session_row_id) + } else { + db.prepare('DELETE FROM search_write_batches WHERE id=?').run(pending.batch_id) + } + db.prepare('DELETE FROM search_pending_deletes WHERE path = ?').run(path) + } + db.exec('COMMIT') + changed() + } catch (error) { + db.exec('ROLLBACK') + throw error + } + await yieldStep() + } + }) + } +} diff --git a/src/main/ai-vault-search/session-search-schema.test.ts b/src/main/ai-vault-search/session-search-schema.test.ts new file mode 100644 index 00000000000..b58a0b44557 --- /dev/null +++ b/src/main/ai-vault-search/session-search-schema.test.ts @@ -0,0 +1,191 @@ +import type * as NodeFs from 'node:fs' +import { mkdtemp, stat, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it, vi } from 'vitest' +import { + removeTree, + WINDOWS_RM_MAX_RETRIES, + WINDOWS_RM_RETRY_DELAY_MS +} from '../../shared/windows-transient-lock-removal' +import SyncDatabase from '../sqlite/sync-database' +import { + SESSION_SEARCH_SCHEMA_VERSION, + openSessionSearchDatabase, + removeSessionSearchDatabase, + VISIBLE_MESSAGES, + VISIBLE_SESSIONS +} from './session-search-schema' + +const recordedRmSync = vi.hoisted(() => vi.fn()) +vi.mock('node:fs', async () => { + const actual = await vi.importActual('node:fs') + return { + ...actual, + rmSync: (...args: Parameters) => { + recordedRmSync(...args) + return actual.rmSync(...args) + } + } +}) + +let roots: string[] = [] + +afterEach(async () => { + await Promise.all(roots.map((root) => removeTree(root))) + roots = [] +}) + +async function tempDatabasePath(): Promise { + const root = await mkdtemp(join(tmpdir(), 'orca-session-search-schema-')) + roots.push(root) + return join(root, 'index.sqlite') +} + +function schemaVersion(db: SyncDatabase): string | undefined { + return ( + db.prepare("SELECT value FROM meta WHERE key = 'schema_version'").get() as + | { value: string } + | undefined + )?.value +} + +describe('openSessionSearchDatabase', () => { + it('keeps a current-version index and its rows', async () => { + const path = await tempDatabasePath() + const first = openSessionSearchDatabase(path) + first.prepare("INSERT INTO files(path,byte_offset,mtime_ms) VALUES ('a',1,1)").run() + first.close() + + const second = openSessionSearchDatabase(path) + expect(schemaVersion(second)).toBe(String(SESSION_SEARCH_SCHEMA_VERSION)) + expect(second.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ c: 1 }) + second.close() + }) + + it('replaces the file on a version mismatch instead of dropping tables in place', async () => { + const path = await tempDatabasePath() + const stale = openSessionSearchDatabase(path) + stale.prepare("INSERT INTO files(path,byte_offset,mtime_ms) VALUES ('a',1,1)").run() + stale + .prepare("UPDATE meta SET value = ? WHERE key = 'schema_version'") + .run(String(SESSION_SEARCH_SCHEMA_VERSION + 1)) + stale.close() + // Why: a stale sidecar must go with the main file, or SQLite replays it into the new one. + await writeFile(`${path}-wal`, 'stale wal bytes') + const before = await stat(path) + + const fresh = openSessionSearchDatabase(path) + expect(schemaVersion(fresh)).toBe(String(SESSION_SEARCH_SCHEMA_VERSION)) + expect(fresh.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ c: 0 }) + fresh.close() + // Why not inode: ext4 hands a freed inode straight back to the next create. + // The planted sidecar is gone (a fresh WAL is checkpointed away on close). + await expect(stat(`${path}-wal`)).rejects.toMatchObject({ code: 'ENOENT' }) + expect((await stat(path)).mtimeMs).toBeGreaterThanOrEqual(before.mtimeMs) + }) + + it('indexes only in-flight batch pointers, not every published message', async () => { + const db = openSessionSearchDatabase(await tempDatabasePath()) + try { + const sql = ( + db + .prepare("SELECT sql FROM sqlite_master WHERE type='index' AND name='messages_batch'") + .get() as { sql: string } + ).sql + // Publish nulls batch_id, so a full index would carry one dead entry per message. + expect(sql).toContain('WHERE batch_id IS NOT NULL') + } finally { + db.close() + } + }) + + it('removes the database with every sidecar', async () => { + const path = await tempDatabasePath() + openSessionSearchDatabase(path).close() + await writeFile(`${path}-shm`, '') + removeSessionSearchDatabase(path) + for (const suffix of ['', '-wal', '-shm']) { + await expect(stat(`${path}${suffix}`)).rejects.toMatchObject({ code: 'ENOENT' }) + } + }) +}) + +it('closes the SQLite handle when corrupt data fails initialization', async () => { + const path = await tempDatabasePath() + await writeFile(path, 'not a SQLite database') + const close = vi.spyOn(SyncDatabase.prototype, 'close') + try { + expect(() => openSessionSearchDatabase(path)).toThrow() + expect(close).toHaveBeenCalledTimes(1) + } finally { + close.mockRestore() + } + removeSessionSearchDatabase(path) + const recovered = openSessionSearchDatabase(path) + recovered.close() +}) + +describe('visibility views', () => { + it('hides a staging session, a tombstoned session and an in-flight batch', async () => { + const db = openSessionSearchDatabase(await tempDatabasePath()) + try { + db.exec(`INSERT INTO sessions(id,index_ready,agent,session_id,file_path,title,resume_command) + VALUES (1,1,'claude','a','a','published',''),(2,0,'claude','b','b','staging',''), + (3,1,'claude','c','c','tombstoned',''); + INSERT INTO search_pending_deletes(path,session_row_id) VALUES ('c',3); + INSERT INTO search_write_batches(id,session_row_id) VALUES (7,1); + INSERT INTO messages(id,session_row_id,batch_id,role) VALUES (1,1,NULL,'user'),(2,1,7,'user')`) + expect(db.prepare(`SELECT title FROM ${VISIBLE_SESSIONS} ORDER BY id`).all()).toEqual([ + { title: 'published' } + ]) + expect(db.prepare(`SELECT id FROM ${VISIBLE_MESSAGES} ORDER BY id`).all()).toEqual([ + { id: 1 } + ]) + // Publish clears the pointer, so visibility never depends on the batch row surviving. + db.exec( + 'UPDATE messages SET batch_id=NULL WHERE batch_id=7; DELETE FROM search_write_batches' + ) + expect(db.prepare(`SELECT count(*) AS n FROM ${VISIBLE_MESSAGES}`).get()).toEqual({ n: 2 }) + } finally { + db.close() + } + }) +}) + +it("retries a Windows lock that outlives rmSync's own retries", async () => { + const path = await tempDatabasePath() + openSessionSearchDatabase(path).close() + vi.spyOn(process, 'platform', 'get').mockReturnValue('win32') + recordedRmSync.mockReset() + const locked = Object.assign(new Error('EPERM: operation not permitted'), { code: 'EPERM' }) + recordedRmSync.mockImplementationOnce(() => { + throw locked + }) + try { + expect(() => removeSessionSearchDatabase(path)).not.toThrow() + expect(recordedRmSync.mock.calls.length).toBe(5) + await expect(stat(path)).rejects.toMatchObject({ code: 'ENOENT' }) + } finally { + recordedRmSync.mockReset() + vi.restoreAllMocks() + } +}) + +it('gives Windows the shared retry options for a late handle release', async () => { + const path = await tempDatabasePath() + vi.spyOn(process, 'platform', 'get').mockReturnValue('win32') + recordedRmSync.mockClear() + try { + removeSessionSearchDatabase(path) + expect(recordedRmSync).toHaveBeenCalled() + for (const [, options] of recordedRmSync.mock.calls) { + expect(options).toMatchObject({ + maxRetries: WINDOWS_RM_MAX_RETRIES, + retryDelay: WINDOWS_RM_RETRY_DELAY_MS + }) + } + } finally { + vi.restoreAllMocks() + } +}) 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..98b8ab5ad48 --- /dev/null +++ b/src/main/ai-vault-search/session-search-schema.ts @@ -0,0 +1,144 @@ +import SyncDatabase from '../sqlite/sync-database' +import { removeTreeSync } from '../../shared/windows-transient-lock-removal' + +// Bump to drop and rebuild: the index is a cache over the transcripts, never a source. +export const SESSION_SEARCH_SCHEMA_VERSION = 1 + +// unicode61 keeps `_ . - /` inside tokens so paths and identifiers match exactly; +// the `identifiers` column carries the split form (see session-search-identifier-split). +// Why: `+` keeps `C++` a token of its own instead of the letter `c`; `#` is +// left out so `#123` still answers a search for `123`. +const TOKENIZER = `tokenize="unicode61 tokenchars '_.-/+'"` + +/** Sessions and messages a read may return: published, not tombstoned. */ +export const VISIBLE_SESSIONS = 'visible_sessions' +export const VISIBLE_MESSAGES = 'visible_messages' + +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, + index_ready INTEGER NOT NULL DEFAULT 1, + 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, + cwd_key TEXT, + branch TEXT, + created_at TEXT, + updated_at TEXT, + message_count INTEGER NOT NULL DEFAULT 0, + resume_command TEXT NOT NULL, + -- Chained digest of the first N messages; forks of one conversation share it. + content_hash TEXT, + content_hash_count INTEGER NOT NULL DEFAULT 0 +); +CREATE INDEX IF NOT EXISTS sessions_agent ON sessions(agent); +CREATE INDEX IF NOT EXISTS sessions_content_hash ON sessions(content_hash); +CREATE INDEX IF NOT EXISTS sessions_updated_at ON sessions(updated_at); +CREATE INDEX IF NOT EXISTS sessions_cwd_key ON sessions(cwd_key); +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 search_pending_deletes( + path TEXT PRIMARY KEY, + session_row_id INTEGER NOT NULL, + batch_id INTEGER +); +-- A row exists only while its batch is in flight; publish clears its messages and deletes it. +CREATE TABLE IF NOT EXISTS search_write_batches( + id INTEGER PRIMARY KEY, + session_row_id INTEGER NOT NULL +); +CREATE INDEX IF NOT EXISTS search_write_batches_session ON search_write_batches(session_row_id); +CREATE TABLE IF NOT EXISTS messages( + id INTEGER PRIMARY KEY, + session_row_id INTEGER NOT NULL, + batch_id INTEGER, + role TEXT NOT NULL, + ts TEXT +); +CREATE INDEX IF NOT EXISTS messages_session ON messages(session_row_id); +-- Partial: publish nulls batch_id, so all but the in-flight rows would be dead entries. +CREATE INDEX IF NOT EXISTS messages_batch ON messages(batch_id) WHERE batch_id IS NOT NULL; +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 +); +-- Why: staged rows must never reach a result. One definition per half, so a new +-- read site cannot forget one; SQLite flattens both into the caller's plan. +CREATE VIEW IF NOT EXISTS ${VISIBLE_SESSIONS} AS SELECT * FROM sessions + WHERE index_ready = 1 + AND id NOT IN (SELECT session_row_id FROM search_pending_deletes WHERE batch_id IS NULL); +CREATE VIEW IF NOT EXISTS ${VISIBLE_MESSAGES} AS SELECT * FROM messages + WHERE batch_id IS NULL; +` + +export function openSessionSearchDatabase(path: string): SyncDatabase { + let db = openWithPragmas(path) + const version = readSchemaVersion(db) + if (version !== null && version !== SESSION_SEARCH_SCHEMA_VERSION) { + // Why: DROP TABLE on a multi-GB FTS index takes minutes and runs inside the + // scanner service's init, past its ready timeout; unlinking is instant. + db.close() + removeSessionSearchDatabase(path) + db = openWithPragmas(path) + } + db.exec(SCHEMA_SQL) + db.prepare('INSERT OR REPLACE INTO meta(key, value) VALUES (?, ?)').run( + 'schema_version', + String(SESSION_SEARCH_SCHEMA_VERSION) + ) + return db +} + +function openWithPragmas(path: string): SyncDatabase { + const db = new SyncDatabase(path) + try { + // Why: only takes effect on an empty file; it is what lets a purge hand pages + // back in bounded steps instead of a full VACUUM. Set before any table exists. + db.pragma('auto_vacuum = INCREMENTAL') + db.pragma('journal_mode = WAL') + db.pragma('synchronous = NORMAL') + db.pragma('journal_size_limit = 8388608') + db.pragma('busy_timeout = 5000') + return db + } catch (error) { + db.close() + throw error + } +} + +export function removeSessionSearchDatabase(path: string): void { + if (path === ':memory:') { + return + } + for (const suffix of ['', '-wal', '-shm', '-journal']) { + removeTreeSync(`${path}${suffix}`) + } +} + +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-staged-write-test-fixture.ts b/src/main/ai-vault-search/session-search-staged-write-test-fixture.ts new file mode 100644 index 00000000000..baa3e2976fe --- /dev/null +++ b/src/main/ai-vault-search/session-search-staged-write-test-fixture.ts @@ -0,0 +1,124 @@ +import { mkdtemp } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { removeTree } from '../../shared/windows-transient-lock-removal' +import type { AiVaultSession } from '../../shared/ai-vault-types' +import type { SessionFileCandidate } from '../ai-vault/session-scanner-types' +import { TranscriptMessageChannel } from '../ai-vault/session-transcript-channel' +import type { + TranscriptMessage, + TranscriptReadOutcome, + TranscriptReadStart +} from '../ai-vault/session-transcript-consumers' +import type SyncDatabase from '../sqlite/sync-database' +import { openSessionSearchDatabase } from './session-search-schema' + +export const SYNTHETIC_TRANSCRIPT = 'synthetic-transcript' + +export function syntheticCandidate( + overrides: Partial = {} +): SessionFileCandidate { + const at = new Date(1740000000000) + return { + agent: 'claude', + codexHome: null, + file: { + path: SYNTHETIC_TRANSCRIPT, + mtimeMs: at.getTime(), + modifiedAt: at.toISOString(), + sizeBytes: 4096, + ...overrides + } + } +} + +export function syntheticSession(overrides: Partial = {}): AiVaultSession { + const at = new Date(1740000000000).toISOString() + return { + id: 'fixture', + executionHostId: 'local', + agent: 'claude', + sessionId: 'fixture', + title: 'fixture session', + cwd: '/fixture', + branch: null, + model: null, + filePath: SYNTHETIC_TRANSCRIPT, + codexHome: null, + createdAt: at, + updatedAt: at, + modifiedAt: at, + messageCount: 0, + totalTokens: 0, + previewMessages: [], + queuedMessageCount: 0, + subagentTranscriptCount: 0, + resumeCommand: '', + subagent: null, + ...overrides + } +} + +export function userMessages(text: string, count: number): TranscriptMessage[] { + return Array.from({ length: count }, (_unused, index) => ({ + role: 'user' as const, + text, + timestamp: new Date(1740000000000 + index * 1000).toISOString() + })) +} + +/** + * Drives one read through the real fan-out channel, so a test exercises the + * registration path the transcript reader uses rather than the consumer alone. + */ +export function replayTranscriptRead(args: { + candidate?: SessionFileCandidate + mode?: TranscriptReadStart['mode'] + previousByteOffset?: number + messages: TranscriptMessage[] + outcome?: Partial +}): void { + const candidate = args.candidate ?? syntheticCandidate() + const mode = args.mode ?? 'replace' + const channel = new TranscriptMessageChannel() + channel.beginRead({ + candidate, + mode, + previousByteOffset: args.previousByteOffset ?? 0 + }) + for (const message of args.messages) { + channel.push(message) + } + channel.finishRead({ + session: syntheticSession(), + byteOffset: 4096, + incomplete: false, + ...args.outcome + }) +} + +export type SessionSearchIndexFile = { + path: string + /** The store keeps its own connection private, so row assertions need this one. */ + db: SyncDatabase + close: () => Promise +} + +/** An on-disk index: `:memory:` is per-connection, so a second reader needs a real file. */ +export async function openSessionSearchIndexFile(name: string): Promise { + const root = await mkdtemp(join(tmpdir(), `${name}-`)) + const path = join(root, 'index.sqlite') + const db = openSessionSearchDatabase(path) + let open = true + return { + path, + db, + close: async () => { + if (open) { + open = false + db.close() + } + await removeTree(root) + } + } +} diff --git a/src/main/ai-vault-search/session-search-synthetic-corpus.ts b/src/main/ai-vault-search/session-search-synthetic-corpus.ts new file mode 100644 index 00000000000..f87e64ad8a6 --- /dev/null +++ b/src/main/ai-vault-search/session-search-synthetic-corpus.ts @@ -0,0 +1,132 @@ +import { mkdtemp, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' + +// Why synthetic and in-repo: the cost model has to be reproducible on any host +// and must never read a real transcript. The shapes here mirror what a Claude +// JSONL transcript actually holds — prose turns, a pasted diff, tool calls and +// their output — because the index's disk cost tracks the mix, not the size. + +const WORDS = [ + 'terminal', + 'reattach', + 'worktree', + 'resolveTerminalPath', + 'src/main/ai-vault/session-transcript-reader.ts', + 'the', + 'index', + 'cursor', + 'byteOffset', + 'publish', + 'staged', + 'transaction', + 'MAX_RETRIES', + 'relay', + 'daemon', + 'pty', + 'snapshot', + 'because' +] + +/** Deterministic: the same seed gives the same corpus on every host and run. */ +function mulberry32(seed: number): () => number { + let state = seed >>> 0 + return () => { + state = (state + 0x6d2b79f5) >>> 0 + let t = Math.imul(state ^ (state >>> 15), 1 | state) + t = (t + Math.imul(t ^ (t >>> 7), 61 | t)) ^ t + return ((t ^ (t >>> 14)) >>> 0) / 4294967296 + } +} + +function words(random: () => number, count: number): string { + const out: string[] = [] + for (let index = 0; index < count; index++) { + out.push(WORDS[Math.floor(random() * WORDS.length)]) + } + return out.join(' ') +} + +export type SyntheticCorpus = { + root: string + files: string[] + /** Total bytes of transcript written, the denominator of write amplification. */ + transcriptBytes: number + messageCount: number +} + +export type SyntheticCorpusOptions = { + sessions?: number + turnsPerSession?: number + seed?: number +} + +/** Writes a corpus of Claude JSONL transcripts and reports what it cost on disk. */ +export async function writeSyntheticTranscriptCorpus( + options: SyntheticCorpusOptions = {} +): Promise { + const sessions = options.sessions ?? 40 + const turns = options.turnsPerSession ?? 60 + const random = mulberry32(options.seed ?? 1) + const root = await mkdtemp(join(tmpdir(), 'orca-search-corpus-')) + const files: string[] = [] + let transcriptBytes = 0 + let messageCount = 0 + + for (let session = 0; session < sessions; session++) { + const sessionId = `00000000-0000-4000-8000-${String(session).padStart(12, '0')}` + const lines: string[] = [] + for (let turn = 0; turn < turns; turn++) { + const at = new Date(1740000000000 + turn * 60_000).toISOString() + lines.push( + JSON.stringify({ + type: 'user', + sessionId, + timestamp: at, + cwd: `/repo/app-${session % 7}`, + gitBranch: 'main', + message: { role: 'user', content: words(random, 40) } + }) + ) + lines.push( + JSON.stringify({ + type: 'assistant', + sessionId, + timestamp: at, + message: { + role: 'assistant', + model: 'claude-fable-5', + content: [ + { type: 'text', text: words(random, 120) }, + { + type: 'tool_use', + name: 'Bash', + input: { command: `rg ${words(random, 3)}` } + } + ] + } + }) + ) + lines.push( + JSON.stringify({ + type: 'user', + sessionId, + timestamp: at, + message: { + role: 'user', + content: [{ type: 'tool_result', tool_use_id: 'toolu_1', content: words(random, 200) }] + } + }) + ) + // One user turn, one assistant turn, one tool call, one tool result. + messageCount += 4 + } + const path = join(root, `${sessionId}.jsonl`) + const body = `${lines.join('\n')}\n` + await writeFile(path, body) + transcriptBytes += Buffer.byteLength(body) + files.push(path) + } + + return { root, files, transcriptBytes, messageCount } +} diff --git a/src/main/ai-vault-search/session-search-transcript-fixtures.ts b/src/main/ai-vault-search/session-search-transcript-fixtures.ts new file mode 100644 index 00000000000..c83b7be07e1 --- /dev/null +++ b/src/main/ai-vault-search/session-search-transcript-fixtures.ts @@ -0,0 +1,111 @@ +import { stat } from 'node:fs/promises' +import { + createSessionParseStats, + parseAgentSessionFileCached, + type SessionParseStats +} from '../ai-vault/session-scanner-parse-cache' +import type { SessionFileCandidate } from '../ai-vault/session-scanner-types' + +// Transcript builders shared by the session-search store tests; each file owns +// its temp directories, this module only shapes records and drives the parser. + +export const CLAUDE_SESSION_ID = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee' +export const CODEX_SESSION_ID = '019f0000-1111-7222-8333-444444444444' +export const CODEX_ROLLOUT_FILE = `rollout-2026-05-01T10-00-00-${CODEX_SESSION_ID}.jsonl` + +const RECORD_EPOCH_MS = 1740000000000 + +export function recordTimestamp(index: number): string { + return new Date(RECORD_EPOCH_MS + index * 60_000).toISOString() +} + +export function userRecord( + index: number, + content: unknown, + sessionId = CLAUDE_SESSION_ID, + cwd = '/repo/app' +): string { + return JSON.stringify({ + type: 'user', + sessionId, + timestamp: recordTimestamp(index), + cwd, + gitBranch: 'main', + message: { role: 'user', content } + }) +} + +export function assistantRecord( + index: number, + content: unknown, + sessionId = CLAUDE_SESSION_ID +): string { + return JSON.stringify({ + type: 'assistant', + sessionId, + timestamp: recordTimestamp(index), + message: { role: 'assistant', model: 'claude-fable-5', content } + }) +} + +export async function sessionCandidate( + agent: SessionFileCandidate['agent'], + path: string, + codexHome: string | null = null +): Promise { + const fileStat = await stat(path) + return { + agent, + codexHome, + file: { + path, + mtimeMs: fileStat.mtimeMs, + modifiedAt: fileStat.mtime.toISOString(), + sizeBytes: fileStat.size, + dev: fileStat.dev, + ino: fileStat.ino + } + } +} + +export async function parseTranscript( + path: string, + agent: SessionFileCandidate['agent'] = 'claude', + codexHome: string | null = null +): Promise<{ stats: SessionParseStats }> { + const stats = createSessionParseStats() + await parseAgentSessionFileCached( + await sessionCandidate(agent, path, codexHome), + process.platform, + stats + ) + return { stats } +} + +function codexLine(record: Record): string { + return JSON.stringify(record) +} + +/** Minimal Codex rollout: meta, one user message, one completed shell command. */ +export function codexRolloutLines(command: string[], output: string, prompt: string): string[] { + return [ + codexLine({ + timestamp: recordTimestamp(0), + type: 'session_meta', + payload: { id: CODEX_SESSION_ID, cwd: '/repo/app', git: { branch: 'main' } } + }), + codexLine({ + timestamp: recordTimestamp(1), + type: 'response_item', + payload: { type: 'message', role: 'user', content: prompt } + }), + codexLine({ + timestamp: recordTimestamp(2), + type: 'event_msg', + payload: { + type: 'item_completed', + item: { type: 'CommandExecution', command, aggregated_output: output } + } + }) + ] +} diff --git a/src/main/ai-vault-search/session-search-wal-budget.ts b/src/main/ai-vault-search/session-search-wal-budget.ts new file mode 100644 index 00000000000..c2bac9771a7 --- /dev/null +++ b/src/main/ai-vault-search/session-search-wal-budget.ts @@ -0,0 +1,24 @@ +import type SyncDatabase from '../sqlite/sync-database' + +export const SEARCH_WAL_PENDING_BYTES = 64 * 1024 * 1024 +export class SearchWalBackpressureError extends Error { + constructor() { + super('Session search indexing is waiting for an older index reader to finish.') + this.name = 'SearchWalBackpressureError' + } +} + +/** PASSIVE never waits for readers; a blocked checkpoint must not grow without a bound. */ +export function assertSearchWalBudget( + db: SyncDatabase, + limitBytes = SEARCH_WAL_PENDING_BYTES +): void { + const row = (db.pragma('wal_checkpoint(PASSIVE)') as { log: number; checkpointed: number }[])[0] + if (!row || row.log < 0) { + return + } + const pageSize = Number(db.pragma('page_size', { simple: true })) + if ((row.log - row.checkpointed) * pageSize >= limitBytes) { + throw new SearchWalBackpressureError() + } +} diff --git a/src/main/observability/redactor.ts b/src/main/observability/redactor.ts index 7ddab42a0c9..8f8b7eedc56 100644 --- a/src/main/observability/redactor.ts +++ b/src/main/observability/redactor.ts @@ -14,7 +14,7 @@ const LABELED_KV = /\b(?:api[-_]?key|token|secret|password|bearer|authorization)\b\s*[:=]\s*(?:Bearer\s+\S+|Token\s+\S+|\S+)/gi // Tagged tokens let triage see what was redacted without the key. Order is most-specific-first: `sk-ant-` before `sk-`, or the Anthropic tag is lost. -const PROVIDER_PATTERNS: { tag: string; re: RegExp }[] = [ +export const PROVIDER_PATTERNS: { tag: string; re: RegExp }[] = [ { tag: 'anthropic-key', re: /sk-ant-[a-zA-Z0-9_-]{40,}/g }, { tag: 'openai-key', re: /sk-(?:proj-)?[a-zA-Z0-9_-]{32,}/g }, { tag: 'github-token', re: /gh[pousr]_[A-Za-z0-9]{36,}/g },