Skip to content

Commit 4119ea6

Browse files
committed
fix(knowledge): serialize the Tin backfill with search-index membership changes
- Each backfill page takes its bases' membership locks, then reads under a fresh snapshot - Remove a base's projection through the embedding knowledge base index instead of scanning the projection
1 parent 48d9d83 commit 4119ea6

2 files changed

Lines changed: 82 additions & 20 deletions

File tree

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

Lines changed: 52 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,9 @@
11
import { readFileSync } from 'node:fs'
22
import path from 'node:path'
3-
import { installProjection } from '@sim/db/script-migrations/0019_tin_keyword_projection'
3+
import {
4+
backfillProjection,
5+
installProjection,
6+
} from '@sim/db/script-migrations/0019_tin_keyword_projection'
47
import { generateId } from '@sim/utils/id'
58
import postgres, { type Sql } from 'postgres'
69
import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest'
@@ -58,6 +61,7 @@ describe.runIf(Boolean(databaseUrl))('Tin keyword projection in PostgreSQL', ()
5861
document_id text NOT NULL, enabled boolean NOT NULL DEFAULT true, content text,
5962
content_tsv tsvector NOT NULL
6063
)`
64+
await sql`CREATE INDEX ON embedding (knowledge_base_id)`
6165
/** The shipped table, with its foreign key pointed at this schema's `embedding`. */
6266
const [createTable] = readFileSync(
6367
path.join(MIGRATIONS, '0366_tin_keyword_projection.sql'),
@@ -161,4 +165,51 @@ describe.runIf(Boolean(databaseUrl))('Tin keyword projection in PostgreSQL', ()
161165
await update
162166
expect(await projected()).toMatchObject([{ id: 'toggled', enabled: false }])
163167
})
168+
it('removes a base through the chunk index instead of scanning the projection', async () => {
169+
await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv)
170+
SELECT 'index-' || n, 'index', 'doc', to_tsvector('english', 'Chunk ' || n)
171+
FROM generate_series(1, 20) n`
172+
/** Pending per-backend counts include earlier statements, so compare before and after. */
173+
const scans = await sql.begin(async (tx) => {
174+
const read = async () => {
175+
const [row] = await tx<Array<{ scans: number }>>`
176+
SELECT seq_scan::int AS scans FROM pg_stat_xact_user_tables
177+
WHERE relid = 'embedding_keyword_tin'::regclass`
178+
return row.scans
179+
}
180+
await tx`SET LOCAL enable_seqscan = off`
181+
const before = await read()
182+
await tx`UPDATE knowledge_base SET is_search_index = false WHERE id = 'index'`
183+
return (await read()) - before
184+
})
185+
expect(scans).toBe(0)
186+
expect(await projected()).toEqual([])
187+
})
188+
189+
it('does not backfill a chunk whose base stops being a search index mid-page', async () => {
190+
await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv)
191+
VALUES ('unprojected', 'index', 'doc', to_tsvector('english', 'Predates the trigger'))`
192+
await sql`DELETE FROM embedding_keyword_tin`
193+
let demoted!: () => void
194+
const demotedSignal = new Promise<void>((resolve) => {
195+
demoted = resolve
196+
})
197+
let release!: () => void
198+
const released = new Promise<void>((resolve) => {
199+
release = resolve
200+
})
201+
const demotion = promoter.begin(async (tx) => {
202+
await tx`UPDATE knowledge_base SET is_search_index = false WHERE id = 'index'`
203+
demoted()
204+
await released
205+
})
206+
await demotedSignal
207+
const [{ pid }] = await sql`SELECT pg_backend_pid() AS pid`
208+
const backfill = backfillProjection(sql)
209+
await waitUntilBlocked(admin, pid)
210+
release()
211+
await demotion
212+
expect(await backfill).toBe(0)
213+
expect(await projected()).toEqual([])
214+
})
164215
})

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

Lines changed: 30 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,9 @@ async function tinAvailable(sql: Sql): Promise<boolean> {
5656
* each other, and a marker change takes it exclusively, so it waits for writers already in flight
5757
* and later writers wait for it. Each then reads under a fresh snapshot, so neither misses the
5858
* other's rows. A row lock cannot do this: a key-share lock on the marker row is compatible with
59-
* the uncommitted marker update, so a writer would read the old marker without waiting.
59+
* the uncommitted marker update, so a writer would read the old marker without waiting. A whole
60+
* base is reached through `embedding`'s knowledge base index, since a projected row always carries
61+
* its chunk's base and the projection keeps no index but Tin's.
6062
*/
6163
export async function installProjection(sql: Sql): Promise<void> {
6264
await sql.begin(async (tx) => {
@@ -102,7 +104,8 @@ export async function installProjection(sql: Sql): Promise<void> {
102104
BEGIN
103105
PERFORM pg_advisory_xact_lock(knowledge_tin_membership_key(NEW.id));
104106
IF NOT NEW.is_search_index THEN
105-
DELETE FROM embedding_keyword_tin WHERE knowledge_base_id = NEW.id;
107+
DELETE FROM embedding_keyword_tin t USING embedding e
108+
WHERE e.knowledge_base_id = NEW.id AND t.id = e.id;
106109
RETURN NEW;
107110
END IF;
108111
INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content)
@@ -122,40 +125,48 @@ export async function installProjection(sql: Sql): Promise<void> {
122125
})
123126
}
124127

125-
/** Fills rows the trigger has not written yet, in independently committed keyset pages. */
126-
async function backfillProjection(sql: Sql): Promise<number> {
128+
/**
129+
* Fills rows the trigger has not written yet, in independently committed keyset pages. Each page
130+
* takes the membership locks of its bases before reading them, as the embedding trigger does, and
131+
* then reads under a fresh snapshot: a base whose search-index marker is changing is read only
132+
* after that change commits, so the backfill never writes back a chunk the change removed.
133+
*/
134+
export async function backfillProjection(sql: Sql): Promise<number> {
127135
let afterId = ''
128136
let inserted = 0
129137
let scanned = 0
130138
const startedAt = Date.now()
131139
for (;;) {
132-
const [page] = await sql.begin(async (tx) => {
140+
const page = await sql.begin(async (tx) => {
133141
await tx.unsafe("SET LOCAL lock_timeout = '5s'")
134142
await tx.unsafe("SET LOCAL statement_timeout = '60s'")
135-
return tx.unsafe<Array<{ after_id: string | null; scanned: number; inserted: number }>>(
136-
`
137-
WITH source_page AS MATERIALIZED (
138-
SELECT id FROM embedding WHERE id > $1 ORDER BY id LIMIT ${BATCH_SIZE}
139-
), batch AS MATERIALIZED (
143+
const rows = await tx<Array<{ id: string; knowledge_base_id: string }>>`
144+
SELECT id, knowledge_base_id FROM embedding WHERE id > ${afterId} ORDER BY id
145+
LIMIT ${BATCH_SIZE}`
146+
if (rows.length === 0) return null
147+
const ids = rows.map((row) => row.id)
148+
const bases = [...new Set(rows.map((row) => row.knowledge_base_id))]
149+
await tx`SELECT pg_advisory_xact_lock_shared(knowledge_tin_membership_key(base))
150+
FROM unnest(${bases}::text[]) AS base ORDER BY base`
151+
const [{ written }] = await tx<Array<{ written: number }>>`
152+
WITH batch AS MATERIALIZED (
140153
SELECT e.id, e.knowledge_base_id, e.document_id, e.enabled, e.content_tsv
141-
FROM source_page p
142-
INNER JOIN embedding e ON e.id = p.id
154+
FROM embedding e
143155
INNER JOIN knowledge_base k ON k.id = e.knowledge_base_id AND k.is_search_index
144-
WHERE NOT EXISTS (SELECT 1 FROM embedding_keyword_tin t WHERE t.id = p.id)
156+
WHERE e.id = ANY(${ids}::text[])
157+
AND NOT EXISTS (SELECT 1 FROM embedding_keyword_tin t WHERE t.id = e.id)
145158
ORDER BY e.id FOR KEY SHARE OF e
146159
), written AS (
147160
INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content)
148161
SELECT id, knowledge_base_id, document_id, enabled,
149162
knowledge_tin_base_token(knowledge_base_id) || ' ' || knowledge_tin_stream(content_tsv)
150163
FROM batch
151164
ON CONFLICT (id) DO NOTHING RETURNING id
152-
) SELECT max(id) AS after_id, count(*)::int AS scanned,
153-
(SELECT count(*)::int FROM written) AS inserted FROM source_page`,
154-
[afterId]
155-
)
165+
) SELECT count(*)::int AS written FROM written`
166+
return { afterId: ids[ids.length - 1], scanned: ids.length, inserted: written }
156167
})
157-
if (!page.after_id) break
158-
afterId = page.after_id
168+
if (!page) break
169+
afterId = page.afterId
159170
inserted += page.inserted
160171
scanned += page.scanned
161172
if (scanned % (BATCH_SIZE * 100) === 0) {

0 commit comments

Comments
 (0)