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() - } -}