mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
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.
This commit is contained in:
@@ -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)
|
||||
})
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
)
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
@@ -0,0 +1,54 @@
|
||||
// Identifier shadow terms: `resolveTerminalPath` → `resolve terminal path`,
|
||||
// `src/main/foo-bar.ts` → `src main foo bar ts`. Stored in a separate FTS5
|
||||
// column so a partial identifier still matches; the largest single accuracy
|
||||
// win measured in the retrieval shoot-out (MRR 0.50 → 0.55).
|
||||
|
||||
const RAW_TOKEN = /[A-Za-z0-9_./-]+/g
|
||||
const CAMEL_PIECE = /[A-Z]+(?![a-z])|[A-Z][a-z0-9]*|[a-z0-9]+/g
|
||||
const SEPARATOR = /[_./-]+/
|
||||
// Worth shadowing: has a separator, a camel boundary, or is SCREAMING_CASE.
|
||||
const INTERESTING = /[_./-]|[a-z0-9][A-Z]|^[A-Z]{2,}[0-9_]*$/
|
||||
const MIN_TOKEN = 3
|
||||
const MAX_TOKEN = 120
|
||||
const MIN_PIECE = 2
|
||||
|
||||
function hasMixedCase(piece: string): boolean {
|
||||
return /[a-z]/.test(piece) && /[A-Z]/.test(piece)
|
||||
}
|
||||
|
||||
export function identifierShadowTerms(text: string, limit = 4000): string[] {
|
||||
const out: string[] = []
|
||||
const seen = new Set<string>()
|
||||
for (const match of text.matchAll(RAW_TOKEN)) {
|
||||
const token = match[0]
|
||||
if (token.length < MIN_TOKEN || token.length > MAX_TOKEN || !INTERESTING.test(token)) {
|
||||
continue
|
||||
}
|
||||
const parts: string[] = []
|
||||
for (const piece of token.split(SEPARATOR)) {
|
||||
if (!piece) {
|
||||
continue
|
||||
}
|
||||
parts.push(piece)
|
||||
if (hasMixedCase(piece)) {
|
||||
parts.push(...(piece.match(CAMEL_PIECE) ?? []))
|
||||
}
|
||||
}
|
||||
for (const part of parts) {
|
||||
const lowered = part.toLowerCase()
|
||||
if (lowered.length < MIN_PIECE || seen.has(lowered)) {
|
||||
continue
|
||||
}
|
||||
seen.add(lowered)
|
||||
out.push(lowered)
|
||||
if (out.length >= limit) {
|
||||
return out
|
||||
}
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
export function identifierShadowText(text: string, limit?: number): string {
|
||||
return identifierShadowTerms(text, limit).join(' ')
|
||||
}
|
||||
@@ -0,0 +1,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<void> {
|
||||
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()
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
})
|
||||
@@ -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<string> {
|
||||
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<TranscriptMessage>
|
||||
): Generator<TranscriptMessage> {
|
||||
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)]
|
||||
}
|
||||
@@ -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<void> {
|
||||
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()
|
||||
}
|
||||
}
|
||||
@@ -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(/\/+$/, '')
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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`)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
@@ -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]')
|
||||
}
|
||||
@@ -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<void> = yieldToEventLoop
|
||||
): Promise<void> {
|
||||
const pending = db.prepare('SELECT path FROM search_pending_deletes').all() as { path: string }[]
|
||||
const expired =
|
||||
cutoffMs === null
|
||||
? []
|
||||
: (db
|
||||
.prepare('SELECT path FROM files WHERE mtime_ms < ? ORDER BY mtime_ms')
|
||||
.all(cutoffMs) as { path: string }[])
|
||||
for (const { path } of [...pending, ...expired]) {
|
||||
if (closed()) {
|
||||
return
|
||||
}
|
||||
await inSessionParseFileLane(path, async () => {
|
||||
if (closed()) {
|
||||
return
|
||||
}
|
||||
db.exec('BEGIN IMMEDIATE')
|
||||
try {
|
||||
const file = db
|
||||
.prepare(`SELECT session_row_id FROM files WHERE path = ? AND mtime_ms < ?
|
||||
AND path NOT IN (SELECT path FROM search_pending_deletes)`)
|
||||
.get(path, cutoffMs ?? -Infinity) as { session_row_id: number | null } | undefined
|
||||
if (file) {
|
||||
if (file.session_row_id !== null) {
|
||||
db.prepare(
|
||||
'INSERT OR IGNORE INTO search_pending_deletes(path, session_row_id) VALUES (?, ?)'
|
||||
).run(path, file.session_row_id)
|
||||
}
|
||||
db.prepare('DELETE FROM files WHERE path = ?').run(path)
|
||||
}
|
||||
db.exec('COMMIT')
|
||||
changed()
|
||||
} catch (error) {
|
||||
db.exec('ROLLBACK')
|
||||
throw error
|
||||
}
|
||||
while (!closed()) {
|
||||
const pending = db
|
||||
.prepare('SELECT session_row_id,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()
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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<typeof NodeFs>('node:fs')
|
||||
return {
|
||||
...actual,
|
||||
rmSync: (...args: Parameters<typeof actual.rmSync>) => {
|
||||
recordedRmSync(...args)
|
||||
return actual.rmSync(...args)
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
let roots: string[] = []
|
||||
|
||||
afterEach(async () => {
|
||||
await Promise.all(roots.map((root) => removeTree(root)))
|
||||
roots = []
|
||||
})
|
||||
|
||||
async function tempDatabasePath(): Promise<string> {
|
||||
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()
|
||||
}
|
||||
})
|
||||
@@ -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
|
||||
}
|
||||
@@ -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['file']> = {}
|
||||
): 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> = {}): 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<TranscriptReadOutcome>
|
||||
}): 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<void>
|
||||
}
|
||||
|
||||
/** An on-disk index: `:memory:` is per-connection, so a second reader needs a real file. */
|
||||
export async function openSessionSearchIndexFile(name: string): Promise<SessionSearchIndexFile> {
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<SyntheticCorpus> {
|
||||
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 }
|
||||
}
|
||||
@@ -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<SessionFileCandidate> {
|
||||
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, unknown>): 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 }
|
||||
}
|
||||
})
|
||||
]
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
@@ -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 },
|
||||
|
||||
Reference in New Issue
Block a user