Skip to content

Commit decf822

Browse files
committed
improvement(knowledge): key-share lock the chunks a search-index adoption projects
1 parent e43c5e0 commit decf822

3 files changed

Lines changed: 88 additions & 23 deletions

File tree

‎packages/db/script-migrations/0019_tin_keyword_projection.ts‎

Lines changed: 30 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,35 @@ export async function installTinChunkSync(tx: Sql | TransactionSql): Promise<voi
8080
$$`)
8181
}
8282

83+
/**
84+
* The knowledge base trigger's body: projects or removes a whole base when its search-index marker
85+
* changes. The chunks are key-share locked as they are read, as the backfill locks its own, so a
86+
* chunk delete in flight is waited for and its chunk skipped rather than failing the adoption on
87+
* the row's foreign key.
88+
*/
89+
export async function installTinMembershipSync(tx: Sql | TransactionSql): Promise<void> {
90+
await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_knowledge_base_keyword_tin()
91+
RETURNS trigger LANGUAGE plpgsql AS $$
92+
BEGIN
93+
PERFORM pg_advisory_xact_lock(knowledge_tin_membership_key(NEW.id));
94+
IF NOT NEW.is_search_index THEN
95+
DELETE FROM embedding_keyword_tin t USING embedding e
96+
WHERE e.knowledge_base_id = NEW.id AND t.id = e.id;
97+
RETURN NEW;
98+
END IF;
99+
INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content)
100+
SELECT id, knowledge_base_id, document_id, enabled,
101+
knowledge_tin_base_token(knowledge_base_id) || ' ' || knowledge_tin_stream(content_tsv)
102+
FROM embedding WHERE knowledge_base_id = NEW.id
103+
ORDER BY id FOR KEY SHARE
104+
ON CONFLICT (id) DO UPDATE SET
105+
knowledge_base_id = EXCLUDED.knowledge_base_id, document_id = EXCLUDED.document_id,
106+
enabled = EXCLUDED.enabled, content = EXCLUDED.content;
107+
RETURN NEW;
108+
END;
109+
$$`)
110+
}
111+
83112
/**
84113
* Installs the stream and base-token functions shared by the triggers and the query path, and the
85114
* embedding and knowledge base triggers, atomically with respect to embedding writers.
@@ -123,25 +152,7 @@ export async function installProjection(sql: Sql): Promise<void> {
123152
await tx.unsafe(`CREATE OR REPLACE TRIGGER embedding_keyword_tin_sync
124153
AFTER INSERT OR UPDATE OF knowledge_base_id, document_id, enabled, content ON embedding
125154
FOR EACH ROW WHEN (${SYNCHRONOUS_PROJECTION_WHEN}) EXECUTE FUNCTION sync_embedding_keyword_tin()`)
126-
await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_knowledge_base_keyword_tin()
127-
RETURNS trigger LANGUAGE plpgsql AS $$
128-
BEGIN
129-
PERFORM pg_advisory_xact_lock(knowledge_tin_membership_key(NEW.id));
130-
IF NOT NEW.is_search_index THEN
131-
DELETE FROM embedding_keyword_tin t USING embedding e
132-
WHERE e.knowledge_base_id = NEW.id AND t.id = e.id;
133-
RETURN NEW;
134-
END IF;
135-
INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content)
136-
SELECT id, knowledge_base_id, document_id, enabled,
137-
knowledge_tin_base_token(knowledge_base_id) || ' ' || knowledge_tin_stream(content_tsv)
138-
FROM embedding WHERE knowledge_base_id = NEW.id
139-
ON CONFLICT (id) DO UPDATE SET
140-
knowledge_base_id = EXCLUDED.knowledge_base_id, document_id = EXCLUDED.document_id,
141-
enabled = EXCLUDED.enabled, content = EXCLUDED.content;
142-
RETURN NEW;
143-
END;
144-
$$`)
155+
await installTinMembershipSync(tx)
145156
await tx.unsafe(`CREATE OR REPLACE TRIGGER knowledge_base_keyword_tin_sync
146157
AFTER UPDATE OF is_search_index ON knowledge_base
147158
FOR EACH ROW WHEN (OLD.is_search_index IS DISTINCT FROM NEW.is_search_index)

‎packages/db/script-migrations/0025_scope_keyword_projections.integration.ts‎

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -241,6 +241,53 @@ describe('keyword projections scoped to search indexes in PostgreSQL', () => {
241241
expect(await keywordIds()).toEqual(['waiting'])
242242
})
243243

244+
/**
245+
* Adopts `legacy` while a second connection holds a delete of one of its chunks open, and resolves
246+
* once both have committed.
247+
*/
248+
async function adoptWhileDeleting(): Promise<void> {
249+
await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) VALUES
250+
('kept', 'legacy', 'doc', to_tsvector('english', 'Kept chunk')),
251+
('deleted', 'legacy', 'doc', to_tsvector('english', 'Deleted chunk'))`
252+
let deleted!: () => void
253+
const deletedSignal = new Promise<void>((resolve) => {
254+
deleted = resolve
255+
})
256+
let release!: () => void
257+
const released = new Promise<void>((resolve) => {
258+
release = resolve
259+
})
260+
const deleter = sql.begin(async (tx) => {
261+
await tx`DELETE FROM embedding WHERE id = 'deleted'`
262+
deleted()
263+
await released
264+
})
265+
await deletedSignal
266+
const [{ pid }] = await other`SELECT pg_backend_pid() AS pid`
267+
const adoption =
268+
other`UPDATE knowledge_base SET is_search_index = true WHERE id = 'legacy'`.execute()
269+
await waitUntilBlocked(admin, pid)
270+
release()
271+
await deleter
272+
await adoption
273+
}
274+
275+
it('adopts a base while one of its chunks is being deleted', async () => {
276+
await adoptWhileDeleting()
277+
expect(await keywordIds()).toEqual(['kept'])
278+
expect(await tinIds()).toEqual(['kept'])
279+
})
280+
281+
it('adopts a base into Tin while a chunk is being deleted, with no keyword trigger ahead of it', async () => {
282+
await sql`ALTER TABLE knowledge_base DISABLE TRIGGER knowledge_base_keyword_search_sync`
283+
try {
284+
await adoptWhileDeleting()
285+
} finally {
286+
await sql`ALTER TABLE knowledge_base ENABLE TRIGGER knowledge_base_keyword_search_sync`
287+
}
288+
expect(await tinIds()).toEqual(['kept'])
289+
})
290+
244291
it('runs no fan-out for an inserted document, and still fans out a changed ACL', async () => {
245292
const projections = ['embedding_search', 'embedding_keyword_tin'] as const
246293
const scans = await sql.begin(async (tx) => {

‎packages/db/script-migrations/0025_scope_keyword_projections.ts‎

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import {
22
installMembershipKey,
33
installTinChunkSync,
4+
installTinMembershipSync,
45
} from '@sim/db/script-migrations/0019_tin_keyword_projection'
56
import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url'
67
import type { ScriptMigration } from '@sim/db/script-migrations/types'
@@ -63,7 +64,9 @@ async function installKeywordSearchSync(tx: TransactionSql): Promise<void> {
6364
/**
6465
* Projects or removes a whole base's GIN keyword rows when its search-index marker changes, under
6566
* the membership lock taken exclusively, as the Tin projection's knowledge base trigger does. It is
66-
* its own trigger because that one exists only where Tin is installed.
67+
* its own trigger because that one exists only where Tin is installed. The chunks are key-share
68+
* locked as they are read, as the projection backfills lock theirs: a chunk delete in flight is
69+
* waited for and its chunk skipped, rather than failing the adoption on the row's foreign key.
6770
*/
6871
async function installKeywordSearchMembership(tx: TransactionSql): Promise<void> {
6972
await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_knowledge_base_keyword_search()
@@ -77,6 +80,7 @@ async function installKeywordSearchMembership(tx: TransactionSql): Promise<void>
7780
INSERT INTO embedding_keyword_search AS s (id, knowledge_base_id, document_id, enabled, content_tsv)
7881
SELECT id, knowledge_base_id, document_id, enabled, content_tsv
7982
FROM embedding WHERE knowledge_base_id = NEW.id
83+
ORDER BY id FOR KEY SHARE
8084
${KEYWORD_SEARCH_UPSERT};
8185
RETURN NEW;
8286
END;
@@ -88,13 +92,16 @@ async function installKeywordSearchMembership(tx: TransactionSql): Promise<void>
8892
}
8993

9094
/**
91-
* The Tin chunk trigger's body as `0019_tin_keyword_projection` now installs it, where that
92-
* migration installed it: without the delete an insert outside a search index ran.
95+
* The Tin trigger bodies as `0019_tin_keyword_projection` now installs them, where that migration
96+
* installed them: the chunk trigger without the delete an insert outside a search index ran, and
97+
* the knowledge base trigger key-share locking the chunks it projects.
9398
*/
9499
async function installTinSync(tx: TransactionSql): Promise<void> {
95100
const [row] = await tx<Array<{ installed: boolean }>>`
96101
SELECT to_regprocedure('sync_embedding_keyword_tin()') IS NOT NULL AS installed`
97-
if (row?.installed) await installTinChunkSync(tx)
102+
if (!row?.installed) return
103+
await installTinChunkSync(tx)
104+
await installTinMembershipSync(tx)
98105
}
99106

100107
/**

0 commit comments

Comments
 (0)