Skip to content

Commit 3ec5f2e

Browse files
committed
fix(knowledge): serialize the projection source with its document
- the set trigger reads the document under FOR SHARE, so a projection write and a document's connector change cannot interleave and leave the older value - the backfill writes a chunk's source only for the document it read that source from; a chunk moved meanwhile is left to its new document's trigger
1 parent 886c5f7 commit 3ec5f2e

1 file changed

Lines changed: 8 additions & 4 deletions

File tree

‎packages/db/script-migrations/0021_embedding_search_connector.ts‎

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,8 @@ export async function installEmbeddingSearchConnector(sql: Sql): Promise<void> {
3434
await tx.unsafe(`CREATE OR REPLACE FUNCTION set_embedding_search_connector()
3535
RETURNS trigger LANGUAGE plpgsql AS $$
3636
BEGIN
37-
SELECT connector_id INTO NEW.connector_id FROM document WHERE id = NEW.document_id;
37+
SELECT connector_id INTO NEW.connector_id FROM document
38+
WHERE id = NEW.document_id FOR SHARE;
3839
RETURN NEW;
3940
END;
4041
$$`)
@@ -65,18 +66,21 @@ export async function backfillEmbeddingSearchConnector(sql: Sql): Promise<number
6566
/**
6667
* The documents are share-locked before their source is copied, so a detachment in flight
6768
* waits for this page to commit and then fans its own change out through the trigger — the
68-
* chunk can never keep a source its document no longer has.
69+
* chunk can never keep a source its document no longer has. A chunk that moved to another
70+
* document in the meantime is left to that document's trigger: the write requires the
71+
* document the source was read for.
6972
*/
7073
const [{ filled }] = await tx<Array<{ filled: number }>>`
7174
WITH page AS (
72-
SELECT s.id, d.connector_id
75+
SELECT s.id, s.document_id, d.connector_id
7376
FROM embedding_search s JOIN document d ON d.id = s.document_id
7477
WHERE s.id = ANY(${ids}::text[])
7578
AND s.connector_id IS NULL AND d.connector_id IS NOT NULL
7679
FOR SHARE OF d
7780
), updated AS (
7881
UPDATE embedding_search s SET connector_id = page.connector_id
79-
FROM page WHERE s.id = page.id AND s.connector_id IS NULL
82+
FROM page
83+
WHERE s.id = page.id AND s.document_id = page.document_id AND s.connector_id IS NULL
8084
RETURNING s.id
8185
) SELECT count(*)::int AS filled FROM updated`
8286
return { afterId: ids[ids.length - 1], scanned: ids.length, filled }

0 commit comments

Comments
 (0)