Skip to content

Commit cf63591

Browse files
committed
fix(search): scale the retirement scan window with the row limit and time out resumed rechecks like completion
A capped page re-reads its scan from its last mutated row, so a fixed 25,000-ID window re-read most of the same IDs on every page once the row limit shrank. Each page now reads at most four IDs per row of its limit, capped at 25,000, and any fast page doubles the limit so sparse stretches widen the window again. A retry that finds retirement already complete revalidates every captured KB, as completion does, so it now runs under the same 30-minute timeout instead of the two-minute page timeout.
1 parent 889204f commit cf63591

3 files changed

Lines changed: 107 additions & 21 deletions

File tree

‎packages/db/script-migrations/0027_retire_search_embeddings.integration.ts‎

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -390,6 +390,84 @@ describe('retiring dormant Search embeddings', () => {
390390
}
391391
}, 60_000)
392392

393+
it('bounds the IDs each page reads by the row limit once the limit shrinks', async () => {
394+
const docs = 3000
395+
const bound = 25
396+
await sql`INSERT INTO document (id, knowledge_base_id)
397+
SELECT 'doc-' || lpad(i::text, 5, '0'), 'search' FROM generate_series(1, ${docs}) i`
398+
/** Statements over `bound` rows time out, which pins the row limit at its 25-row floor. */
399+
await sql.unsafe(`CREATE FUNCTION bound_document_update() RETURNS trigger LANGUAGE plpgsql AS $$
400+
BEGIN
401+
IF (SELECT count(*) FROM changed_rows) > ${bound} THEN
402+
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
403+
END IF;
404+
RETURN NULL;
405+
END $$`)
406+
await sql`CREATE TRIGGER bound_document_update AFTER UPDATE ON document
407+
REFERENCING NEW TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_document_update()`
408+
/**
409+
* On a table this small the planner may answer any page with a sequential scan, which reads
410+
* every row whatever the window; production pages use the primary key, so the test does too.
411+
*/
412+
await sql`SET enable_seqscan = off`
413+
/** Document rows read by any scan, counted across committed and rolled-back pages alike. */
414+
async function documentReads() {
415+
await sql`SELECT pg_stat_force_next_flush()`
416+
await admin`SELECT pg_stat_clear_snapshot()`
417+
const [row] = await admin`SELECT (seq_tup_read + coalesce(idx_tup_fetch, 0))::int AS n
418+
FROM pg_stat_user_tables WHERE schemaname = ${schema} AND relname = 'document'`
419+
return row.n
420+
}
421+
try {
422+
const before = await documentReads()
423+
expect(await pass()).toBe(true)
424+
const reads = (await documentReads()) - before
425+
/**
426+
* About 120 pages retire the 3,001 documents 25 at a time, each after one rolled-back attempt
427+
* at 50 rows. A window of four IDs per row reads about 300 IDs and 75 update lookups per page,
428+
* roughly 15 reads per document, plus the chunk phase and completion rechecks. A fixed
429+
* 25,000-ID window re-reads the rest of the table on every attempt, over 100 per document.
430+
*/
431+
expect(reads).toBeLessThan(30 * docs)
432+
expect(
433+
(
434+
await sql`SELECT count(*)::int AS n FROM document
435+
WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
436+
)[0].n
437+
).toBe(0)
438+
} finally {
439+
await sql`RESET enable_seqscan`
440+
await sql`DROP TRIGGER IF EXISTS bound_document_update ON document`
441+
await sql`DROP FUNCTION bound_document_update()`
442+
}
443+
}, 120_000)
444+
445+
it('gives the completed-retirement recheck the completion timeout when it resumes', async () => {
446+
expect(await pass()).toBe(true)
447+
await sql`DELETE FROM script_migrations`
448+
/** Stands in for a recheck that outlasts the two-minute page timeout on a large target set. */
449+
await sql`CREATE FUNCTION require_recheck_timeout() RETURNS boolean LANGUAGE plpgsql AS $$
450+
BEGIN
451+
IF current_setting('statement_timeout') <> '30min' THEN
452+
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
453+
END IF;
454+
RETURN true;
455+
END $$`
456+
await sql`ALTER TABLE search_embedding_cleanup_targets RENAME TO captured_targets`
457+
await sql`CREATE VIEW search_embedding_cleanup_targets AS
458+
SELECT knowledge_base_id FROM captured_targets WHERE require_recheck_timeout()`
459+
try {
460+
expect(await sql`SELECT phase FROM search_embedding_cleanup_progress`).toEqual([
461+
{ phase: 'done' },
462+
])
463+
expect(await pass()).toBe(true)
464+
} finally {
465+
await sql`DROP VIEW IF EXISTS search_embedding_cleanup_targets`
466+
await sql`ALTER TABLE IF EXISTS captured_targets RENAME TO search_embedding_cleanup_targets`
467+
await sql`DROP FUNCTION require_recheck_timeout()`
468+
}
469+
})
470+
393471
it('scans past a run of already-retired documents longer than the row limit in one page', async () => {
394472
/** 6,000 retired Search documents exceed the initial 2,000-row limit but fit one 25,000-ID scan. */
395473
await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled)

‎packages/db/script-migrations/0027_retire_search_embeddings.ts‎

Lines changed: 21 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,14 @@ import { sleep } from '@sim/utils/helpers'
77
import postgres, { type Sql, type TransactionSql } from 'postgres'
88

99
const logger = createLogger('RetireSearchEmbeddings')
10-
/** IDs one page reads in primary-key order; reading is cheap next to the mutation. */
10+
/** Most IDs one page reads in primary-key order; reading is cheap next to the mutation. */
1111
const SCAN_PAGE_SIZE = 25_000
12+
/**
13+
* A page reads at most this many IDs per row it may mutate. A page that reaches the row limit is
14+
* re-read from its last mutated row, so a window far wider than the limit would be read again and
15+
* again as the limit shrinks, exactly when the database is already slow.
16+
*/
17+
const SCAN_ROWS_PER_MUTATION = 4
1218
/**
1319
* Rows one page may update or delete. Every retired document is a non-HOT update touching each of
1420
* its indexes, and every deleted chunk cascades into its projections, so the write cost of a page,
@@ -17,7 +23,7 @@ const SCAN_PAGE_SIZE = 25_000
1723
const ROW_LIMIT = { initial: 2_000, min: 25, max: 8_000 } as const
1824
/** A page slower than this halves the row limit. */
1925
const SLOW_PAGE_MS = 30_000
20-
/** A page faster than this that reached the row limit doubles it. */
26+
/** A page faster than this doubles the row limit, widening its scan window with it. */
2127
const FAST_PAGE_MS = SLOW_PAGE_MS / 4
2228
/** Each page is followed by a pause as long as the page, up to this, to leave the primary headroom. */
2329
const MAX_PAGE_PAUSE_MS = 5_000
@@ -33,8 +39,6 @@ interface PageResult {
3339
afterId: string
3440
/** Rows the page updated or deleted. */
3541
mutated: number
36-
/** The page stopped at the row limit rather than the end of its scan. */
37-
limited: boolean
3842
/** A phase change the page committed. */
3943
transition?: 'embeddings' | 'documents_rescan' | 'embeddings_rescan'
4044
}
@@ -162,7 +166,7 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = {
162166
pageMs: Math.round(pageMs),
163167
rowLimit,
164168
})
165-
} else if (page.limited && pageMs < FAST_PAGE_MS) {
169+
} else if (pageMs < FAST_PAGE_MS) {
166170
rowLimit = Math.min(ROW_LIMIT.max, rowLimit * 2)
167171
}
168172
if (batches % 10 === 0) {
@@ -186,14 +190,16 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = {
186190
}
187191

188192
/**
189-
* Retires one page: the target rows among the next `SCAN_PAGE_SIZE` IDs, capped at `rowLimit`.
193+
* Retires one page: the target rows among the next `rowLimit * SCAN_ROWS_PER_MUTATION` IDs (at
194+
* most `SCAN_PAGE_SIZE`), capped at `rowLimit`.
190195
* A capped page advances the cursor only to its last mutated row, so the rest of the scan is
191196
* read again by the next page; an uncapped page advances past its whole scan.
192197
*/
193198
async function retirePage(sql: Sql, rowLimit: number): Promise<PageResult> {
194199
return retryOnLockTimeout(
195200
() =>
196201
sql.begin(async (tx) => {
202+
const scanLimit = Math.min(SCAN_PAGE_SIZE, rowLimit * SCAN_ROWS_PER_MUTATION)
197203
await tx`SET LOCAL statement_timeout = '120s'`
198204
await tx`SET LOCAL lock_timeout = '1s'`
199205
const [progress] = await tx<Progress[]>`
@@ -203,24 +209,25 @@ async function retirePage(sql: Sql, rowLimit: number): Promise<PageResult> {
203209
phase: progress.phase,
204210
afterId: progress.after_id,
205211
mutated: 0,
206-
limited: false,
207212
...page,
208213
})
209214
if (progress.phase === 'done') {
215+
/** A retry after maintenance failed rechecks every captured KB, as completion did. */
216+
await tx`SET LOCAL statement_timeout = '30min'`
210217
await validateTargetMarkers(tx)
211218
return result({ done: true })
212219
}
213220

214221
if (progress.phase === 'documents') {
215222
/** Already-retired documents are skipped so they never spend the row limit. */
216223
const [page] = await pageMutation(tx<
217-
{ after_id: string; limited: boolean; mutated: number; invalid_target: boolean }[]
224+
{ after_id: string; mutated: number; invalid_target: boolean }[]
218225
>`
219226
WITH source_page AS MATERIALIZED (
220227
SELECT id, knowledge_base_id,
221228
(NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL
222229
OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL) AS unretired
223-
FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${SCAN_PAGE_SIZE}
230+
FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${scanLimit}
224231
), target_page AS MATERIALIZED (
225232
SELECT p.id, p.knowledge_base_id FROM source_page p
226233
JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = p.knowledge_base_id
@@ -246,7 +253,7 @@ async function retirePage(sql: Sql, rowLimit: number): Promise<PageResult> {
246253
OR d.processing_queued_at IS NOT NULL OR d.processing_deferred_until IS NOT NULL)
247254
RETURNING d.id
248255
) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id,
249-
e.limited, (SELECT count(*) FROM retired)::int AS mutated,
256+
(SELECT count(*) FROM retired)::int AS mutated,
250257
EXISTS (SELECT 1 FROM invalid_target) AS invalid_target
251258
FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target`)
252259
if (!page) {
@@ -256,21 +263,20 @@ async function retirePage(sql: Sql, rowLimit: number): Promise<PageResult> {
256263
if (page.invalid_target)
257264
throw new Error('Cleanup target is no longer a Search knowledge base')
258265
await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1`
259-
return result({ afterId: page.after_id, mutated: page.mutated, limited: page.limited })
266+
return result({ afterId: page.after_id, mutated: page.mutated })
260267
}
261268

262269
/** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */
263270
const [page] = await pageMutation(tx<
264271
{
265272
after_id: string
266-
limited: boolean
267273
mutated: number
268274
unretired: boolean
269275
invalid_target: boolean
270276
}[]
271277
>`
272278
WITH source_page AS MATERIALIZED (
273-
SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${SCAN_PAGE_SIZE}
279+
SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${scanLimit}
274280
), target_page AS MATERIALIZED (
275281
SELECT p.* FROM source_page p
276282
JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = p.knowledge_base_id
@@ -295,7 +301,7 @@ async function retirePage(sql: Sql, rowLimit: number): Promise<PageResult> {
295301
AND NOT EXISTS (SELECT 1 FROM unretired) AND NOT EXISTS (SELECT 1 FROM invalid_target)
296302
RETURNING e.id
297303
) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id,
298-
e.limited, (SELECT count(*) FROM deleted)::int AS mutated,
304+
(SELECT count(*) FROM deleted)::int AS mutated,
299305
EXISTS (SELECT 1 FROM unretired) AS unretired,
300306
EXISTS (SELECT 1 FROM invalid_target) AS invalid_target
301307
FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target`)
@@ -329,7 +335,7 @@ async function retirePage(sql: Sql, rowLimit: number): Promise<PageResult> {
329335
return result({ done: true })
330336
}
331337
await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1`
332-
return result({ afterId: page.after_id, mutated: page.mutated, limited: page.limited })
338+
return result({ afterId: page.after_id, mutated: page.mutated })
333339
}),
334340
{
335341
budgetMs: LOCK_RETRY_BUDGET_MS,

‎packages/db/script-migrations/search-embedding-retirement.md‎

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -30,18 +30,20 @@ resume the saved scope, phase, cursor and maintenance checkpoints.
3030
The existing maintenance implementation rebuilds HNSW indexes and vacuums affected tables before
3131
deployment continues.
3232

33-
Each page reads at most 25,000 IDs and mutates at most a row limit of them. Pages execute
33+
Each page mutates at most a row limit of target rows and reads at most four IDs per row of that
34+
limit, never more than 25,000 IDs. Pages execute
3435
sequentially, and each is followed by a pause as long as the page took, up to five seconds, to
3536
leave the primary headroom. Retiring a document is a non-HOT update that writes every index on
3637
`document`, and deleting a chunk cascades into its projections, so a page's cost follows the target
3738
rows it mutates, not the IDs it reads. A page that reaches the row limit advances the cursor only to
38-
its last mutated row; the rest of its scan is read again by the next page. Documents that are
39-
already retired never count against the limit.
39+
its last mutated row; the rest of its scan is read again by the next page. Tying the scan window to
40+
the limit keeps that re-reading proportional to the work, even after the limit shrinks. Documents
41+
that are already retired never count against the limit.
4042

4143
The row limit starts at 2,000 rows. A page is timed from the start of its transaction through its
4244
commit, including the synchronous-replication wait and any lock-timeout retries. A page slower than
43-
30 seconds halves the limit. A fast page, one under 7.5 seconds, that reached the limit doubles it,
44-
up to 8,000. The limit never drops below 25 rows. Phase changes do not adjust it.
45+
30 seconds halves the limit. A fast page, one under 7.5 seconds, doubles it up to 8,000, which
46+
also widens the scan window, so sparse stretches are not crawled in small windows. The limit never drops below 25 rows. Phase changes do not adjust it.
4547

4648
Materialized SQL pages keep the IDs inside PostgreSQL; the migration process receives only a cursor
4749
and a validation result. Each page uses a two-minute statement timeout and a one-second lock
@@ -88,7 +90,7 @@ unrelated rows. Before completion, the cleanup checks for unretired documents an
8890
behind either cursor and restarts the affected phase if needed. A final bounded pass validates all
8991
captured KB markers, including empty KBs and KBs whose rows were already scanned, holding shared
9092
marker locks until the completion checkpoint commits. Resuming a completed cleanup before maintenance
91-
also revalidates the captured set. Keep target writers stopped and
93+
also revalidates the captured set, with the same 30-minute timeout as the completion recheck. Keep target writers stopped and
9294
do not change their Search markers during the pass.
9395

9496
Inspect progress with:

0 commit comments

Comments
 (0)