mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
fix(ai-vault-search): hide a tombstoned session's messages, not only its row
`visible_sessions` already subtracted session-keyed tombstones; the message half filtered on the batch pointer alone. A published row outlives its session row until the cleanup lane reaches it, so between a `replace` publish and that drain both generations answered, and a removed file kept answering after its session was gone. Both views now subtract the same set, over a partial index so neither read scans the tombstone table.
This commit is contained in:
@@ -272,10 +272,13 @@ describe('visibility views', () => {
|
||||
(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')`)
|
||||
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 }
|
||||
])
|
||||
|
||||
@@ -67,6 +67,10 @@ CREATE TABLE IF NOT EXISTS search_pending_deletes(
|
||||
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,
|
||||
@@ -91,11 +95,17 @@ CREATE VIRTUAL TABLE IF NOT EXISTS conversation_fts USING fts5(
|
||||
);
|
||||
-- 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;
|
||||
WHERE batch_id IS NULL
|
||||
AND session_row_id NOT IN
|
||||
(SELECT session_row_id FROM search_pending_deletes WHERE batch_id IS NULL);
|
||||
`
|
||||
|
||||
/**
|
||||
|
||||
@@ -31,6 +31,22 @@ afterEach(async () => {
|
||||
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<string, number> {
|
||||
const one = (sql: string): number => (db.prepare(sql).get() as { n: number }).n
|
||||
return {
|
||||
@@ -69,18 +85,9 @@ it('writes both FTS tables for every published conversational row', async () =>
|
||||
// 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)
|
||||
const matches = (table: string, term: string): number =>
|
||||
(
|
||||
index.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
|
||||
expect(matches('messages_fts', 'gamma')).toBe(1)
|
||||
expect(matches('conversation_fts', 'gamma')).toBe(0)
|
||||
expect(matches('conversation_fts', 'beta')).toBe(1)
|
||||
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 () => {
|
||||
@@ -178,15 +185,31 @@ it('replaces the previous generation without ever showing both', async () => {
|
||||
|
||||
expect(counts(index.db).sessions).toBe(1)
|
||||
expect(counts(index.db).messages).toBe(10)
|
||||
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
|
||||
expect(hits('firstgeneration')).toBe(0)
|
||||
expect(hits('secondgeneration')).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()
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user