mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
fix(ai-vault-search): drain the rows a read that never published staged
Only `writePublished` and `removeFile` scheduled the cleanup lane. A read that staged rows and then declined to publish them tombstoned its batch and told nobody, so 256 rows of a 300-message incomplete read sat in `messages` and both FTS tables until some unrelated write happened to schedule a pass. Open-time recovery had the same hole: it wrote the tombstones and drained none. The abandoned branch now schedules the drain the published branch already got, and opening schedules one pass for what recovery just tombstoned. Recovery no longer adds a batch tombstone when the session-keyed one already covers those rows, which it did once more on every reopen.
This commit is contained in:
@@ -106,7 +106,7 @@ class SessionSearchReadConsumer implements TranscriptReadConsumer {
|
||||
this.store.writePublished(candidate)
|
||||
return
|
||||
}
|
||||
this.store.markStale(candidate)
|
||||
this.store.writeAbandoned(candidate)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -28,8 +28,13 @@ export function discardSearchBatch(
|
||||
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:'||id,session_row_id,id FROM search_write_batches`)
|
||||
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)`)
|
||||
}
|
||||
|
||||
@@ -255,9 +255,29 @@ it('retires an unfinished batch as part of opening, not of using', async () => {
|
||||
|
||||
const reopened = openSessionSearchDatabase(path)
|
||||
try {
|
||||
expect(reopened.prepare('SELECT count(*) AS n FROM search_pending_deletes').get()).toEqual({
|
||||
n: 2
|
||||
})
|
||||
// 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' }
|
||||
])
|
||||
} finally {
|
||||
reopened.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)
|
||||
try {
|
||||
expect(reopened.prepare('SELECT path,batch_id FROM search_pending_deletes').all()).toEqual([
|
||||
{ path: '\u0000batch:7', batch_id: 7 }
|
||||
])
|
||||
} finally {
|
||||
reopened.close()
|
||||
}
|
||||
|
||||
@@ -176,6 +176,52 @@ it('retires a batch appended onto a live session when the writer dies', async ()
|
||||
}
|
||||
})
|
||||
|
||||
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()
|
||||
|
||||
@@ -55,6 +55,10 @@ export class SessionSearchStore {
|
||||
) {
|
||||
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()
|
||||
}
|
||||
|
||||
setAcceptingWrites(accept: boolean): void {
|
||||
@@ -117,6 +121,17 @@ export class SessionSearchStore {
|
||||
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)
|
||||
@@ -263,8 +278,15 @@ export class SessionSearchStore {
|
||||
})
|
||||
}
|
||||
|
||||
/** Tests only: the cleanup lane is fire-and-forget everywhere else. */
|
||||
settled(): Promise<void> {
|
||||
return this.cleanup ?? Promise.resolve()
|
||||
/**
|
||||
* 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<void> {
|
||||
while (this.cleanup) {
|
||||
await this.cleanup
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user