From c00a8b9fa203d6df85b1269cc33db2ad1d5dfae3 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Thu, 10 Sep 2026 14:20:30 -0400 Subject: [PATCH] refactor(ai-vault-search): commit a file's rows and its cursor in one transaction The staged-publish model is replaced by one SQLite transaction per file in WAL mode. A read buffers its decoded rows and writes them, its session and its cursor together; a reader on another handle sees the last committed state, which is the "never a torn session" guarantee the staging machinery was built to provide. A crash rolls the whole file back and it is re-read. Gone with it: `search_write_batches`, `search_pending_deletes`, `messages.batch_id`, `sessions.index_ready`, both `visible_*` views, the open-time recovery pass, the cleanup lane and `settled()`, and the ratchet test that made every query site read through a view. Rounds 2 through 6 were all seams between those pieces. A file whose rows exceed `SESSION_SEARCH_COMMIT_CHARS` is cut into chunks. The reader only hands out a byte offset when a read finishes, so a chunk records one no append can continue from: its rows answer searches as a coherent prefix of the session, and the next whole read replaces them. Retention keeps no record of unfinished work. It cuts a session loose from its file in one small transaction, which is what stops it answering, then reclaims its rows in bounded batches and hands the freed pages back as it goes. Rows whose session row is gone are the record of what an interrupted purge left. Deleted for want of a caller in PR 2, 3 or 4: the WAL budget (a buffered write has no staging window to grow one across), the compaction module (folded into the drain), the `sessions_content_hash` index (fork folding runs in JS over rows PR 4 already holds), `lastWriteAt`, `failures`, `SessionSearchStoreOptions`, `openStageCount` and `chunkMessageText`. --- .../session-search-content-hash.test.ts | 2 +- .../session-search-file-records.ts | 20 +- .../session-search-file-write.test.ts | 328 ++++++++++++++++++ .../session-search-index-compaction.ts | 22 -- .../session-search-index-consumer.test.ts | 170 +++++---- .../session-search-index-consumer.ts | 39 +-- ...s => session-search-index-test-fixture.ts} | 0 .../session-search-index-writer.test.ts | 132 +++---- .../session-search-index-writer.ts | 294 +++++++--------- .../session-search-live-transcript.test.ts | 10 +- .../session-search-message-rows.test.ts | 27 +- .../session-search-message-rows.ts | 31 +- .../session-search-pending-deletes.ts | 40 --- .../session-search-retention-delete.test.ts | 164 +++------ .../session-search-retention-delete.ts | 145 ++++---- .../session-search-schema.test.ts | 157 +++------ .../ai-vault-search/session-search-schema.ts | 47 +-- .../session-search-staged-write.test.ts | 269 -------------- .../ai-vault-search/session-search-store.ts | 97 +----- ...ession-search-visible-read-ratchet.test.ts | 153 -------- .../session-search-wal-budget.test.ts | 93 ----- .../session-search-wal-budget.ts | 24 -- 22 files changed, 880 insertions(+), 1384 deletions(-) create mode 100644 src/main/ai-vault-search/session-search-file-write.test.ts delete mode 100644 src/main/ai-vault-search/session-search-index-compaction.ts rename src/main/ai-vault-search/{session-search-staged-write-test-fixture.ts => session-search-index-test-fixture.ts} (100%) delete mode 100644 src/main/ai-vault-search/session-search-pending-deletes.ts delete mode 100644 src/main/ai-vault-search/session-search-staged-write.test.ts delete mode 100644 src/main/ai-vault-search/session-search-visible-read-ratchet.test.ts delete mode 100644 src/main/ai-vault-search/session-search-wal-budget.test.ts delete 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 index bff194e77d4..0b1a55e95ae 100644 --- a/src/main/ai-vault-search/session-search-content-hash.test.ts +++ b/src/main/ai-vault-search/session-search-content-hash.test.ts @@ -5,7 +5,7 @@ import { foldContentHash, isCollapsibleContentHash } from './session-search-content-hash' -import { userMessages } from './session-search-staged-write-test-fixture' +import { userMessages } from './session-search-index-test-fixture' it('reaches the same digest whether the prefix arrives whole or in two appends', () => { const messages = userMessages('turn', 5) diff --git a/src/main/ai-vault-search/session-search-file-records.ts b/src/main/ai-vault-search/session-search-file-records.ts index bd563e400aa..950996c7fa6 100644 --- a/src/main/ai-vault-search/session-search-file-records.ts +++ b/src/main/ai-vault-search/session-search-file-records.ts @@ -22,11 +22,19 @@ export function cwdKey(cwd: string | null): string | null { export class SessionSearchFileRecords { constructor(private readonly db: SyncDatabase) {} - createStagingSession(candidate: SessionFileCandidate): number { + /** + * The row a read hangs its messages off, before the parser has said what the + * session is. Never visible on its own: the same transaction that creates it + * either fills it in or, for a chunked read, leaves it holding that read's + * own rows and a cursor no append can continue from. + */ + createSessionRow(candidate: SessionFileCandidate): number { return Number( this.db - .prepare(`INSERT INTO sessions(index_ready,agent,session_id,file_path,title,resume_command) - VALUES (0,?,'',?,'','')`) + .prepare( + `INSERT INTO sessions(agent,session_id,file_path,title,resume_command) + VALUES (?,'',?,'','')` + ) .run(candidate.agent, candidate.file.path).lastInsertRowid ) } @@ -56,9 +64,11 @@ export class SessionSearchFileRecords { contentHash.count ] this.db - .prepare(`UPDATE sessions SET agent = ?, session_id = ?, file_path = ?, codex_home = ?, title = ?, + .prepare( + `UPDATE sessions SET agent = ?, session_id = ?, file_path = ?, codex_home = ?, title = ?, cwd = ?, cwd_key = ?, branch = ?, created_at = ?, updated_at = ?, message_count = ?, resume_command = ?, - content_hash = ?, content_hash_count = ? WHERE id = ?`) + content_hash = ?, content_hash_count = ? WHERE id = ?` + ) .run(...values, rowId) } diff --git a/src/main/ai-vault-search/session-search-file-write.test.ts b/src/main/ai-vault-search/session-search-file-write.test.ts new file mode 100644 index 00000000000..02402f065cc --- /dev/null +++ b/src/main/ai-vault-search/session-search-file-write.test.ts @@ -0,0 +1,328 @@ +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers' +import SyncDatabase from '../sqlite/sync-database' +import { registerSessionSearchIndexConsumer } from './session-search-index-consumer' +import { SessionSearchIndexWriter } from './session-search-index-writer' +import { + openSessionSearchIndexFile, + replayTranscriptRead, + syntheticCandidate, + syntheticSession, + SYNTHETIC_TRANSCRIPT, + userMessages, + type SessionSearchIndexFile +} from './session-search-index-test-fixture' +import { SessionSearchStore } from './session-search-store' + +// Every assertion here reads through `index.db`, a second connection to the same +// file. That is the whole consistency model: one transaction per file in WAL +// mode, so another handle sees the last committed state and never a session part +// way through being rewritten. + +let index: SessionSearchIndexFile +let store: SessionSearchStore +let errors: unknown[] + +beforeEach(async () => { + index = await openSessionSearchIndexFile('ss-file-write') + errors = [] + store = new SessionSearchStore(index.path, (error) => errors.push(error)) + registerSessionSearchIndexConsumer(store) +}) + +afterEach(async () => { + vi.restoreAllMocks() + resetTranscriptConsumersForTests() + store.close() + await index.close() +}) + +function matches(db: SyncDatabase, table: string, term: string): number { + return ( + db + .prepare( + `SELECT count(*) AS n FROM ${table} JOIN messages m ON m.id = ${table}.rowid + JOIN sessions s ON s.id = m.session_row_id WHERE ${table} MATCH ?` + ) + .get(term) as { n: number } + ).n +} + +/** Fails the nth statement matching `pick`, wherever the writer prepares it. */ +function failOnStatement(pick: (sql: string) => boolean, nth: number): void { + const prepare = SyncDatabase.prototype.prepare + let seen = 0 + vi.spyOn(SyncDatabase.prototype, 'prepare').mockImplementation(function ( + this: SyncDatabase, + sql: string + ) { + if (pick(sql) && ++seen === nth) { + throw new Error('index write crashed mid transaction') + } + return prepare.call(this, sql) + }) +} + +function counts(db: SyncDatabase): Record { + const one = (sql: string): number => (db.prepare(sql).get() as { n: number }).n + return { + sessions: one('SELECT count(*) AS n FROM sessions'), + messages: one('SELECT count(*) AS n FROM messages'), + files: one('SELECT count(*) AS n FROM files'), + full: one('SELECT count(*) AS n FROM messages_fts'), + conversation: one('SELECT count(*) AS n FROM conversation_fts') + } +} + +it('writes a whole read in one transaction', () => { + replayTranscriptRead({ messages: userMessages('needle text', 300) }) + + const after = counts(index.db) + expect(after.sessions).toBe(1) + expect(after.messages).toBe(300) + expect(after.full).toBe(300) + expect(errors).toEqual([]) +}) + +it('writes both FTS tables for every conversational row', () => { + replayTranscriptRead({ + messages: [ + { role: 'user', text: 'alpha question', timestamp: null }, + { role: 'assistant', text: 'beta answer', timestamp: null }, + { role: 'tool', text: 'gamma tool output', timestamp: null } + ] + }) + + // messages_fts carries every row; conversation_fts is the tool-free half. + expect(counts(index.db).full).toBe(3) + expect(counts(index.db).conversation).toBe(2) + expect(matches(index.db, 'messages_fts', 'gamma')).toBe(1) + expect(matches(index.db, 'conversation_fts', 'gamma')).toBe(0) + expect(matches(index.db, 'conversation_fts', 'beta')).toBe(1) +}) + +it('leaves the index exactly as it found it when a read never finishes', () => { + const write = store.beginWrite(syntheticCandidate(), 'replace', 0)! + for (const message of userMessages('neverfinished', 200)) { + write.add(message) + } + // The process dies here: the rows only ever existed in this buffer. + expect(counts(index.db)).toMatchObject({ + sessions: 0, + messages: 0, + files: 0 + }) +}) + +it('rolls a whole file back when a write throws part way through its transaction', () => { + replayTranscriptRead({ + messages: userMessages('firstgeneration', 3), + outcome: { byteOffset: 40 } + }) + const before = counts(index.db) + + failOnStatement((sql) => sql.startsWith('INSERT INTO messages('), 50) + replayTranscriptRead({ + messages: userMessages('crashedgeneration', 100), + outcome: { byteOffset: 900 } + }) + vi.restoreAllMocks() + + // Not one of the 49 rows that were already inserted survived, the previous + // generation is untouched, and the cursor still describes what is really here. + expect(counts(index.db)).toEqual(before) + expect(matches(index.db, 'messages_fts', 'crashedgeneration')).toBe(0) + expect(matches(index.db, 'messages_fts', 'firstgeneration')).toBe(3) + expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(40) + expect(errors).toHaveLength(1) + // The file is owed a re-read, which is the only reason anything was lost. + expect(store.takeStale().map((candidate) => candidate.file.path)).toEqual([SYNTHETIC_TRANSCRIPT]) + + // And the connection is usable again: a transaction left open by the failure + // would take down every write after it, not just the one that threw. + replayTranscriptRead({ + messages: userMessages('afterthecrash', 2), + outcome: { byteOffset: 900 } + }) + expect(matches(index.db, 'messages_fts', 'afterthecrash')).toBe(2) + expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(900) +}) + +it('takes the rows back when recording the cursor is what fails', () => { + replayTranscriptRead({ + messages: userMessages('firstgeneration', 3), + outcome: { byteOffset: 40 } + }) + + // The cursor is written last, so this is the crash point that would leave rows + // no cursor describes: a later append would continue from an offset those rows + // already cover, and index the same span twice. + failOnStatement((sql) => sql.startsWith('INSERT INTO files('), 1) + replayTranscriptRead({ + messages: userMessages('crashedgeneration', 5), + outcome: { byteOffset: 900 } + }) + vi.restoreAllMocks() + + expect(counts(index.db)).toMatchObject({ sessions: 1, messages: 3 }) + expect(matches(index.db, 'messages_fts', 'firstgeneration')).toBe(3) + expect(matches(index.db, 'messages_fts', 'crashedgeneration')).toBe(0) + expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(40) +}) + +it('shows a reader on another handle one generation or the other, never a mixture', () => { + replayTranscriptRead({ + messages: userMessages('firstgeneration', 3), + outcome: { byteOffset: 40 } + }) + expect(counts(index.db).messages).toBe(3) + + const write = store.beginWrite(syntheticCandidate(), 'replace', 0)! + for (const message of userMessages('secondgeneration', 7)) { + write.add(message) + // Every point at which the other handle could issue a query mid-read. + expect(counts(index.db).messages).toBe(3) + expect(matches(index.db, 'messages_fts', 'secondgeneration')).toBe(0) + } + expect( + write.commit({ + session: syntheticSession(), + byteOffset: 900, + incomplete: false + }) + ).toBe(true) + + expect(counts(index.db).messages).toBe(7) + expect(matches(index.db, 'messages_fts', 'firstgeneration')).toBe(0) + expect(matches(index.db, 'messages_fts', 'secondgeneration')).toBe(7) +}) + +// Four of these fill the 400-char ceiling the two tests below construct. +const CHUNKED_MESSAGE = `chunkedneedle ${'filler '.repeat(12)}nd` + +it('leaves the session consistent after every chunk of a file too large for one transaction', () => { + expect(CHUNKED_MESSAGE.length).toBe(100) + const writer = new SessionSearchIndexWriter(index.db, 400) + const write = writer.beginWrite(syntheticCandidate(), 'replace', 0)! + for (const [position, message] of userMessages(CHUNKED_MESSAGE, 10).entries()) { + write.add(message) + const rows = counts(index.db).messages + // Four messages per chunk, and nothing else reaches the file between them. + expect(rows).toBe(Math.floor((position + 1) / 4) * 4) + // Whatever landed is a coherent prefix of this session and answers searches. + expect(matches(index.db, 'messages_fts', 'chunkedneedle')).toBe(rows) + if (rows > 0) { + // The cursor a chunk leaves refuses every append rather than inventing an + // offset the reader never gave it. + expect(writer.indexedFile(SYNTHETIC_TRANSCRIPT, null)).toBeNull() + expect(writer.beginWrite(syntheticCandidate(), 'append', 0)).toBeNull() + } + } + expect(counts(index.db).messages).toBe(8) + + expect( + write.commit({ + session: syntheticSession(), + byteOffset: 4096, + incomplete: false + }) + ).toBe(true) + expect(counts(index.db)).toMatchObject({ + sessions: 1, + messages: 10, + full: 10 + }) + expect(writer.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(4096) +}) + +it('re-reads a chunked file whole when its writer died between chunks', () => { + const writer = new SessionSearchIndexWriter(index.db, 400) + const abandoned = writer.beginWrite(syntheticCandidate(), 'replace', 0)! + for (const message of userMessages(CHUNKED_MESSAGE, 10)) { + abandoned.add(message) + } + expect(counts(index.db).messages).toBe(8) + + // Nothing can continue that prefix, so the only way forward is a whole re-read, + // and that replaces every row the dead writer left. + expect(writer.indexedFile(SYNTHETIC_TRANSCRIPT, null)).toBeNull() + const replacement = writer.beginWrite(syntheticCandidate(), 'replace', 0)! + replacement.add(userMessages('wholereread', 1)[0]!) + expect( + replacement.commit({ + session: syntheticSession(), + byteOffset: 4096, + incomplete: false + }) + ).toBe(true) + expect(counts(index.db)).toMatchObject({ sessions: 1, messages: 1 }) + expect(matches(index.db, 'messages_fts', 'chunkedneedle')).toBe(0) +}) + +it('replaces the previous generation without ever showing both', () => { + replayTranscriptRead({ messages: userMessages('firstgeneration', 10) }) + replayTranscriptRead({ messages: userMessages('secondgeneration', 10) }) + + expect(counts(index.db)).toMatchObject({ + sessions: 1, + messages: 10, + full: 10 + }) + expect(matches(index.db, 'messages_fts', 'firstgeneration')).toBe(0) + expect(matches(index.db, 'messages_fts', 'secondgeneration')).toBe(10) +}) + +it('continues a session across an append rather than replaying it', () => { + replayTranscriptRead({ + messages: userMessages('openingturn', 3), + outcome: { byteOffset: 40 } + }) + replayTranscriptRead({ + messages: userMessages('laterturn', 2), + mode: 'append', + previousByteOffset: 40, + outcome: { byteOffset: 90 } + }) + + expect(counts(index.db)).toMatchObject({ sessions: 1, messages: 5 }) + expect(matches(index.db, 'messages_fts', 'openingturn')).toBe(3) + expect(matches(index.db, 'messages_fts', 'laterturn')).toBe(2) + expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(90) +}) + +it('stops answering for a removed file the moment it is removed', () => { + replayTranscriptRead({ messages: userMessages('removedneedle', 3) }) + store.removeFile(SYNTHETIC_TRANSCRIPT) + + expect(counts(index.db)).toMatchObject({ + sessions: 0, + messages: 0, + files: 0, + full: 0 + }) + expect(matches(index.db, 'messages_fts', 'removedneedle')).toBe(0) +}) + +it('writes nothing for an incomplete read and owes the file a whole re-read', () => { + replayTranscriptRead({ + messages: userMessages('incompleteread', 300), + outcome: { incomplete: true } + }) + + expect(counts(index.db)).toMatchObject({ + sessions: 0, + messages: 0, + files: 0, + full: 0 + }) + expect(store.pendingFileCount).toBe(1) + expect(errors).toEqual([]) +}) + +it('closes twice without turning the second call into an error', () => { + store.close() + // node:sqlite throws ERR_INVALID_STATE on a second close of one handle, and a + // store is closed both by whoever owns it and by a teardown that cannot know. + expect(() => store.close()).not.toThrow() + store = new SessionSearchStore(index.path, (error) => errors.push(error)) +}) diff --git a/src/main/ai-vault-search/session-search-index-compaction.ts b/src/main/ai-vault-search/session-search-index-compaction.ts deleted file mode 100644 index 93d2b2649d0..00000000000 --- a/src/main/ai-vault-search/session-search-index-compaction.ts +++ /dev/null @@ -1,22 +0,0 @@ -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-index-consumer.test.ts b/src/main/ai-vault-search/session-search-index-consumer.test.ts index 3fe71d89a13..e5d42673f35 100644 --- a/src/main/ai-vault-search/session-search-index-consumer.test.ts +++ b/src/main/ai-vault-search/session-search-index-consumer.test.ts @@ -9,7 +9,7 @@ import { SYNTHETIC_TRANSCRIPT, userMessages, type SessionSearchIndexFile -} from './session-search-staged-write-test-fixture' +} from './session-search-index-test-fixture' import { SessionSearchStore, STALE_PATH_LIMIT } from './session-search-store' let index: SessionSearchIndexFile @@ -29,8 +29,12 @@ afterEach(async () => { await index.close() }) -function visibleMessages(): number { - return (index.db.prepare('SELECT count(*) AS n FROM visible_messages').get() as { n: number }).n +function indexedMessages(): number { + return ( + index.db.prepare('SELECT count(*) AS n FROM messages').get() as { + n: number + } + ).n } function cursor(): number | undefined { @@ -38,10 +42,12 @@ function cursor(): number | undefined { } it('appends onto its own cursor and carries the content hash forward', async () => { - replayTranscriptRead({ messages: userMessages('first half', 3), outcome: { byteOffset: 100 } }) - await store.settled() + replayTranscriptRead({ + messages: userMessages('first half', 3), + outcome: { byteOffset: 100 } + }) const first = index.db - .prepare('SELECT content_hash AS hash, content_hash_count AS count FROM visible_sessions') + .prepare('SELECT content_hash AS hash, content_hash_count AS count FROM sessions') .get() as { hash: string; count: number } replayTranscriptRead({ @@ -50,12 +56,11 @@ it('appends onto its own cursor and carries the content hash forward', async () messages: userMessages('second half', 2), outcome: { byteOffset: 220 } }) - await store.settled() - expect(visibleMessages()).toBe(5) + expect(indexedMessages()).toBe(5) expect(cursor()).toBe(220) const second = index.db - .prepare('SELECT content_hash AS hash, content_hash_count AS count FROM visible_sessions') + .prepare('SELECT content_hash AS hash, content_hash_count AS count FROM sessions') .get() as { hash: string; count: number } expect(second.count).toBe(first.count + 2) expect(second.hash).not.toBe(first.hash) @@ -70,7 +75,6 @@ it('appends onto a file it read through and decoded no session from', async () = messages: userMessages('excluded span', 3), outcome: { session: null, byteOffset: 100 } }) - await store.settled() expect(cursor()).toBe(100) expect(store.takeStale()).toEqual([]) @@ -80,16 +84,17 @@ it('appends onto a file it read through and decoded no session from', async () = messages: userMessages('decoded at last', 2), outcome: { byteOffset: 220 } }) - await store.settled() - expect(visibleMessages()).toBe(2) + expect(indexedMessages()).toBe(2) expect(cursor()).toBe(220) expect(store.takeStale()).toEqual([]) }) it('declines an append that starts past its own cursor and records the file', async () => { - replayTranscriptRead({ messages: userMessages('indexed span', 3), outcome: { byteOffset: 100 } }) - await store.settled() + replayTranscriptRead({ + messages: userMessages('indexed span', 3), + outcome: { byteOffset: 100 } + }) // The session list read further than this index did, so the appended span // continues from bytes the index never saw. @@ -99,9 +104,8 @@ it('declines an append that starts past its own cursor and records the file', as messages: userMessages('unseen span', 4), outcome: { byteOffset: 1200 } }) - await store.settled() - expect(visibleMessages()).toBe(3) + expect(indexedMessages()).toBe(3) expect(cursor()).toBe(100) expect(store.takeStale().map((candidate) => candidate.file.path)).toEqual([SYNTHETIC_TRANSCRIPT]) }) @@ -113,7 +117,6 @@ it('declines a file whose identity changed under the same path', async () => { messages: userMessages('original file', 2), outcome: { byteOffset: 100 } }) - await store.settled() replayTranscriptRead({ candidate: syntheticCandidate({ dev: 1, ino: 77 }), @@ -122,15 +125,16 @@ it('declines a file whose identity changed under the same path', async () => { messages: userMessages('replacement file', 2), outcome: { byteOffset: 200 } }) - await store.settled() - expect(visibleMessages()).toBe(2) + expect(indexedMessages()).toBe(2) expect(store.takeStale()).toHaveLength(1) }) it('never advances the cursor for an incomplete read', async () => { - replayTranscriptRead({ messages: userMessages('complete span', 3), outcome: { byteOffset: 100 } }) - await store.settled() + replayTranscriptRead({ + messages: userMessages('complete span', 3), + outcome: { byteOffset: 100 } + }) replayTranscriptRead({ mode: 'append', @@ -138,12 +142,16 @@ it('never advances the cursor for an incomplete read', async () => { messages: userMessages('partial span', 5), outcome: { byteOffset: 400, incomplete: true } }) - await store.settled() - await store.purgeOlderThan(null) - expect(visibleMessages()).toBe(3) + expect(indexedMessages()).toBe(3) expect(cursor()).toBe(100) - expect((index.db.prepare('SELECT count(*) AS n FROM messages').get() as { n: number }).n).toBe(3) + expect( + ( + index.db.prepare('SELECT count(*) AS n FROM messages').get() as { + n: number + } + ).n + ).toBe(3) expect(store.takeStale()).toHaveLength(1) }) @@ -152,34 +160,40 @@ it('indexes nothing at all from a read that was incomplete from the start', asyn messages: userMessages('unreachable', 4), outcome: { byteOffset: 0, incomplete: true } }) - await store.settled() - await store.purgeOlderThan(null) - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) - expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 0 + }) + expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ + n: 0 + }) expect(cursor()).toBeUndefined() }) it('drops a file whose parser returned no session', async () => { - replayTranscriptRead({ messages: userMessages('was indexed', 3), outcome: { byteOffset: 100 } }) - await store.settled() + replayTranscriptRead({ + messages: userMessages('was indexed', 3), + outcome: { byteOffset: 100 } + }) replayTranscriptRead({ messages: userMessages('now rejected', 2), outcome: { session: null, byteOffset: 300 } }) - await store.settled() - await store.purgeOlderThan(null) - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) - expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 0 + }) + expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ + n: 0 + }) // The file is still read through, so a later scan does not re-read it. expect(cursor()).toBe(300) }) -it('stages nothing for a source whose parser cannot reach the channel', async () => { +it('writes nothing for a source whose parser cannot reach the channel', async () => { // An OpenCode SQLite candidate decodes in a worker, so every read of it is - // incomplete; opening a batch per scan would tombstone rows forever. + // incomplete, and no re-read would help. const candidate = { ...syntheticCandidate({ path: '/opencode/opencode.db#session-1' }), agent: 'opencode' as const @@ -189,10 +203,8 @@ it('stages nothing for a source whose parser cannot reach the channel', async () messages: [], outcome: { byteOffset: 0, incomplete: true } }) - await store.settled() - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) - expect(index.db.prepare('SELECT count(*) AS n FROM search_pending_deletes').get()).toEqual({ + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) expect(store.takeStale()).toEqual([]) @@ -201,18 +213,20 @@ it('stages nothing for a source whose parser cannot reach the channel', async () it('ignores a candidate older than the retention cutoff', async () => { store.setRetentionCutoffMs(Date.now()) replayTranscriptRead({ messages: userMessages('too old', 3) }) - await store.settled() - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 0 + }) expect(store.takeStale()).toEqual([]) }) it('stops writing while the store refuses writes, but remembers what it skipped', async () => { store.setAcceptingWrites(false) replayTranscriptRead({ messages: userMessages('paused', 3) }) - await store.settled() - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 0 + }) expect(errors).toEqual([]) // A pause is exactly the window in which every read is declined. Forgetting // them would leave the whole paused span unindexed with nothing to replay it. @@ -222,7 +236,6 @@ it('stops writing while the store refuses writes, but remembers what it skipped' it('keeps the paused re-read set when the retention window is reconfigured', async () => { store.setAcceptingWrites(false) replayTranscriptRead({ messages: userMessages('paused', 2) }) - await store.settled() expect(store.pendingFileCount).toBe(1) // The set records what still has to be read, not what is worth keeping. A @@ -234,8 +247,9 @@ it('keeps the paused re-read set when the retention window is reconfigured', asy store.setAcceptingWrites(true) expect(store.takeStale()).toHaveLength(1) replayTranscriptRead({ messages: userMessages('outside the window now', 2) }) - await store.settled() - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 0 + }) expect(store.pendingFileCount).toBe(0) }) @@ -253,8 +267,10 @@ it('drops the oldest record rather than growing without a bound, and says so', ( }) it('keeps the session list running when the index write fails', async () => { - replayTranscriptRead({ messages: userMessages('healthy', 2), outcome: { byteOffset: 100 } }) - await store.settled() + replayTranscriptRead({ + messages: userMessages('healthy', 2), + outcome: { byteOffset: 100 } + }) index.db.exec('DROP TABLE messages_fts') expect(() => @@ -265,30 +281,36 @@ it('keeps the session list running when the index write fails', async () => { outcome: { byteOffset: 500 } }) ).not.toThrow() - expect(store.failures).toBeGreaterThan(0) + expect(errors.length).toBeGreaterThan(0) expect(store.takeStale()).toHaveLength(1) }) it('unregisters cleanly, leaving later reads unindexed', async () => { resetTranscriptConsumersForTests() replayTranscriptRead({ messages: userMessages('after unregister', 3) }) - await store.settled() - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 0 + }) }) it('drops a removed source and keeps its cursor gone', async () => { - replayTranscriptRead({ messages: userMessages('present', 3), outcome: { byteOffset: 100 } }) - await store.settled() + replayTranscriptRead({ + messages: userMessages('present', 3), + outcome: { byteOffset: 100 } + }) store.removeFile(SYNTHETIC_TRANSCRIPT) - await store.purgeOlderThan(null) expect(cursor()).toBeUndefined() - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) - expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ n: 0 }) + expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 0 + }) + expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ + n: 0 + }) }) -it('publishes the session metadata the read decoded', async () => { +it('writes the session metadata the read decoded', async () => { replayTranscriptRead({ messages: userMessages('metadata', 1), outcome: { @@ -303,13 +325,10 @@ it('publishes the session metadata the read decoded', async () => { byteOffset: 42 } }) - await store.settled() expect( index.db - .prepare( - 'SELECT session_id, title, cwd, cwd_key, branch, resume_command FROM visible_sessions' - ) + .prepare('SELECT session_id, title, cwd, cwd_key, branch, resume_command FROM sessions') .get() ).toEqual({ session_id: 'abc-123', @@ -328,7 +347,6 @@ it('keeps a proven file identity when a later read cannot stat it', async () => messages: userMessages('first', 2), outcome: { byteOffset: 100 } }) - await store.settled() // A host that cannot prove identity re-reads the same file. replayTranscriptRead({ @@ -338,8 +356,7 @@ it('keeps a proven file identity when a later read cannot stat it', async () => messages: userMessages('second', 2), outcome: { byteOffset: 200 } }) - await store.settled() - expect(visibleMessages()).toBe(4) + expect(indexedMessages()).toBe(4) // The stored identity survived, so a rename-replace is still detectable. replayTranscriptRead({ @@ -349,29 +366,8 @@ it('keeps a proven file identity when a later read cannot stat it', async () => messages: userMessages('replacement', 2), outcome: { byteOffset: 300 } }) - await store.settled() - expect(visibleMessages()).toBe(4) + expect(indexedMessages()).toBe(4) expect(cursor()).toBe(200) expect(store.takeStale()).toHaveLength(1) }) - -it('leaves no batch on disk when a rejected read is the last one before shutdown', async () => { - replayTranscriptRead({ messages: userMessages('indexed', 3), outcome: { byteOffset: 100 } }) - await store.settled() - - // The parser rejects the file, so this read publishes a cursor and tombstones - // its own staged rows. Nothing writes after it. - replayTranscriptRead({ - messages: userMessages('rejected', 4), - outcome: { session: null, byteOffset: 300 } - }) - await store.settled() - - expect(index.db.prepare('SELECT count(*) AS n FROM search_write_batches').get()).toEqual({ n: 0 }) - expect(index.db.prepare('SELECT count(*) AS n FROM search_pending_deletes').get()).toEqual({ - n: 0 - }) - expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ n: 0 }) - expect(index.db.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ n: 0 }) -}) diff --git a/src/main/ai-vault-search/session-search-index-consumer.ts b/src/main/ai-vault-search/session-search-index-consumer.ts index f2d10bb2f87..f74558fb7bb 100644 --- a/src/main/ai-vault-search/session-search-index-consumer.ts +++ b/src/main/ai-vault-search/session-search-index-consumer.ts @@ -8,7 +8,7 @@ import { type TranscriptReadStart } from '../ai-vault/session-transcript-consumers' import { fileIdentity } from './session-search-file-cursor' -import type { SessionSearchStagedWrite } from './session-search-index-writer' +import type { SessionSearchFileWrite } from './session-search-index-writer' import type { SessionSearchStore } from './session-search-store' /** @@ -21,8 +21,8 @@ import type { SessionSearchStore } from './session-search-store' * * - `beginRead` returns null when this index's cursor is behind the offset an * `append` continues from, or when the file's identity changed. - * - a staging failure stops the read's rows without failing the session list. - * - an `incomplete` outcome never publishes; those rows are not the whole span. + * - a buffering failure stops the read's rows without failing the session list. + * - an `incomplete` outcome never commits; those rows are not the whole span. */ export class SessionSearchIndexConsumer implements TranscriptConsumer { constructor(private readonly store: SessionSearchStore) {} @@ -51,12 +51,12 @@ export class SessionSearchIndexConsumer implements TranscriptConsumer { return null } } - const staged = this.store.beginWrite(candidate, start.mode, start.previousByteOffset) - if (!staged) { + const write = this.store.beginWrite(candidate, start.mode, start.previousByteOffset) + if (!write) { this.store.markStale(candidate) return null } - return new SessionSearchReadConsumer(this.store, start, staged) + return new SessionSearchReadConsumer(this.store, start, write) } } @@ -66,7 +66,7 @@ class SessionSearchReadConsumer implements TranscriptReadConsumer { constructor( private readonly store: SessionSearchStore, private readonly start: TranscriptReadStart, - private readonly staged: SessionSearchStagedWrite + private readonly write: SessionSearchFileWrite ) {} message(message: TranscriptMessage): void { @@ -74,11 +74,12 @@ class SessionSearchReadConsumer implements TranscriptReadConsumer { return } try { - this.staged.add(message) + this.write.add(message) } catch (error) { // Never throws back into the reader: the channel would drop this consumer - // for the rest of the read and `finish` would never run, stranding the - // staged batch. Failing here keeps the cleanup on one path. + // for the rest of the read and `finish` would never run. Failing here + // keeps the whole read on one path — the buffer is dropped and the file is + // re-read. this.failed = true this.store.reportWriteFailure(error) } @@ -86,27 +87,19 @@ class SessionSearchReadConsumer implements TranscriptReadConsumer { finish(outcome: TranscriptReadOutcome): void { const { candidate } = this.start - let published = false + let committed = false try { // An incomplete read's rows are not the whole span, so the cursor must not // move past them; the file is re-read whole instead. - published = !this.failed && !outcome.incomplete && this.staged.publish(outcome) + committed = !this.failed && !outcome.incomplete && this.write.commit(outcome) } catch (error) { this.store.reportWriteFailure(error) - } finally { - // Before the store is told anything. `discard` is what tombstones the - // staged rows of a read that decoded no session, and the store's cleanup - // lane reads the tombstone table the moment it is scheduled — telling the - // store first left that batch on disk until some later write happened to - // schedule another pass, which for the last read before a shutdown is - // never. - this.staged.discard() } - if (published) { - this.store.writePublished(candidate) + if (committed) { + this.store.writeCommitted(candidate) return } - this.store.writeAbandoned(candidate) + this.store.markStale(candidate) } } diff --git a/src/main/ai-vault-search/session-search-staged-write-test-fixture.ts b/src/main/ai-vault-search/session-search-index-test-fixture.ts similarity index 100% rename from src/main/ai-vault-search/session-search-staged-write-test-fixture.ts rename to src/main/ai-vault-search/session-search-index-test-fixture.ts diff --git a/src/main/ai-vault-search/session-search-index-writer.test.ts b/src/main/ai-vault-search/session-search-index-writer.test.ts index 6c081040495..f257ec7b4ce 100644 --- a/src/main/ai-vault-search/session-search-index-writer.test.ts +++ b/src/main/ai-vault-search/session-search-index-writer.test.ts @@ -1,6 +1,5 @@ import { afterEach, beforeEach, expect, it } from 'vitest' import { SessionSearchIndexConsumer } from './session-search-index-consumer' -import { SessionSearchIndexWriter } from './session-search-index-writer' import { openSessionSearchIndexFile, syntheticCandidate, @@ -8,7 +7,7 @@ import { SYNTHETIC_TRANSCRIPT, userMessages, type SessionSearchIndexFile -} from './session-search-staged-write-test-fixture' +} from './session-search-index-test-fixture' import { SessionSearchStore } from './session-search-store' // The store is driven directly here. Every guard below is also shadowed by the @@ -31,75 +30,83 @@ afterEach(async () => { }) function count(table: string): number { - return (index.db.prepare(`SELECT count(*) AS n FROM ${table}`).get() as { n: number }).n + return ( + index.db.prepare(`SELECT count(*) AS n FROM ${table}`).get() as { + n: number + } + ).n } -function publishRead(previousByteOffset: number, byteOffset: number, text: string): boolean { - const staged = store.beginWrite( +function indexRead(previousByteOffset: number, byteOffset: number, text: string): boolean { + const write = store.beginWrite( syntheticCandidate(), previousByteOffset === 0 ? 'replace' : 'append', previousByteOffset ) - if (!staged) { + if (!write) { return false } for (const message of userMessages(text, 2)) { - staged.add(message) + write.add(message) } - const published = staged.publish({ session: syntheticSession(), byteOffset, incomplete: false }) - staged.discard() - return published + return write.commit({ + session: syntheticSession(), + byteOffset, + incomplete: false + }) } -it('refuses an append whose predecessor offset is not the published cursor', () => { - expect(publishRead(0, 100, 'first')).toBe(true) +it('refuses an append whose predecessor offset is not the committed cursor', () => { + expect(indexRead(0, 100, 'first')).toBe(true) expect(store.beginWrite(syntheticCandidate(), 'append', 900)).toBeNull() expect(store.beginWrite(syntheticCandidate(), 'append', 99)).toBeNull() - // The one offset that does continue the published span is accepted. + // The one offset that does continue the committed span is accepted. expect(store.beginWrite(syntheticCandidate(), 'append', 100)).not.toBeNull() }) -it('refuses to publish a stage whose cursor moved underneath it', async () => { - const candidate = syntheticCandidate() - const stale = store.beginWrite(candidate, 'replace', 0)! +it('refuses to commit a write whose cursor moved underneath it', () => { + const stale = store.beginWrite(syntheticCandidate(), 'replace', 0)! for (const message of userMessages('stalegeneration', 40)) { stale.add(message) } // A second read of the same path finishes first. Without the parse file lane // this is the overlap that would otherwise resurrect the stale rows. - expect(publishRead(0, 200, 'winninggeneration')).toBe(true) - - expect(stale.publish({ session: syntheticSession(), byteOffset: 100, incomplete: false })).toBe( - false - ) - stale.discard() - await store.purgeOlderThan(null) + expect(indexRead(0, 200, 'winninggeneration')).toBe(true) + expect( + stale.commit({ + session: syntheticSession(), + byteOffset: 100, + incomplete: false + }) + ).toBe(false) expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(200) - expect(count('visible_sessions')).toBe(1) + expect(count('sessions')).toBe(1) expect(count('messages')).toBe(2) - expect(count('search_pending_deletes')).toBe(0) expect(errors).toEqual([]) }) -it('refuses to publish a stage whose file was removed mid-read', async () => { - expect(publishRead(0, 100, 'firstgeneration')).toBe(true) - const staged = store.beginWrite(syntheticCandidate(), 'append', 100)! +it('refuses to commit a write whose file was removed mid-read', () => { + expect(indexRead(0, 100, 'firstgeneration')).toBe(true) + const write = store.beginWrite(syntheticCandidate(), 'append', 100)! for (const message of userMessages('afterremoval', 10)) { - staged.add(message) + write.add(message) } store.removeFile(SYNTHETIC_TRANSCRIPT) - expect(staged.publish({ session: syntheticSession(), byteOffset: 300, incomplete: false })).toBe( - false - ) - staged.discard() - await store.purgeOlderThan(null) - + // Committing here would put a source back that its owner proved was deleted. + expect( + write.commit({ + session: syntheticSession(), + byteOffset: 300, + incomplete: false + }) + ).toBe(false) expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, null)).toBeNull() expect(count('sessions')).toBe(0) expect(count('messages')).toBe(0) + expect(count('files')).toBe(0) }) it('declines a behind cursor in beginRead before it ever reaches the store', () => { @@ -109,34 +116,51 @@ it('declines a behind cursor in beginRead before it ever reaches the store', () indexedFile: () => ({ byteOffset: 100, mtimeMs: 1, sizeBytes: 1 }), beginWrite: (_candidate: unknown, _mode: unknown, previousByteOffset: number) => { attempted.push(previousByteOffset) - return { add: () => undefined, publish: () => true, discard: () => undefined } + return { add: () => undefined, commit: () => true } }, markStale: () => undefined } as unknown as SessionSearchStore const consumer = new SessionSearchIndexConsumer(stub) expect( - consumer.beginRead({ candidate: syntheticCandidate(), mode: 'append', previousByteOffset: 900 }) + consumer.beginRead({ + candidate: syntheticCandidate(), + mode: 'append', + previousByteOffset: 900 + }) ).toBeNull() // The store was never asked, so the writer's own guard cannot be what refused. expect(attempted).toEqual([]) expect( - consumer.beginRead({ candidate: syntheticCandidate(), mode: 'append', previousByteOffset: 100 }) + consumer.beginRead({ + candidate: syntheticCandidate(), + mode: 'append', + previousByteOffset: 100 + }) ).not.toBeNull() expect(attempted).toEqual([100]) }) -it('treats half a recorded identity as no identity at all', async () => { +it('treats half a recorded identity as no identity at all', () => { // A host that could stat dev but not ino: `remote-session-file-stat` spreads // the two independently, and `upsertFile` preserves the half it was given. - const partial = { ...syntheticCandidate({ dev: 7 }), agent: 'claude' as const } - const staged = store.beginWrite(partial, 'replace', 0)! - for (const message of userMessages('halfidentity', 2)) { - staged.add(message) + const partial = { + ...syntheticCandidate({ dev: 7 }), + agent: 'claude' as const } - staged.publish({ session: syntheticSession(), byteOffset: 100, incomplete: false }) - staged.discard() - expect(index.db.prepare('SELECT dev, ino FROM files').get()).toEqual({ dev: 7, ino: null }) + const write = store.beginWrite(partial, 'replace', 0)! + for (const message of userMessages('halfidentity', 2)) { + write.add(message) + } + write.commit({ + session: syntheticSession(), + byteOffset: 100, + incomplete: false + }) + expect(index.db.prepare('SELECT dev, ino FROM files').get()).toEqual({ + dev: 7, + ino: null + }) // One matching number is not proof of sameness, and one mismatching number is // not proof of replacement. Neither compares, so neither declines. @@ -144,19 +168,3 @@ it('treats half a recorded identity as no identity at all', async () => { expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, { dev: 8, ino: 99 })?.byteOffset).toBe(100) expect(store.beginWrite(syntheticCandidate({ dev: 8, ino: 99 }), 'append', 100)).not.toBeNull() }) - -it('forgets a finished read rather than growing a stage per path', () => { - const writer = new SessionSearchIndexWriter(index.db) - for (const path of ['/a.jsonl', '/b.jsonl', '/a.jsonl']) { - const staged = writer.beginWrite(syntheticCandidate({ path }), 'replace', 0)! - staged.add(userMessages('leakcheck', 1)[0]) - staged.publish({ session: syntheticSession(), byteOffset: 10, incomplete: false }) - staged.discard() - expect(writer.openStageCount).toBe(0) - } - // An abandoned read is still a finished one once it is discarded. - const abandoned = writer.beginWrite(syntheticCandidate({ path: '/c.jsonl' }), 'replace', 0)! - expect(writer.openStageCount).toBe(1) - abandoned.discard() - expect(writer.openStageCount).toBe(0) -}) diff --git a/src/main/ai-vault-search/session-search-index-writer.ts b/src/main/ai-vault-search/session-search-index-writer.ts index f43b5377482..8364eabcf5f 100644 --- a/src/main/ai-vault-search/session-search-index-writer.ts +++ b/src/main/ai-vault-search/session-search-index-writer.ts @@ -10,20 +10,36 @@ import type { SessionSearchIndexedFile } from './session-search-file-cursor' import { SessionSearchFileRecords } from './session-search-file-records' -import { insertSearchMessage, searchMessageRows } from './session-search-message-rows' -import { discardSearchBatch, retireSearchSession } from './session-search-pending-deletes' -import { assertSearchWalBudget, SEARCH_WAL_PENDING_BYTES } from './session-search-wal-budget' +import { + deleteSearchMessages, + insertSearchMessage, + searchMessageRows +} from './session-search-message-rows' -// Why bounded rather than streamed: the transcript reader pushes messages -// synchronously, so a staged write cannot make the producer wait. Rows are -// buffered to one of these two ceilings and then written in a single -// transaction, which caps both the retained bytes and the length of one stall. -export const SEARCH_WRITE_ROWS_PER_STEP = 128 -export const SEARCH_WRITE_CHARS_PER_STEP = 256 * 1024 -// Why sampled: the checkpoint costs more than the step it guards, and the backlog only grows -// while a second connection pins a snapshot — a killed scanner child whose handle outlives the -// replacement fork, not two processes the app runs on purpose. -const WAL_BUDGET_EVERY_STEPS = 16 +/** + * How much decoded text one transaction may carry. + * + * A file's rows are buffered in memory and written in one transaction, so the + * whole read is either in the index or not. The ceiling is what keeps that + * promise affordable: at the measured 26 MB of transcript per second it caps a + * single commit near a second and the WAL it produces near 64 MB, and it is far + * above the largest real transcript (the 40-session benchmark corpus is 10.5 MB + * in total), so an ordinary file never reaches it. Above the ceiling the read is + * cut into chunks that each leave the index consistent — see `chunked`. + */ +export const SESSION_SEARCH_COMMIT_CHARS = 32 * 1024 * 1024 + +/** + * The cursor of a file whose rows are a prefix, written by a chunk of a read + * that has not reached the end of the file. + * + * The reader hands out byte offsets only when a read finishes, so a chunk has + * no honest offset to record. This one is unusable on purpose: `indexedFile` + * reports no cursor for it, so an append is declined and the file is re-read + * whole. The rows are still a coherent prefix of that session and answer + * searches until the re-read replaces them. + */ +const PARTIAL_FILE_CURSOR = -1 type FileRow = { dev: number | null @@ -34,51 +50,38 @@ type FileRow = { session_row_id: number | null } -type ExistingFile = Pick +type FileCursor = Pick -export type SessionSearchStagedWrite = { - /** Buffers one message, flushing a full batch into the staging area. */ +export type SessionSearchFileWrite = { + /** Buffers one message, committing a chunk when the buffer reaches the ceiling. */ add(message: TranscriptMessage): void /** - * Makes every staged row visible in one transaction, or drops the file's rows - * when the read decoded no session. False when the file was invalidated or - * its published cursor moved while this write was staging. + * Writes this file's rows, its session and its cursor in one transaction. + * False when the file's record changed under this read — it was removed, or + * another writer moved the cursor these rows continue from. A read that never + * calls this leaves the index exactly as it found it, unless it chunked. */ - publish(outcome: TranscriptReadOutcome): boolean - /** Tombstones whatever is still staged; safe after `publish` and after a failure. */ - discard(): void + commit(outcome: TranscriptReadOutcome): boolean } export class SessionSearchIndexWriter { private readonly records: SessionSearchFileRecords - // One open stage per path. Reads of one transcript are serialized by the parse - // file lane, so this only ever tracks the read in flight; `removeFile` can - // therefore invalidate the open stage by path. If the lane is ever bypassed, - // the later stage takes the slot and the earlier one loses its invalidation - // hook, but it still cannot publish: `publishable` re-reads the cursor and - // refuses. Both halves are pinned in session-search-index-writer.test.ts. - private readonly staging = new Map() constructor( private readonly db: SyncDatabase, - private readonly walBudgetBytes: number = SEARCH_WAL_PENDING_BYTES + private readonly commitChars: number = SESSION_SEARCH_COMMIT_CHARS ) { this.records = new SessionSearchFileRecords(db) } - /** Tests only: an entry surviving a finished read leaks for the store's life. */ - get openStageCount(): number { - return this.staging.size - } - - /** This index's own cursor, or null when the file is unknown or its identity changed. */ + /** This index's own cursor, or null when the file is unknown, changed, or half written. */ indexedFile(path: string, identity: SessionSearchFileIdentity): SessionSearchIndexedFile | null { const row = this.db .prepare( 'SELECT dev, ino, byte_offset, mtime_ms, size_bytes, session_row_id FROM files WHERE path = ?' ) .get(path) as FileRow | undefined - if (!row) { + if (!row || row.byte_offset === PARTIAL_FILE_CURSOR) { return null } // A recorded identity that no longer matches is a different file at the same @@ -95,47 +98,46 @@ export class SessionSearchIndexWriter { return null } } - return { byteOffset: row.byte_offset, mtimeMs: row.mtime_ms, sizeBytes: row.size_bytes } + return { + byteOffset: row.byte_offset, + mtimeMs: row.mtime_ms, + sizeBytes: row.size_bytes + } } /** - * Opens a staging batch for one read, or returns null when the read cannot - * extend what is published: an `append` whose predecessor byte offset is not - * this index's own cursor covers a span the index never saw. + * Opens a buffered write for one read, or returns null when the read cannot + * extend what the index holds: an `append` whose predecessor byte offset is + * not this index's own cursor covers a span the index never saw. */ beginWrite( candidate: SessionFileCandidate, mode: 'replace' | 'append', previousByteOffset: number - ): SessionSearchStagedWrite | null { + ): SessionSearchFileWrite | null { const path = candidate.file.path - const existing = this.file(path) - if (mode === 'append' && existing?.byte_offset !== previousByteOffset) { + const cursor = this.cursor(path) + if (mode === 'append' && cursor?.byte_offset !== previousByteOffset) { return null } // A file the index read through and decoded no session from still has a // cursor worth continuing: it has no session row to hang new rows off, so // this read makes one. Declining instead would force a whole re-read of // that file on every pass for as long as it grows. - return this.stage( - candidate, - existing, - mode === 'append' ? (existing?.session_row_id ?? null) : null - ) + return this.buffered(candidate, cursor, mode === 'append') } - /** Invalidation hides the generation immediately; cleanup does the expensive deletes later. */ + /** + * Drops a source: its session, its rows and its file record, in one + * transaction. Unbounded on purpose — the caller has proven this one file is + * gone and expects it out of results when the call returns, and a read of it + * that is still in flight is fenced by the cursor its commit re-reads. + */ removeFile(path: string): void { - const open = this.staging.get(path) - if (open) { - open.invalidated = true - } - const existing = this.file(path) + const cursor = this.cursor(path) this.db.exec('BEGIN IMMEDIATE') try { - if (existing?.session_row_id != null) { - retireSearchSession(this.db, existing.session_row_id) - } + this.dropSession(cursor?.session_row_id ?? null) this.db.prepare('DELETE FROM files WHERE path = ?').run(path) this.db.exec('COMMIT') } catch (error) { @@ -144,144 +146,112 @@ export class SessionSearchIndexWriter { } } - private file(path: string): ExistingFile | undefined { + private cursor(path: string): FileCursor | undefined { return this.db .prepare('SELECT session_row_id,byte_offset FROM files WHERE path = ?') - .get(path) as ExistingFile | undefined + .get(path) as FileCursor | undefined } - /** `resumed` is the session row this read continues, or null when it starts one. */ - private stage( + private buffered( candidate: SessionFileCandidate, - existing: ExistingFile | undefined, - resumed: number | null - ): SessionSearchStagedWrite { + opened: FileCursor | undefined, + append: boolean + ): SessionSearchFileWrite { const db = this.db const path = candidate.file.path - let hash = resumed === null ? EMPTY_CONTENT_HASH : this.records.contentHash(resumed) - let sessionId: number - let batchId: number - db.exec('BEGIN IMMEDIATE') - try { - sessionId = resumed ?? this.records.createStagingSession(candidate) - batchId = Number( - db.prepare('INSERT INTO search_write_batches(session_row_id) VALUES (?)').run(sessionId) - .lastInsertRowid - ) - db.exec('COMMIT') - } catch (error) { - db.exec('ROLLBACK') - throw error - } - - const open = { invalidated: false } - this.staging.set(path, open) const buffer: TranscriptMessage[] = [] let bufferedChars = 0 - let steps = 0 + // What this write believes the file record holds. Re-read inside every + // transaction: a `removeFile` or another writer between two chunks means + // these rows no longer continue anything, and committing on top of that + // would resurrect a deleted source or duplicate a span. + let expected = opened + // The session row is reused across re-reads of one file, so a `replace` + // swaps a session's rows rather than minting a second generation of it. + let session = opened?.session_row_id ?? null + let hash = append && session !== null ? this.records.contentHash(session) : EMPTY_CONTENT_HASH + // A replace owns the session's whole row set, so the old generation goes in + // the same transaction as the first of the new one. Chunk two onwards must + // not repeat it. + let replaced = append - // A writer that published its own batch for this path moved the cursor these - // staged rows continue from; publishing on top would duplicate or skip a span. - const publishable = (): boolean => { - if (open.invalidated) { - return false - } - const current = this.file(path) + const current = (): boolean => { + const row = this.cursor(path) return ( - current?.session_row_id === existing?.session_row_id && - current?.byte_offset === existing?.byte_offset + row?.session_row_id === expected?.session_row_id && + row?.byte_offset === expected?.byte_offset ) } - const flush = (): void => { - if (buffer.length === 0) { - return - } - if (steps++ % WAL_BUDGET_EVERY_STEPS === 0) { - assertSearchWalBudget(db, this.walBudgetBytes) - } + /** `outcome` is null for a chunk of a read that has not reached the file's end. */ + const write = (outcome: TranscriptReadOutcome | null): boolean => { + const decoded = outcome?.session ?? null db.exec('BEGIN IMMEDIATE') try { - for (const row of buffer) { - insertSearchMessage(db, sessionId, batchId, row) + if (!current()) { + db.exec('ROLLBACK') + return false + } + if (outcome && !decoded) { + // Read through, but nothing to search: the cursor advances so the file + // is not re-read whole on every pass, and whatever generation was here + // — including this read's own committed chunks — goes with it. + this.dropSession(session) + session = null + this.records.upsertFile(candidate, outcome.byteOffset, null) + } else { + session ??= this.records.createSessionRow(candidate) + if (!replaced) { + deleteSearchMessages(db, session) + replaced = true + } + for (const row of buffer) { + insertSearchMessage(db, session, row) + } + if (decoded) { + this.records.updateSession(decoded, session, hash) + } + this.records.upsertFile( + candidate, + outcome ? outcome.byteOffset : PARTIAL_FILE_CURSOR, + session + ) } db.exec('COMMIT') } catch (error) { db.exec('ROLLBACK') throw error } + expected = { + session_row_id: session, + byte_offset: outcome ? outcome.byteOffset : PARTIAL_FILE_CURSOR + } buffer.length = 0 bufferedChars = 0 - } - - const close = (): void => { - if (this.staging.get(path) === open) { - this.staging.delete(path) - } + return true } return { add: (message) => { - if (open.invalidated) { - return - } hash = foldContentHash(hash, [message]) for (const row of searchMessageRows([message])) { buffer.push(row) bufferedChars += row.text.length - if ( - buffer.length >= SEARCH_WRITE_ROWS_PER_STEP || - bufferedChars >= SEARCH_WRITE_CHARS_PER_STEP - ) { - flush() - } + } + if (bufferedChars >= this.commitChars) { + write(null) } }, - publish: (outcome) => { - if (!publishable()) { - return false - } - flush() - db.exec('BEGIN IMMEDIATE') - try { - if (outcome.session) { - this.records.updateSession(outcome.session, sessionId, hash) - if (resumed === null && existing?.session_row_id != null) { - retireSearchSession(db, existing.session_row_id) - } - db.prepare('UPDATE sessions SET index_ready=1 WHERE id=?').run(sessionId) - // Clearing the pointer before dropping the batch is what makes a recycled - // rowid harmless: no published row can name a later in-flight batch. - db.prepare('UPDATE messages SET batch_id=NULL WHERE batch_id=?').run(batchId) - db.prepare('DELETE FROM search_write_batches WHERE id=?').run(batchId) - this.records.upsertFile(candidate, outcome.byteOffset, sessionId) - } else { - // No session: the file is read through but holds nothing to search, - // so the cursor advances and the old generation's rows are retired. - // `discard` then tombstones this write's own staging rows. - if (existing?.session_row_id != null) { - retireSearchSession(db, existing.session_row_id) - } - this.records.upsertFile(candidate, outcome.byteOffset, null) - } - db.exec('COMMIT') - } catch (error) { - db.exec('ROLLBACK') - throw error - } - return true - }, - discard: () => { - buffer.length = 0 - close() - // A surviving batch row means publish never made these rows visible, - // whatever ended the stage. - if (db.prepare('SELECT 1 FROM search_write_batches WHERE id=?').get(batchId)) { - // Owning the session means this read created it, so retiring it takes - // the whole staging generation with it. - discardSearchBatch(db, sessionId, batchId, resumed === null) - } - } + commit: (outcome) => write(outcome) } } + + /** Caller's transaction: drops a session and every row that hangs off it. */ + private dropSession(sessionRowId: number | null): void { + if (sessionRowId === null) { + return + } + deleteSearchMessages(this.db, sessionRowId) + this.db.prepare('DELETE FROM sessions WHERE id = ?').run(sessionRowId) + } } diff --git a/src/main/ai-vault-search/session-search-live-transcript.test.ts b/src/main/ai-vault-search/session-search-live-transcript.test.ts index b8d02835180..fe36dcb84c0 100644 --- a/src/main/ai-vault-search/session-search-live-transcript.test.ts +++ b/src/main/ai-vault-search/session-search-live-transcript.test.ts @@ -47,14 +47,14 @@ async function makeTempDir(): Promise { return root } -/** Sessions a query over the published views would return for one FTS term. */ +/** Sessions a query would return for one FTS term, read on a second handle. */ function sessionsMatching(term: string, table = 'messages_fts'): string[] { return ( reader .prepare( `SELECT DISTINCT s.session_id AS id FROM ${table} - JOIN visible_messages m ON m.id = ${table}.rowid - JOIN visible_sessions s ON s.id = m.session_row_id + JOIN messages m ON m.id = ${table}.rowid + JOIN sessions s ON s.id = m.session_row_id WHERE ${table} MATCH ? ORDER BY s.session_id` ) .all(term) as { id: string }[] @@ -84,7 +84,9 @@ it('indexes a Claude transcript through the reader and resumes on append', async expect(errors).toEqual([]) expect(sessionsMatching('zygomorphic')).toEqual([SESSION_ID]) // An append extends one session rather than creating a second. - expect(reader.prepare('SELECT count(*) AS n FROM visible_sessions').get()).toEqual({ n: 1 }) + expect(reader.prepare('SELECT count(*) AS n FROM sessions').get()).toEqual({ + n: 1 + }) }) it('keeps a tool result searchable but out of the conversation half', async () => { 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 index f86b3f96cf7..15612909a6e 100644 --- a/src/main/ai-vault-search/session-search-message-rows.test.ts +++ b/src/main/ai-vault-search/session-search-message-rows.test.ts @@ -1,23 +1,18 @@ import { expect, it } from 'vitest' import type { TranscriptMessage } from '../ai-vault/session-transcript-consumers' -import { - chunkMessageText, - insertSearchMessage, - searchMessageRows -} from './session-search-message-rows' +import { insertSearchMessage, searchMessageRows } from './session-search-message-rows' import { openSessionSearchIndexFile, type SessionSearchIndexFile -} from './session-search-staged-write-test-fixture' +} from './session-search-index-test-fixture' /** Every column of both FTS tables, so an assertion cannot miss the shadow terms. */ async function indexedColumns( index: SessionSearchIndexFile, message: TranscriptMessage ): Promise { - index.db.exec('INSERT INTO search_write_batches(id,session_row_id) VALUES (1,1)') for (const row of searchMessageRows([message])) { - insertSearchMessage(index.db, 1, 1, row) + insertSearchMessage(index.db, 1, row) } const full = index.db .prepare('SELECT user_text, assistant_text, tool_text, identifiers FROM messages_fts') @@ -31,7 +26,9 @@ async function indexedColumns( 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) + const chunks = [...searchMessageRows([{ role: 'user', text, timestamp: null }])].map( + (row) => row.text + ) expect(chunks.length).toBeGreaterThan(1) expect(chunks.join('')).toBe(text) @@ -42,17 +39,17 @@ it('splits an oversized message on a line boundary and keeps every character', ( }) it('leaves a message that fits as a single row', () => { - expect(chunkMessageText('short enough')).toEqual(['short enough']) + const rows = [...searchMessageRows([{ role: 'user', text: 'short enough', timestamp: null }])] + expect(rows.map((row) => row.text)).toEqual(['short enough']) }) 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)') for (const row of searchMessageRows([ { role: 'tool', text: 'rg pericardium', timestamp: null } ])) { - insertSearchMessage(index.db, 1, 1, row) + insertSearchMessage(index.db, 1, row) } 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 }) @@ -65,7 +62,11 @@ it('stores a chunk exactly as the transcript wrote it', async () => { const index = await openSessionSearchIndexFile('ss-rows-verbatim') try { const text = 'deploy with AKIAIOSFODNN7EXAMPLE and the resolveTerminalPath fix' - const stored = await indexedColumns(index, { role: 'assistant', text, timestamp: null }) + const stored = await indexedColumns(index, { + role: 'assistant', + text, + timestamp: null + }) // The index is a second copy of content the user already holds in plaintext, // so it neither rewrites nor drops any of it. diff --git a/src/main/ai-vault-search/session-search-message-rows.ts b/src/main/ai-vault-search/session-search-message-rows.ts index da5e3e4b7ab..81acda996d0 100644 --- a/src/main/ai-vault-search/session-search-message-rows.ts +++ b/src/main/ai-vault-search/session-search-message-rows.ts @@ -36,20 +36,19 @@ export function* searchMessageRows( /** * 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 + * transaction, so a 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 = 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 + .prepare('INSERT INTO messages(session_row_id, role, ts) VALUES (?, ?, ?)') + .run(sessionId, message.role, message.timestamp).lastInsertRowid const user = message.role === 'user' ? text : '' const assistant = message.role === 'assistant' ? text : '' const tool = message.role === 'tool' ? text : '' @@ -65,6 +64,24 @@ export function insertSearchMessage( } } -export function chunkMessageText(text: string): string[] { - return [...textChunks(text)] +/** + * Deletes up to `limit` of a session's rows from `messages` and both FTS + * tables, in the caller's transaction, and reports how many went. Bounded + * because a retention sweep must not hold one transaction over a whole + * session; a replace passes no limit, since its rows and their replacements + * have to land together. + */ +export function deleteSearchMessages(db: SyncDatabase, sessionId: number, limit = -1): number { + const ids = db + .prepare('SELECT id FROM messages WHERE session_row_id = ? LIMIT ?') + .all(sessionId, limit) 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) + } + return ids.length } diff --git a/src/main/ai-vault-search/session-search-pending-deletes.ts b/src/main/ai-vault-search/session-search-pending-deletes.ts deleted file mode 100644 index 7de70a37940..00000000000 --- a/src/main/ai-vault-search/session-search-pending-deletes.ts +++ /dev/null @@ -1,40 +0,0 @@ -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. - // The batch half skips a batch whose session is already retired: that tombstone - // covers the same rows, and without the guard every reopen adds one more pass - // over them. - 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:'||b.id,b.session_row_id,b.id FROM search_write_batches b - WHERE NOT EXISTS (SELECT 1 FROM search_pending_deletes d - WHERE d.session_row_id=b.session_row_id AND d.batch_id IS NULL)`) -} diff --git a/src/main/ai-vault-search/session-search-retention-delete.test.ts b/src/main/ai-vault-search/session-search-retention-delete.test.ts index d9e15ac12d6..add6dcb3d57 100644 --- a/src/main/ai-vault-search/session-search-retention-delete.test.ts +++ b/src/main/ai-vault-search/session-search-retention-delete.test.ts @@ -1,21 +1,17 @@ import { expect, it } from 'vitest' -import { inSessionParseFileLane } from '../ai-vault/session-parse-file-lane' import type SyncDatabase from '../sqlite/sync-database' -import { discardSearchBatch } from './session-search-pending-deletes' import { deleteExpiredSearchFiles, RETENTION_DELETE_ROWS_PER_STEP } from './session-search-retention-delete' -import { openSessionSearchIndexFile } from './session-search-staged-write-test-fixture' +import { openSessionSearchIndexFile } from './session-search-index-test-fixture' import { SessionSearchStore } from './session-search-store' function seed(db: SyncDatabase, id: number, rows: number, mtime: number): void { - db.prepare(`INSERT INTO sessions(id,agent,session_id,file_path,title,cwd,cwd_key,resume_command) - VALUES (?, 'claude', ?, ?, 'synthetic retention', '/fixture', '/fixture', '')`).run( - id, - String(id), - String(id) - ) + db.prepare( + `INSERT INTO sessions(id,agent,session_id,file_path,title,cwd,cwd_key,resume_command) + VALUES (?, 'claude', ?, ?, 'synthetic retention', '/fixture', '/fixture', '')` + ).run(id, String(id), String(id)) db.prepare('INSERT INTO files(path,byte_offset,mtime_ms,session_row_id) VALUES (?,1,?,?)').run( String(id), mtime, @@ -35,14 +31,18 @@ function seed(db: SyncDatabase, id: number, rows: number, mtime: number): void { db.exec('COMMIT') } -/** Sessions a search over the visible views would still return. */ +/** + * Sessions a search would still return. Every retrieval joins a message to its + * session, which is what makes cutting the session loose enough to hide the + * whole thing while its rows are still being reclaimed. + */ function visibleSessionIds(db: SyncDatabase): string[] { return ( db .prepare( `SELECT DISTINCT s.session_id AS id FROM messages_fts - JOIN visible_messages m ON m.id = messages_fts.rowid - JOIN visible_sessions s ON s.id = m.session_row_id + JOIN messages m ON m.id = messages_fts.rowid + JOIN sessions s ON s.id = m.session_row_id WHERE messages_fts MATCH 'retentionneedle' ORDER BY s.session_id` ) .all() as { id: string }[] @@ -72,32 +72,31 @@ it('seeks the expiring end of the file list instead of scanning it', async () => } }) -it('yields within a large file while hiding partial rows and preserving unrelated sessions', async () => { +it('hides an expiring session at once, then reclaims its rows in bounded steps', async () => { const index = await openSessionSearchIndexFile('ss-retention-yield') seed(index.db, 1, 1025, 1) seed(index.db, 2, 1, 200) let previous = 1025 - let steps = 0 + const steps: number[] = [] try { await deleteExpiredSearchFiles( index.db, 100, () => false, - () => undefined, async () => { const left = count(index.db, 'messages WHERE session_row_id=1') - expect(previous - left).toBeLessThanOrEqual(RETENTION_DELETE_ROWS_PER_STEP) - expect(previous - left).toBeGreaterThan(0) + steps.push(previous - left) previous = left - steps++ - // The expiring session is hidden from the first step, never half-deleted. + // Cut loose in the very first transaction, so no query ever sees it with + // some of its messages already gone. expect(visibleSessionIds(index.db)).toEqual(['2']) } ) - expect(steps).toBe(5) + // The file transaction, then one bounded batch per step until the rows are gone. + expect(steps).toEqual([0, RETENTION_DELETE_ROWS_PER_STEP, 256, 256, 256, 1]) expect(count(index.db, 'messages_fts')).toBe(1) expect(count(index.db, 'conversation_fts')).toBe(1) - expect(count(index.db, 'search_pending_deletes')).toBe(0) + expect(count(index.db, 'sessions')).toBe(1) } finally { await index.close() } @@ -107,26 +106,34 @@ it('finishes an interrupted deletion after reopening', async () => { const index = await openSessionSearchIndexFile('ss-retention-resume') let store = new SessionSearchStore(index.path) let closed = false + let steps = 0 try { seed(index.db, 1, 513, 1) await deleteExpiredSearchFiles( index.db, 100, () => closed, - () => undefined, async () => { - store.close() - closed = true + if (++steps === 2) { + store.close() + closed = true + } } ) + // Some rows went, the rest did not, and nothing recorded that anywhere. + const stranded = count(index.db, 'messages') + expect(stranded).toBeGreaterThan(0) + expect(stranded).toBeLessThan(513) + expect(visibleSessionIds(index.db)).toEqual([]) + store = new SessionSearchStore(index.path) closed = false - expect(visibleSessionIds(index.db)).toEqual([]) - // A durable tombstone survives the restart, so the rest goes even with - // retention now unlimited. + // Rows nothing points at are the whole record of unfinished work, so the + // rest goes even with retention now unlimited. await store.purgeOlderThan(null) + expect(count(index.db, 'messages')).toBe(0) expect(count(index.db, 'messages_fts')).toBe(0) - expect(count(index.db, 'sessions')).toBe(0) + expect(count(index.db, 'conversation_fts')).toBe(0) } finally { if (!closed) { store.close() @@ -135,33 +142,6 @@ it('finishes an interrupted deletion after reopening', async () => { } }) -it('does not orphan a replacement file when resuming an older deletion for the same path', async () => { - const index = await openSessionSearchIndexFile('ss-retention-reused-path') - const store = new SessionSearchStore(index.path) - try { - seed(index.db, 1, 2, 1) - index.db.exec( - "INSERT INTO search_pending_deletes(path,session_row_id) VALUES ('1',1); DELETE FROM files WHERE path='1'" - ) - seed(index.db, 2, 2, 1) - index.db.exec("UPDATE files SET path='1' WHERE path='2'") - await store.purgeOlderThan(100) - for (const table of [ - 'messages', - 'messages_fts', - 'conversation_fts', - 'sessions', - 'files', - 'search_pending_deletes' - ]) { - expect(count(index.db, table)).toBe(0) - } - } finally { - store.close() - await index.close() - } -}) - it('cancels retention between batches and resumes without exposing a partial session', async () => { const index = await openSessionSearchIndexFile('ss-retention-cancel') const store = new SessionSearchStore(index.path) @@ -183,65 +163,31 @@ it('cancels retention between batches and resumes without exposing a partial ses } }) -it('never leaves a batch tombstone behind the batch row it names', async () => { - const index = await openSessionSearchIndexFile('ss-retention-orphan-batch') - const db = index.db +it('keeps a file a read refreshed after the expiry list was taken', async () => { + const index = await openSessionSearchIndexFile('ss-retention-refreshed') try { - seed(db, 1, 1, 1) - const batchId = Number( - db.prepare('INSERT INTO search_write_batches(session_row_id) VALUES (?)').run(1) - .lastInsertRowid - ) - db.prepare("INSERT INTO messages(session_row_id,batch_id,role) VALUES (1,?,'user')").run( - batchId - ) - - // Retention snapshots an empty tombstone set, then blocks on this path's - // parse lane, which an in-flight read of the same file holds. - let release = (): void => undefined - const held = new Promise((resolve) => { - release = resolve - }) - const lane = inSessionParseFileLane('1', () => held) - const purge = deleteExpiredSearchFiles( - db, + seed(index.db, 1, 2, 1) + seed(index.db, 2, 2, 2) + let refreshed = false + // The scan of `files` happens once, up front. A read of the second transcript + // lands while the first is being deleted, which makes it new enough to keep. + await deleteExpiredSearchFiles( + index.db, 100, () => false, - () => undefined + async () => { + if (!refreshed) { + refreshed = true + index.db.prepare('UPDATE files SET mtime_ms = 500 WHERE path = ?').run('2') + } + } ) - await new Promise((resolve) => setImmediate(resolve)) - // That read ends without publishing. It appended onto a session that is - // live, so discard writes a batch-keyed tombstone instead of retiring one. - discardSearchBatch(db, 1, batchId, false) - release() - await lane - await purge - - // Retiring the session took the batch row with it, which frees the rowid. - expect(count(db, 'search_write_batches')).toBe(0) - expect(count(db, 'search_pending_deletes')).toBe(0) - - // The next read takes that rowid. A surviving tombstone for it would delete - // this batch's staged rows, and the session would publish missing messages. - seed(db, 2, 1, 100) - const reused = Number( - db.prepare('INSERT INTO search_write_batches(session_row_id) VALUES (?)').run(2) - .lastInsertRowid - ) - expect(reused).toBe(batchId) - const staged = Number( - db - .prepare("INSERT INTO messages(session_row_id,batch_id,role) VALUES (2,?,'user')") - .run(reused).lastInsertRowid - ) - await deleteExpiredSearchFiles( - db, - null, - () => false, - () => undefined - ) - expect(db.prepare('SELECT id FROM messages WHERE id=?').get(staged)).toBeDefined() + // Only the per-file transaction re-reading the mtime it is about to act on + // keeps that session; the list it came from says both should go. + expect(count(index.db, 'files')).toBe(1) + expect(visibleSessionIds(index.db)).toEqual(['2']) + expect(count(index.db, 'messages')).toBe(2) } finally { await index.close() } diff --git a/src/main/ai-vault-search/session-search-retention-delete.ts b/src/main/ai-vault-search/session-search-retention-delete.ts index 99957a87299..271e63b374a 100644 --- a/src/main/ai-vault-search/session-search-retention-delete.ts +++ b/src/main/ai-vault-search/session-search-retention-delete.ts @@ -1,102 +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' +import { deleteSearchMessages } from './session-search-message-rows' export const RETENTION_DELETE_ROWS_PER_STEP = 256 +// Why in step with the deletes rather than one sweep at the end: `auto_vacuum = +// INCREMENTAL` holds every freed page until something asks for it back, and +// asking for a whole purge's worth at once is one long stall (40 ms per 22 MB +// freed, measured) instead of many short ones. +const RECLAIM_PAGES_PER_STEP = 2000 -/** A durable tombstone hides partial deletes and lets a reopened store finish them. */ +/** + * Drops every file older than the cutoff, then hands its rows back in bounded + * steps. + * + * The two halves are separate on purpose. Cutting a session loose from its file + * is one small transaction, and it is what makes the session stop answering + * searches — every read joins `sessions`, so a row whose session is gone is + * already unreachable. Reclaiming those rows is the expensive half, and it can + * be paused, interrupted or resumed at any point without a reader ever seeing a + * session that is half deleted. A crash in the middle leaves rows nothing + * points at, and `drainOrphanedMessages` finds them on the next pass. + */ 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 (cutoffMs !== null) { + const expired = db + .prepare('SELECT path FROM files WHERE mtime_ms < ? ORDER BY mtime_ms') + .all(cutoffMs) as { path: string }[] + for (const { path } of expired) { if (closed()) { return } db.exec('BEGIN IMMEDIATE') try { + // Re-read under the lock: a read of this file may have landed since the + // list was taken, which makes it new enough to keep. 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 + .prepare('SELECT session_row_id FROM files WHERE path = ? AND mtime_ms < ?') + .get(path, cutoffMs) 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 sessions WHERE id = ?').run(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 - ) - // Invariant: a batch tombstone never outlives its batch row, or the - // freed rowid comes back and the tombstone deletes a live batch's rows. - db.prepare( - 'DELETE FROM search_pending_deletes WHERE session_row_id=? AND batch_id IS NOT NULL' - ).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() - } - }) + await yieldStep() + } } + await drainOrphanedMessages(db, closed, yieldStep) +} + +/** + * Deletes rows whose session no longer exists, a bounded batch per transaction. + * + * That set is exactly what retention, a removed source and an interrupted + * earlier drain leave behind, so the index needs no record of unfinished work + * beyond the rows themselves. + */ +export async function drainOrphanedMessages( + db: SyncDatabase, + closed: () => boolean, + yieldStep: () => Promise = yieldToEventLoop +): Promise { + // Ordered by session so one call to this walks a session's rows to the end + // before paying for the scan that finds the next one. + const nextOrphan = db.prepare( + `SELECT session_row_id FROM messages + WHERE session_row_id NOT IN (SELECT id FROM sessions) LIMIT 1` + ) + let orphan = (nextOrphan.get() as { session_row_id: number } | undefined)?.session_row_id + while (orphan !== undefined && !closed()) { + db.exec('BEGIN IMMEDIATE') + let deleted = 0 + try { + deleted = deleteSearchMessages(db, orphan, RETENTION_DELETE_ROWS_PER_STEP) + db.exec('COMMIT') + } catch (error) { + db.exec('ROLLBACK') + throw error + } + db.pragma(`incremental_vacuum(${RECLAIM_PAGES_PER_STEP})`) + if (deleted < RETENTION_DELETE_ROWS_PER_STEP) { + orphan = (nextOrphan.get() as { session_row_id: number } | undefined)?.session_row_id + } + await yieldStep() + } + // A `removeFile` frees its pages outside this loop and may leave none to drain. + db.pragma(`incremental_vacuum(${RECLAIM_PAGES_PER_STEP})`) } diff --git a/src/main/ai-vault-search/session-search-schema.test.ts b/src/main/ai-vault-search/session-search-schema.test.ts index da85360f0bd..eb09a4ce86e 100644 --- a/src/main/ai-vault-search/session-search-schema.test.ts +++ b/src/main/ai-vault-search/session-search-schema.test.ts @@ -12,9 +12,7 @@ import SyncDatabase from '../sqlite/sync-database' import { SESSION_SEARCH_SCHEMA_VERSION, openSessionSearchDatabase, - removeSessionSearchDatabase, - VISIBLE_MESSAGES, - VISIBLE_SESSIONS + removeSessionSearchDatabase } from './session-search-schema' const recordedRmSync = vi.hoisted(() => vi.fn()) @@ -59,7 +57,9 @@ describe('openSessionSearchDatabase', () => { 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 }) + expect(second.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ + c: 1 + }) second.close() }) @@ -77,7 +77,9 @@ describe('openSessionSearchDatabase', () => { 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 }) + 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). @@ -85,28 +87,15 @@ describe('openSessionSearchDatabase', () => { 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' }) + await expect(stat(`${path}${suffix}`)).rejects.toMatchObject({ + code: 'ENOENT' + }) } }) }) @@ -124,7 +113,9 @@ it('rebuilds a file too corrupt to open instead of refusing forever', async () = const rebuilt = openSessionSearchDatabase(path) try { expect(schemaVersion(rebuilt)).toBe(String(SESSION_SEARCH_SCHEMA_VERSION)) - expect(rebuilt.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ c: 0 }) + expect(rebuilt.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ + c: 0 + }) } finally { rebuilt.close() } @@ -147,7 +138,9 @@ it('gives up rather than looping when a fresh file still cannot be opened', asyn await writeFile(path, 'not a SQLite database') // Every open of this path fails, so the one permitted retry is exhausted. const open = vi.spyOn(SyncDatabase.prototype, 'pragma').mockImplementation(() => { - throw Object.assign(new Error('database disk image is malformed'), { code: 'SQLITE_CORRUPT' }) + throw Object.assign(new Error('database disk image is malformed'), { + code: 'SQLITE_CORRUPT' + }) }) try { expect(() => openSessionSearchDatabase(path)).toThrow(/malformed/) @@ -165,7 +158,9 @@ it('surfaces the unlink failure itself when a stale index cannot be removed', as stale.close() recordedRmSync.mockReset() recordedRmSync.mockImplementation(() => { - throw Object.assign(new Error('EPERM: operation not permitted, unlink'), { code: 'EPERM' }) + throw Object.assign(new Error('EPERM: operation not permitted, unlink'), { + code: 'EPERM' + }) }) try { // The stale handle is closed before the unlink, so the failure path must not @@ -204,7 +199,9 @@ it('rebuilds a newer index rather than reading a schema it does not know', async const rebuilt = openSessionSearchDatabase(path) try { expect(schemaVersion(rebuilt)).toBe(String(SESSION_SEARCH_SCHEMA_VERSION)) - expect(rebuilt.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ c: 0 }) + expect(rebuilt.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ + c: 0 + }) } finally { rebuilt.close() } @@ -222,7 +219,9 @@ it('rebuilds when meta exists but its version row is gone', async () => { const rebuilt = openSessionSearchDatabase(path) try { expect(schemaVersion(rebuilt)).toBe(String(SESSION_SEARCH_SCHEMA_VERSION)) - expect(rebuilt.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ c: 0 }) + expect(rebuilt.prepare('SELECT COUNT(*) AS c FROM files').get()).toEqual({ + c: 0 + }) } finally { rebuilt.close() } @@ -245,98 +244,50 @@ it('opens with the pragmas the write path depends on', async () => { } }) -it('retires an unfinished batch as part of opening, not of using', async () => { - const path = await tempDatabasePath() - const crashed = openSessionSearchDatabase(path) - crashed.exec(`INSERT INTO sessions(id,index_ready,agent,session_id,file_path,title,resume_command) - VALUES (1,0,'claude','a','a','staging',''); - INSERT INTO search_write_batches(id,session_row_id) VALUES (7,1)`) - crashed.close() - - const reopened = openSessionSearchDatabase(path) +it("walks a session's rows through an index rather than scanning the table", async () => { + const db = openSessionSearchDatabase(await tempDatabasePath()) try { - // One tombstone, not two: the session-keyed one already covers every row the - // batch holds, and a per-reopen duplicate would replay the same deletes. - expect(reopened.prepare('SELECT path FROM search_pending_deletes').all()).toEqual([ - { path: '\u0000session:1' } - ]) + // The replace delete and the orphan drain both take this path, once per file. + const plan = ( + db + .prepare('EXPLAIN QUERY PLAN SELECT id FROM messages WHERE session_row_id = ? LIMIT ?') + .all(1, 1) as { detail: string }[] + ) + .map((row) => row.detail) + .join(' ') + expect(plan).toContain('messages_session') } finally { - reopened.close() + db.close() } }) -it('retires a batch whose session survived it', async () => { - const path = await tempDatabasePath() - const crashed = openSessionSearchDatabase(path) - crashed.exec(`INSERT INTO sessions(id,index_ready,agent,session_id,file_path,title,resume_command) - VALUES (1,1,'claude','a','a','published',''); - INSERT INTO search_write_batches(id,session_row_id) VALUES (7,1)`) - crashed.close() - - const reopened = openSessionSearchDatabase(path) +it('keeps only the session indexes a retrieval query can seek', async () => { + const db = openSessionSearchDatabase(await tempDatabasePath()) try { - expect(reopened.prepare('SELECT path,batch_id FROM search_pending_deletes').all()).toEqual([ - { path: '\u0000batch:7', batch_id: 7 } - ]) + const names = ( + db + .prepare("SELECT name FROM sqlite_master WHERE type='index' AND tbl_name='sessions'") + .all() as { name: string }[] + ) + .map((row) => row.name) + .sort() + // One per shape PR 4's retrieval seeks: the agent filter, the newest-first + // order and date window, and the folder-prefix range scan. Fork folding reads + // `content_hash` off rows it already holds, so that column is not indexed. + expect(names).toEqual(['sessions_agent', 'sessions_cwd_key', 'sessions_updated_at']) } finally { - reopened.close() + db.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'), - (3,3,NULL,'user')`) - expect(db.prepare(`SELECT title FROM ${VISIBLE_SESSIONS} ORDER BY id`).all()).toEqual([ - { title: 'published' } - ]) - // Row 3 is the one a batch-pointer filter alone would return: its batch is - // long gone and only the session-keyed tombstone retires it. - 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('subtracts the tombstones through an index rather than scanning them', async () => { - const db = openSessionSearchDatabase(await tempDatabasePath()) - try { - for (const view of [VISIBLE_MESSAGES, VISIBLE_SESSIONS]) { - const plan = ( - db.prepare(`EXPLAIN QUERY PLAN SELECT id FROM ${view}`).all() as { detail: string }[] - ) - .map((row) => row.detail) - .join(' ') - // Without the partial index this reads "SCAN search_pending_deletes", - // once per statement, on every read either view serves. - expect(plan).toContain('search_pending_deletes_session') - } - } 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' }) + const locked = Object.assign(new Error('EPERM: operation not permitted'), { + code: 'EPERM' + }) recordedRmSync.mockImplementationOnce(() => { throw locked }) diff --git a/src/main/ai-vault-search/session-search-schema.ts b/src/main/ai-vault-search/session-search-schema.ts index 3d04a985852..7ad414aa3aa 100644 --- a/src/main/ai-vault-search/session-search-schema.ts +++ b/src/main/ai-vault-search/session-search-schema.ts @@ -2,7 +2,6 @@ import { mkdirSync } from 'node:fs' import { dirname } from 'node:path' import SyncDatabase from '../sqlite/sync-database' import { removeTreeSync } from '../../shared/windows-transient-lock-removal' -import { recoverSearchWrites } from './session-search-pending-deletes' // The index stores transcript content as written, with no redaction. A secret in // a transcript is already plaintext under the user's home directory and is @@ -11,7 +10,7 @@ import { recoverSearchWrites } from './session-search-pending-deletes' // policy, decided where the wire is. // Bump to drop and rebuild: the index is a cache over the transcripts, never a source. -export const SESSION_SEARCH_SCHEMA_VERSION = 1 +export const SESSION_SEARCH_SCHEMA_VERSION = 2 // unicode61 keeps `_ . - /` inside tokens so paths and identifiers match exactly; // the `identifiers` column carries the split form (see session-search-identifier-split). @@ -19,15 +18,10 @@ export const SESSION_SEARCH_SCHEMA_VERSION = 1 // 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, -- The transcript this session was decoded from. Not unique: OpenCode's SQLite @@ -48,7 +42,6 @@ CREATE TABLE IF NOT EXISTS sessions( 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( @@ -62,50 +55,20 @@ CREATE TABLE IF NOT EXISTS files( ); -- Retention walks the expiring end of this column; without it that is a full scan and a sort. CREATE INDEX IF NOT EXISTS files_mtime ON files(mtime_ms); -CREATE TABLE IF NOT EXISTS search_pending_deletes( - path TEXT PRIMARY KEY, - session_row_id INTEGER NOT NULL, - batch_id INTEGER -); --- Both views subtract this set on every read; without it each one scans the table. --- Partial because a batch-keyed tombstone names a session that is still visible. -CREATE INDEX IF NOT EXISTS search_pending_deletes_session - ON search_pending_deletes(session_row_id) WHERE batch_id IS NULL; --- 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 ); +-- Both the replace delete and the orphan drain walk a session's rows through this. 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. --- Both halves subtract the same session-keyed tombstones. A message outlives its --- session row until the cleanup lane reaches it, so filtering messages on the --- batch pointer alone would show a replaced generation beside its successor and --- would keep answering for a file that was already removed. -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 - AND session_row_id NOT IN - (SELECT session_row_id FROM search_pending_deletes WHERE batch_id IS NULL); ` /** @@ -154,9 +117,6 @@ function openExisting(path: string): SyncDatabase { 'schema_version', String(SESSION_SEARCH_SCHEMA_VERSION) ) - // Part of opening, not of using: a batch or staging session that outlived its - // writer has to be tombstoned before anything can read or write past it. - recoverSearchWrites(db) return db } catch (error) { db?.close() @@ -186,6 +146,9 @@ function openWithPragmas(path: string): SyncDatabase { // 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') + // The whole consistency model: a file's rows and its cursor land in one + // transaction, and a reader on another handle sees the last committed state + // of the index rather than a session half way through being rewritten. db.pragma('journal_mode = WAL') db.pragma('synchronous = NORMAL') db.pragma('journal_size_limit = 8388608') diff --git a/src/main/ai-vault-search/session-search-staged-write.test.ts b/src/main/ai-vault-search/session-search-staged-write.test.ts deleted file mode 100644 index cf287fdca13..00000000000 --- a/src/main/ai-vault-search/session-search-staged-write.test.ts +++ /dev/null @@ -1,269 +0,0 @@ -import { afterEach, beforeEach, expect, it } from 'vitest' -import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers' -import type SyncDatabase from '../sqlite/sync-database' -import { registerSessionSearchIndexConsumer } from './session-search-index-consumer' -import { SEARCH_WRITE_ROWS_PER_STEP } from './session-search-index-writer' -import { - openSessionSearchIndexFile, - replayTranscriptRead, - syntheticCandidate, - syntheticSession, - SYNTHETIC_TRANSCRIPT, - userMessages, - type SessionSearchIndexFile -} from './session-search-staged-write-test-fixture' -import { SessionSearchStore } from './session-search-store' - -let index: SessionSearchIndexFile -let store: SessionSearchStore -let errors: unknown[] - -beforeEach(async () => { - index = await openSessionSearchIndexFile('ss-staged-write') - errors = [] - store = new SessionSearchStore(index.path, (error) => errors.push(error)) - registerSessionSearchIndexConsumer(store) -}) - -afterEach(async () => { - resetTranscriptConsumersForTests() - store.close() - await index.close() -}) - -/** - * The read shape every query site has to use: an FTS table joined to the - * visible half. Both views subtract tombstones, so a hit here is a row a search - * may return right now, with no cleanup pass required first. - */ -function matches(db: SyncDatabase, table: string, term: string): number { - return ( - db - .prepare( - `SELECT count(*) AS n FROM ${table} JOIN visible_messages m ON m.id = ${table}.rowid - WHERE ${table} MATCH ?` - ) - .get(term) as { n: number } - ).n -} - -function counts(db: SyncDatabase): Record { - const one = (sql: string): number => (db.prepare(sql).get() as { n: number }).n - return { - sessions: one('SELECT count(*) AS n FROM visible_sessions'), - messages: one('SELECT count(*) AS n FROM visible_messages'), - rawMessages: one('SELECT count(*) AS n FROM messages'), - batches: one('SELECT count(*) AS n FROM search_write_batches'), - tombstones: one('SELECT count(*) AS n FROM search_pending_deletes'), - full: one('SELECT count(*) AS n FROM messages_fts'), - conversation: one('SELECT count(*) AS n FROM conversation_fts') - } -} - -it('publishes a whole read atomically and leaves no batch behind', async () => { - replayTranscriptRead({ messages: userMessages('needle text', 300) }) - await store.settled() - - const after = counts(index.db) - expect(after.sessions).toBe(1) - expect(after.messages).toBe(300) - expect(after.batches).toBe(0) - expect(after.tombstones).toBe(0) - expect(errors).toEqual([]) -}) - -it('writes both FTS tables for every published conversational row', async () => { - replayTranscriptRead({ - messages: [ - { role: 'user', text: 'alpha question', timestamp: null }, - { role: 'assistant', text: 'beta answer', timestamp: null }, - { role: 'tool', text: 'gamma tool output', timestamp: null } - ] - }) - await store.settled() - - // messages_fts carries every row; conversation_fts is the tool-free half. - expect(counts(index.db).full).toBe(3) - expect(counts(index.db).conversation).toBe(2) - expect(matches(index.db, 'messages_fts', 'gamma')).toBe(1) - expect(matches(index.db, 'conversation_fts', 'gamma')).toBe(0) - expect(matches(index.db, 'conversation_fts', 'beta')).toBe(1) -}) - -it('hides staged rows from both halves until the read finishes', async () => { - const candidate = syntheticCandidate() - const staged = store.beginWrite(candidate, 'replace', 0)! - for (const message of userMessages('stagedneedle', 200)) { - staged.add(message) - } - - // Rows really are on disk mid-read; only the views hold them back. - const mid = counts(index.db) - expect(mid.rawMessages).toBe(SEARCH_WRITE_ROWS_PER_STEP) - expect(mid.messages).toBe(0) - expect(mid.sessions).toBe(0) - - expect(staged.publish({ session: syntheticSession(), byteOffset: 4096, incomplete: false })).toBe( - true - ) - staged.discard() - expect(counts(index.db).messages).toBe(200) - expect(counts(index.db).sessions).toBe(1) -}) - -it('tombstones an abandoned batch instead of publishing half a file', async () => { - const staged = store.beginWrite(syntheticCandidate(), 'replace', 0)! - for (const message of userMessages('abandoned', 50)) { - staged.add(message) - } - staged.discard() - - expect(counts(index.db).messages).toBe(0) - expect(counts(index.db).sessions).toBe(0) - await store.purgeOlderThan(null) - expect(counts(index.db).rawMessages).toBe(0) - expect(counts(index.db).tombstones).toBe(0) -}) - -it('recovers a batch its writer never finished when the store reopens', async () => { - const staged = store.beginWrite(syntheticCandidate(), 'replace', 0)! - for (const message of userMessages('crashed', 20)) { - staged.add(message) - } - // No discard and no publish: the process died mid-write. - store.close() - - const reopened = new SessionSearchStore(index.path, (error) => errors.push(error)) - try { - expect(counts(index.db).sessions).toBe(0) - expect(counts(index.db).messages).toBe(0) - await reopened.purgeOlderThan(null) - expect(counts(index.db).rawMessages).toBe(0) - expect(counts(index.db).batches).toBe(0) - } finally { - reopened.close() - store = new SessionSearchStore(index.path, (error) => errors.push(error)) - } -}) - -it('retires a batch appended onto a live session when the writer dies', async () => { - replayTranscriptRead({ messages: userMessages('published', 3), outcome: { byteOffset: 100 } }) - await store.settled() - - const staged = store.beginWrite(syntheticCandidate(), 'append', 100)! - for (const message of userMessages('crashedappend', 200)) { - staged.add(message) - } - // No discard and no publish: the process died mid-append. The staged rows hang - // off a session that is published and must stay readable. - store.close() - - const reopened = new SessionSearchStore(index.path, (error) => errors.push(error)) - try { - expect(counts(index.db).messages).toBe(3) - expect(counts(index.db).sessions).toBe(1) - await reopened.purgeOlderThan(null) - // The appended rows are gone and the published generation survived them. - expect(counts(index.db).rawMessages).toBe(3) - expect(counts(index.db).batches).toBe(0) - expect(counts(index.db).tombstones).toBe(0) - expect(counts(index.db).full).toBe(3) - expect(counts(index.db).conversation).toBe(3) - expect(reopened.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(100) - } finally { - reopened.close() - store = new SessionSearchStore(index.path, (error) => errors.push(error)) - } -}) - -it('closes twice without turning the second call into an error', () => { - store.close() - // node:sqlite throws ERR_INVALID_STATE on a second close of one handle, and a - // store is closed both by whoever owns it and by a teardown that cannot know. - expect(() => store.close()).not.toThrow() - store = new SessionSearchStore(index.path, (error) => errors.push(error)) -}) - -it('drains the rows an incomplete read staged, without waiting for another write', async () => { - replayTranscriptRead({ - messages: userMessages('incompleteread', 300), - outcome: { incomplete: true } - }) - // Two flushes reached disk before the read turned out not to be publishable. - expect(counts(index.db).rawMessages).toBe(2 * SEARCH_WRITE_ROWS_PER_STEP) - await store.settled() - - const after = counts(index.db) - expect(after.rawMessages).toBe(0) - expect(after.full).toBe(0) - expect(after.conversation).toBe(0) - expect(after.tombstones).toBe(0) - expect(after.batches).toBe(0) - // The file is still owed a whole re-read; only its staged rows are gone. - expect(store.pendingFileCount).toBe(1) - expect(errors).toEqual([]) -}) - -it('drains what a dead writer left, under one tombstone, when the store reopens', async () => { - const staged = store.beginWrite(syntheticCandidate(), 'replace', 0)! - for (const message of userMessages('crashedrows', 200)) { - staged.add(message) - } - // No discard and no publish: the process died mid-write. - store.close() - - const reopened = new SessionSearchStore(index.path, (error) => errors.push(error)) - try { - // One tombstone covers the staging session and its batch alike; a second - // would replay the same deletes, and would be written again on every reopen. - expect(counts(index.db).tombstones).toBe(1) - await reopened.settled() - const after = counts(index.db) - expect(after.rawMessages).toBe(0) - expect(after.full).toBe(0) - expect(after.tombstones).toBe(0) - expect(after.batches).toBe(0) - expect(errors).toEqual([]) - } finally { - reopened.close() - store = new SessionSearchStore(index.path, (error) => errors.push(error)) - } -}) - -it('replaces the previous generation without ever showing both', async () => { - replayTranscriptRead({ messages: userMessages('firstgeneration', 10) }) - await store.settled() - replayTranscriptRead({ messages: userMessages('secondgeneration', 10) }) - await store.settled() - await store.purgeOlderThan(null) - - expect(counts(index.db).sessions).toBe(1) - expect(counts(index.db).messages).toBe(10) - expect(matches(index.db, 'messages_fts', 'firstgeneration')).toBe(0) - expect(matches(index.db, 'messages_fts', 'secondgeneration')).toBe(10) -}) - -// The cleanup lane cannot start before the first microtask, so everything below -// reads the index in the state a search issued in the same tick would see. -it('answers for only the new generation the moment a replace publishes', async () => { - replayTranscriptRead({ messages: userMessages('firstgeneration', 3) }) - replayTranscriptRead({ messages: userMessages('secondgeneration', 3) }) - - expect(matches(index.db, 'messages_fts', 'firstgeneration')).toBe(0) - expect(matches(index.db, 'messages_fts', 'secondgeneration')).toBe(3) - expect(counts(index.db).messages).toBe(3) - // Both generations really are still on disk; only the views subtract one. - expect(counts(index.db).rawMessages).toBe(6) - await store.settled() -}) - -it('stops answering for a removed file the moment it is removed', async () => { - replayTranscriptRead({ messages: userMessages('removedneedle', 3) }) - store.removeFile(SYNTHETIC_TRANSCRIPT) - - expect(counts(index.db).sessions).toBe(0) - expect(counts(index.db).messages).toBe(0) - expect(matches(index.db, 'messages_fts', 'removedneedle')).toBe(0) - expect(counts(index.db).rawMessages).toBe(3) - await store.settled() -}) diff --git a/src/main/ai-vault-search/session-search-store.ts b/src/main/ai-vault-search/session-search-store.ts index 7f112217704..8e49ff92ce1 100644 --- a/src/main/ai-vault-search/session-search-store.ts +++ b/src/main/ai-vault-search/session-search-store.ts @@ -1,13 +1,12 @@ import type SyncDatabase from '../sqlite/sync-database' import type { SessionFileCandidate } from '../ai-vault/session-scanner-types' -import { compactSessionSearchIndex } from './session-search-index-compaction' import type { SessionSearchFileIdentity, SessionSearchIndexedFile } from './session-search-file-cursor' import { SessionSearchIndexWriter, - type SessionSearchStagedWrite + type SessionSearchFileWrite } from './session-search-index-writer' import { deleteExpiredSearchFiles } from './session-search-retention-delete' import { openSessionSearchDatabase } from './session-search-schema' @@ -17,11 +16,6 @@ import { openSessionSearchDatabase } from './session-search-schema' // that silently forgets is worse than one that says it is incomplete. export const STALE_PATH_LIMIT = 20_000 -export type SessionSearchStoreOptions = { - /** The WAL backlog a staging write refuses to grow past. Only tests narrow it. */ - walBudgetBytes?: number -} - /** * Owns the index database. PR 2 scope: the write half only — the transcript * consumer writes through it and nothing reads from it yet. Lifecycle (who @@ -33,10 +27,6 @@ export class SessionSearchStore { private closed = false private acceptingWrites = true private retentionCutoffMs: number | null = null - private cleanupRequested = false - private cleanup: Promise | null = null - private lastIndexedAt: string | null = null - private writeFailures = 0 // Files this index knows it is behind on. Filled by a declined or abandoned // read; PR 3's indexer drains it. Nothing here schedules the re-read. private readonly stale = new Map() @@ -48,15 +38,10 @@ export class SessionSearchStore { console.warn( '[ai-vault-search] index write failed:', error instanceof Error ? error.name : 'IndexError' - ), - options: SessionSearchStoreOptions = {} + ) ) { this.db = openSessionSearchDatabase(path) - this.writer = new SessionSearchIndexWriter(this.db, options.walBudgetBytes) - // Opening tombstones whatever a dead writer left staged. Nothing else will - // schedule that drain: a store that is only ever read from, or one whose - // next read declines, would carry those rows for the life of the index. - this.scheduleCleanup() + this.writer = new SessionSearchIndexWriter(this.db) } setAcceptingWrites(accept: boolean): void { @@ -92,7 +77,7 @@ export class SessionSearchStore { candidate: SessionFileCandidate, mode: 'replace' | 'append', previousByteOffset: number - ): SessionSearchStagedWrite | null { + ): SessionSearchFileWrite | null { if (!this.acceptsCandidate(candidate)) { return null } @@ -104,27 +89,13 @@ export class SessionSearchStore { } } - writePublished(candidate: SessionFileCandidate): void { + writeCommitted(candidate: SessionFileCandidate): void { // Why: a list scan queues every file the backfill has not reached yet; once // one lands, a later pass must not re-read the whole queue. this.stale.delete(candidate.file.path) - this.lastIndexedAt = new Date().toISOString() - this.scheduleCleanup() - } - - /** - * A read that staged rows and then could not publish them. The tombstone its - * discard wrote needs the same drain a publish gets, or the staged rows sit in - * `messages` and both FTS tables until some unrelated write happens to - * schedule a pass — which for the last read before a shutdown is never. - */ - writeAbandoned(candidate: SessionFileCandidate): void { - this.markStale(candidate) - this.scheduleCleanup() } reportWriteFailure(error: unknown): void { - this.writeFailures += 1 this.onError(error) } @@ -180,14 +151,6 @@ export class SessionSearchStore { return this.stale.size } - get lastWriteAt(): string | null { - return this.lastIndexedAt - } - - get failures(): number { - return this.writeFailures - } - /** * Drops a source's rows. Only a proven deletion may call this: an unreadable * source is `unverifiable`, not `missing`, and keeps its rows @@ -197,24 +160,19 @@ export class SessionSearchStore { this.stale.delete(path) try { this.writer.removeFile(path) - this.scheduleCleanup() } catch (error) { this.onError(error) } } - /** Hides expired sessions immediately, then removes their rows in resumable batches. */ + /** Cuts expired sessions loose at once, then reclaims their rows in resumable batches. */ async purgeOlderThan(cutoffMs: number | null, signal?: AbortSignal): Promise { try { await deleteExpiredSearchFiles( this.db, cutoffMs, - () => this.closed || signal?.aborted === true, - () => undefined + () => this.closed || signal?.aborted === true ) - if (!this.closed && !signal?.aborted) { - await compactSessionSearchIndex(this.db, () => this.closed || signal?.aborted === true) - } } catch (error) { if (!this.closed) { this.onError(error) @@ -231,45 +189,4 @@ export class SessionSearchStore { this.closed = true this.db.close() } - - /** Drains tombstones left by a publish or a removal, one file at a time. */ - scheduleCleanup(): void { - if (this.closed) { - return - } - if (this.cleanup) { - this.cleanupRequested = true - return - } - this.cleanupRequested = false - this.cleanup = deleteExpiredSearchFiles( - this.db, - null, - () => this.closed, - () => undefined - ) - .catch((error) => { - if (!this.closed) { - this.onError(error) - } - }) - .finally(() => { - this.cleanup = null - if (this.cleanupRequested) { - this.scheduleCleanup() - } - }) - } - - /** - * Tests only: the cleanup lane is fire-and-forget everywhere else. Loops - * because a pass requested while one was running is scheduled from the - * finished pass's own continuation, so awaiting a single promise would return - * with work still queued. - */ - async settled(): Promise { - while (this.cleanup) { - await this.cleanup - } - } } diff --git a/src/main/ai-vault-search/session-search-visible-read-ratchet.test.ts b/src/main/ai-vault-search/session-search-visible-read-ratchet.test.ts deleted file mode 100644 index 98fdeacf98d..00000000000 --- a/src/main/ai-vault-search/session-search-visible-read-ratchet.test.ts +++ /dev/null @@ -1,153 +0,0 @@ -import { readdir, readFile } from 'node:fs/promises' -import { join } from 'node:path' -import { expect, it } from 'vitest' -import { VISIBLE_MESSAGES, VISIBLE_SESSIONS } from './session-search-schema' - -// Why a ratchet and not a type: publish makes a read's rows visible atomically, -// but it does that by flipping `batch_id` and `index_ready`, not by moving the -// FTS rows — those land in the staging flushes. An FTS table on its own still -// holds staged and tombstoned rows, and only the views subtract them. Every read -// site has to join one, and nothing in SQL can force that. -// -// The unit is one SQL statement, not one file. A file is far too coarse: the -// query module this exists for will name both views somewhere, and that would -// whitelist every raw read in it. - -const FTS_TABLE = /\b(?:messages_fts|conversation_fts)\b/ -const SELECT = /\bSELECT\b/i -const VISIBLE_VIEW = new RegExp( - `\\b(?:${VISIBLE_MESSAGES}|${VISIBLE_SESSIONS}|VISIBLE_MESSAGES|VISIBLE_SESSIONS)\\b` -) -// A table name this scan cannot read: `FROM ' + table` or `FROM ${table}`. -const DYNAMIC_TABLE = /\b(?:FROM|JOIN)\s*(?:\$\{|$)/i -// A literal, plus any literals concatenated onto it: one statement, not several. -const SQL_FRAGMENT = - /(?:`(?:[^`\\]|\\[\s\S])*`|'(?:[^'\\\n]|\\.)*'|"(?:[^"\\\n]|\\.)*")(?:\s*\+\s*(?:`(?:[^`\\]|\\[\s\S])*`|'(?:[^'\\\n]|\\.)*'|"(?:[^"\\\n]|\\.)*"))*/g - -/** Statement-sized spans of SQL text, with the quotes and the `+` joins removed. */ -export function sqlStatements(source: string): string[] { - const statements: string[] = [] - for (const [fragment] of source.matchAll(SQL_FRAGMENT)) { - const sql = fragment - .split(/\s*\+\s*/) - .map((part) => part.slice(1, -1)) - .join('') - statements.push(...sql.split(';')) - } - return statements -} - -/** - * Reads that could return a staged or tombstoned row. A statement whose table - * name is assembled at runtime counts only in a file that names an FTS table - * somewhere, which is the shape a concatenated or interpolated read takes. - */ -export function unguardedFtsReads(source: string): string[] { - const fileNamesFts = FTS_TABLE.test(source) - return sqlStatements(source).filter((statement) => { - if (!SELECT.test(statement) || VISIBLE_VIEW.test(statement)) { - return false - } - return FTS_TABLE.test(statement) || (fileNamesFts && DYNAMIC_TABLE.test(statement)) - }) -} - -/** - * Statements that read an FTS table without subtracting staged rows and are - * allowed to. Every entry needs a reason, and the census fails on one that has - * stopped being necessary. - */ -const ALLOWED: Record = {} - -const ROOTS = ['src/main', 'src/relay', 'src/cli', 'src/shared', 'src/preload'] - -async function sourceFiles(root: string): Promise<{ name: string; text: string }[]> { - const out: { name: string; text: string }[] = [] - const walk = async (dir: string): Promise => { - for (const entry of await readdir(dir, { withFileTypes: true })) { - const path = join(dir, entry.name) - if (entry.isDirectory()) { - await walk(path) - continue - } - if ( - !entry.name.endsWith('.ts') || - entry.name.endsWith('.test.ts') || - entry.name.endsWith('-test-fixture.ts') - ) { - continue - } - out.push({ name: entry.name, text: await readFile(path, 'utf-8') }) - } - } - await walk(root) - return out -} - -it('flags a raw read and the two shapes that hide the table name', () => { - expect( - unguardedFtsReads( - String.raw`db.prepare('SELECT rowid FROM messages_fts WHERE messages_fts MATCH ?')` - ) - ).toHaveLength(1) - - // Concatenated: the table name is never inside the SELECT literal. - expect( - unguardedFtsReads(String.raw`const t = 'messages_fts'; db.prepare('SELECT rowid FROM ' + t)`) - ).toHaveLength(1) - - // Interpolated: same, through a template. - expect( - unguardedFtsReads( - // A plain double-quoted string, so the ${...} reaches the checker verbatim. - "const TABLE = 'conversation_fts'; db.prepare(`SELECT rowid FROM ${TABLE}`)" - ) - ).toHaveLength(1) - - // Concatenated literals are one statement, so the join is seen. - expect( - unguardedFtsReads( - String.raw`db.prepare('SELECT rowid FROM messages_fts JOIN ' + 'visible_messages m ON m.id = rowid')` - ) - ).toEqual([]) -}) - -it('flags a raw read in a file that names a view somewhere else entirely', () => { - // The shape a file-level scan cannot see: the view name is real, but it is in - // an unrelated string, so it guards nothing. - const source = String.raw` - const LABEL = 'visible_messages is the published half' - export function hits(db: Db) { - return db.prepare('SELECT rowid FROM messages_fts WHERE messages_fts MATCH ?').all() - } - ` - const offenders = unguardedFtsReads(source) - expect(offenders).toHaveLength(1) - expect(offenders[0]).toContain('messages_fts') -}) - -it('leaves a joined read alone even next to a raw table name in a delete', () => { - const source = String.raw` - const purge = db.prepare('DELETE FROM messages_fts WHERE rowid = ?') - const read = db.prepare('SELECT m.id FROM conversation_fts JOIN visible_messages m ON m.id = conversation_fts.rowid') - ` - expect(unguardedFtsReads(source)).toEqual([]) -}) - -it('reads every FTS table through a visibility view across every bundled root', async () => { - const files = (await Promise.all(ROOTS.map((root) => sourceFiles(root)))).flat() - expect(files.length).toBeGreaterThan(500) - // The scan really reaches the modules that name these tables. - expect(files.some((file) => FTS_TABLE.test(file.text))).toBe(true) - - const offenders = files - .filter((file) => unguardedFtsReads(file.text).length > 0 && !(file.name in ALLOWED)) - .map((file) => file.name) - expect(offenders).toEqual([]) - - // A stale exemption is an unguarded read waiting to happen. - const unused = Object.keys(ALLOWED).filter( - (name) => !files.some((file) => file.name === name && unguardedFtsReads(file.text).length > 0) - ) - expect(unused).toEqual([]) -}) diff --git a/src/main/ai-vault-search/session-search-wal-budget.test.ts b/src/main/ai-vault-search/session-search-wal-budget.test.ts deleted file mode 100644 index 9e7d71333b8..00000000000 --- a/src/main/ai-vault-search/session-search-wal-budget.test.ts +++ /dev/null @@ -1,93 +0,0 @@ -import { afterEach, expect, it } from 'vitest' -import { resetTranscriptConsumersForTests } from '../ai-vault/session-transcript-consumers' -import SyncDatabase from '../sqlite/sync-database' -import { registerSessionSearchIndexConsumer } from './session-search-index-consumer' -import { - openSessionSearchIndexFile, - replayTranscriptRead, - SYNTHETIC_TRANSCRIPT, - userMessages -} from './session-search-staged-write-test-fixture' -import { SessionSearchStore } from './session-search-store' -import { assertSearchWalBudget, SearchWalBackpressureError } from './session-search-wal-budget' - -afterEach(() => { - resetTranscriptConsumersForTests() -}) - -it('backpressures a pinned snapshot and resumes checkpoints after that reader releases', async () => { - const index = await openSessionSearchIndexFile('ss-wal-budget') - const reader = new SyncDatabase(index.path, { readonly: true }) - try { - assertSearchWalBudget(index.db) - reader.exec('BEGIN') - reader.prepare('SELECT count(*) FROM sessions').get() - index.db - .prepare( - 'INSERT INTO sessions(agent,session_id,file_path,title,resume_command) VALUES (?,?,?,?,?)' - ) - .run('claude', 'a', 'a', 'synthetic'.repeat(10000), '') - expect(() => assertSearchWalBudget(index.db, 4096)).toThrow(SearchWalBackpressureError) - reader.exec('COMMIT') - expect(() => assertSearchWalBudget(index.db, 4096)).not.toThrow() - } finally { - reader.close() - await index.close() - } -}) - -it('retains the searchable generation and retries a backpressured read once the reader releases', async () => { - const index = await openSessionSearchIndexFile('ss-wal-write') - const errors: unknown[] = [] - const store = new SessionSearchStore(index.path, (error) => errors.push(error), { - walBudgetBytes: 4096 - }) - registerSessionSearchIndexConsumer(store) - const reader = new SyncDatabase(index.path, { readonly: true }) - const hits = (term: string): number => - ( - index.db - .prepare( - `SELECT count(*) AS n FROM messages_fts JOIN visible_messages m ON m.id = messages_fts.rowid - WHERE messages_fts MATCH ?` - ) - .get(term) as { n: number } - ).n - try { - replayTranscriptRead({ messages: userMessages('oldneedle', 1), outcome: { byteOffset: 1 } }) - await store.settled() - reader.exec('BEGIN') - reader.prepare('SELECT count(*) FROM messages').get() - - replayTranscriptRead({ - mode: 'append', - previousByteOffset: 1, - messages: userMessages('newneedle', 1000), - outcome: { byteOffset: 2 } - }) - await store.settled() - expect(store.failures).toBeGreaterThan(0) - expect(store.pendingFileCount).toBe(1) - // The published generation is untouched and the cursor has not moved. - expect(hits('oldneedle')).toBe(1) - expect(hits('newneedle')).toBe(0) - expect(store.indexedFile(SYNTHETIC_TRANSCRIPT, null)?.byteOffset).toBe(1) - - reader.exec('COMMIT') - replayTranscriptRead({ - mode: 'append', - previousByteOffset: 1, - messages: userMessages('newneedle', 1000), - outcome: { byteOffset: 2 } - }) - await store.settled() - await store.purgeOlderThan(null) - expect(hits('newneedle')).toBe(1000) - expect(store.pendingFileCount).toBe(0) - expect(index.db.prepare('SELECT count(*) AS n FROM messages').get()).toEqual({ n: 1001 }) - } finally { - reader.close() - store.close() - await index.close() - } -}) diff --git a/src/main/ai-vault-search/session-search-wal-budget.ts b/src/main/ai-vault-search/session-search-wal-budget.ts deleted file mode 100644 index c2bac9771a7..00000000000 --- a/src/main/ai-vault-search/session-search-wal-budget.ts +++ /dev/null @@ -1,24 +0,0 @@ -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() - } -}