Skip to content

Commit 294c505

Browse files
authored
fix(search): bound Search retirement pages so the migration finishes under its statement timeout (#8460)
* fix(search): bound retirement pages by mutated rows so they finish under the statement timeout A retirement page read 25,000 IDs and updated or deleted every target row among them in one statement. Retiring a document is a non-HOT update that writes every index on `document`, and a deleted chunk cascades into its projections, so on a KB that dominates the table a page's write cost, not its scan, outran the two-minute statement timeout and failed the deploy migration. Each page now mutates at most a row limit of its target rows. A page that reaches the limit advances the cursor only to its last mutated row, and already-retired documents never spend the limit. The limit starts at 2,000, halves after a slow page or a statement timeout (the timed-out page rolls back with its cursor and is retried), and doubles after a fast full page. The completion rechecks, which walk every captured KB once, run with a 30-minute timeout. The retirement stays idempotent and resumes from its saved cursor. * fix(search): shrink only timed-out page mutations and pace Search retirement pages Only a statement timeout from a page's mutating statement now halves the row limit and retries the rolled-back page. Any other timeout, such as a completion recheck, fails the run at once instead of repeating the same statement at every smaller limit. Pages are timed around the whole call, commit included, so the synchronous-replication wait counts toward the slow-page threshold. Each page is followed by a pause as long as the page, up to five seconds, and the row limit is capped at 8,000. Phase changes no longer adjust the limit. The progress log now carries the phase, cursor and rows mutated, and the migration logs slow-page halvings, phase changes and the start of the completion recheck. * test(search): prove retired documents never spend the retirement row limit A run of already-retired Search documents longer than the row limit must be crossed in one page. The test counts documents-phase statements and fails if the page filter on unretired rows is removed. * 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. * fix(search): never grow a Search retirement page back to a size that timed out, and pause after a timeout A timed-out page halved the row limit, but one fast page doubled it straight back, so the run alternated between the size that timed out and half of it, rolling back a full statement-timeout page each time. The limit now grows only up to half of the smallest size that timed out, and a timed-out page is followed by the same pause as any other page.
1 parent cebe8a3 commit 294c505

3 files changed

Lines changed: 451 additions & 40 deletions

File tree

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

Lines changed: 235 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -257,19 +257,31 @@ describe('retiring dormant Search embeddings', () => {
257257
await sql`INSERT INTO embedding VALUES ('26003', 'second-search', 'second-search-doc')`
258258
await sql`CREATE TABLE deletion_blocker (id text REFERENCES embedding(id))`
259259
await sql`INSERT INTO deletion_blocker VALUES ('26002')`
260+
/**
261+
* Pages split by mutated rows under an adaptive limit, so how far each failed run gets depends
262+
* on page geometry. The geometry-free invariant: every committed delete sits at or behind the
263+
* cursor, and a failed page leaves every row past it (IDs 00001-26003 are contiguous).
264+
*/
265+
async function expectRolledBackPastCursor() {
266+
const [progress] = await sql`SELECT phase, after_id FROM search_embedding_cleanup_progress`
267+
expect(progress.phase).toBe('embeddings')
268+
const [beyond] =
269+
await sql`SELECT count(*)::int AS n FROM embedding WHERE id > ${progress.after_id}`
270+
expect(beyond.n).toBe(26003 - Number(progress.after_id))
271+
expect(progress.after_id < '26002').toBe(true)
272+
}
260273
await expect(pass()).rejects.toThrow()
261-
const before = await sql`SELECT * FROM search_embedding_cleanup_progress`
262274
const [remaining] = await sql`SELECT count(*)::int AS n FROM embedding`
263275
expect(remaining.n).toBeGreaterThan(501)
264276
expect(remaining.n).toBeLessThan(26002)
265277
expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0)
278+
await expectRolledBackPastCursor()
266279
await sql`UPDATE knowledge_base SET is_search_index = false WHERE id = 'second-search'`
267280
await expect(pass()).rejects.toThrow('no longer a Search knowledge base')
268-
expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(remaining.n)
281+
await expectRolledBackPastCursor()
269282
await sql`UPDATE knowledge_base SET is_search_index = true WHERE id = 'second-search'`
270283
await expect(pass()).rejects.toThrow()
271-
expect(await sql`SELECT * FROM search_embedding_cleanup_progress`).toEqual(before)
272-
expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(remaining.n)
284+
await expectRolledBackPastCursor()
273285
await sql`DROP TABLE deletion_blocker`
274286
await sql`INSERT INTO document (id, knowledge_base_id) VALUES ('aaa-late-document', 'search')`
275287
await sql`INSERT INTO embedding VALUES ('00000', 'search', 'aaa-late-document')`
@@ -309,6 +321,225 @@ describe('retiring dormant Search embeddings', () => {
309321
).toBe(0)
310322
}, 60_000)
311323

324+
it('splits pages whose writes would outrun the statement timeout and resumes at the split', async () => {
325+
const bound = 300
326+
await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled)
327+
SELECT 'doc-' || lpad(i::text, 5, '0'), CASE WHEN i % 3 = 0 THEN 'ordinary' ELSE 'search' END,
328+
i % 7 = 0, i % 7 <> 0
329+
FROM generate_series(1, 4000) i`
330+
await sql`INSERT INTO embedding
331+
SELECT lpad(i::text, 5, '0'), CASE WHEN i % 3 = 0 THEN 'ordinary' ELSE 'search' END,
332+
CASE WHEN i % 3 = 0 THEN 'ordinary-doc' ELSE 'search-doc' END
333+
FROM generate_series(1003, 5002) i`
334+
await sql`CREATE TABLE committed_statement (rows integer NOT NULL)`
335+
/** Sequences are not transactional, so this counts the timed-out statements that rolled back. */
336+
await sql`CREATE SEQUENCE timed_out_statement`
337+
/** Stands in for write cost: a statement touching more than `bound` rows times out and rolls back. */
338+
await sql.unsafe(`CREATE FUNCTION bound_statement_rows() RETURNS trigger LANGUAGE plpgsql AS $$
339+
DECLARE touched integer;
340+
BEGIN
341+
SELECT count(*) INTO touched FROM changed_rows;
342+
IF touched > ${bound} THEN
343+
PERFORM nextval('timed_out_statement');
344+
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
345+
END IF;
346+
INSERT INTO committed_statement VALUES (touched);
347+
RETURN NULL;
348+
END $$`)
349+
await sql`CREATE TRIGGER bound_document_update AFTER UPDATE ON document
350+
REFERENCING NEW TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_statement_rows()`
351+
await sql`CREATE TRIGGER bound_embedding_delete AFTER DELETE ON embedding
352+
REFERENCING OLD TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_statement_rows()`
353+
try {
354+
expect(await pass()).toBe(true)
355+
const [{ largest, total }] = await sql`SELECT max(rows)::int AS largest,
356+
sum(rows)::int AS total FROM committed_statement`
357+
expect(largest).toBeLessThanOrEqual(bound)
358+
expect(largest).toBeGreaterThan(0)
359+
/**
360+
* The limit halves 2,000 → 1,000 → 500 → 250 on the first document page and never grows back
361+
* to a size that timed out, in either phase: exactly three rolled-back statements. A page
362+
* that read past its row limit in either phase would time out again.
363+
*/
364+
const [{ timeouts }] = await sql`SELECT last_value::int AS timeouts FROM timed_out_statement`
365+
expect(timeouts).toBe(3)
366+
/**
367+
* Documents: i % 3 <> 0 gives 2,667 Search rows, of which i % 7 = 0 leaves 381 retired, so
368+
* 2,286 are updated, plus `search-doc`. Chunks: 501 Search rows from the fixture plus the
369+
* 2,667 with i % 3 <> 0 among 1003-5002. Equality proves no row was mutated twice.
370+
*/
371+
expect(total).toBe(2287 + 3168)
372+
expect(
373+
(
374+
await sql`SELECT count(*)::int AS n FROM document
375+
WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
376+
)[0].n
377+
).toBe(0)
378+
expect(
379+
(
380+
await sql`SELECT count(*)::int AS n FROM document
381+
WHERE knowledge_base_id = 'ordinary' AND NOT user_excluded AND enabled`
382+
)[0].n
383+
/** `ordinary-doc` plus the 1,333 i % 3 = 0 rows, less the 190 of them seeded retired. */
384+
).toBe(1 + 1333 - 190)
385+
expect(
386+
(await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0]
387+
.n
388+
).toBe(0)
389+
expect(
390+
(
391+
await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'ordinary'`
392+
)[0].n
393+
/** 501 fixture chunks plus the 1,333 i % 3 = 0 rows among 1003-5002. */
394+
).toBe(501 + 1333)
395+
} finally {
396+
await sql`DROP TRIGGER IF EXISTS bound_document_update ON document`
397+
await sql`DROP TRIGGER IF EXISTS bound_embedding_delete ON embedding`
398+
await sql`DROP FUNCTION bound_statement_rows()`
399+
await sql`DROP TABLE committed_statement`
400+
await sql`DROP SEQUENCE timed_out_statement`
401+
}
402+
}, 60_000)
403+
404+
it('bounds the IDs each page reads by the row limit once the limit shrinks', async () => {
405+
const docs = 3000
406+
const bound = 25
407+
await sql`INSERT INTO document (id, knowledge_base_id)
408+
SELECT 'doc-' || lpad(i::text, 5, '0'), 'search' FROM generate_series(1, ${docs}) i`
409+
/** Statements over `bound` rows time out, which pins the row limit at its 25-row floor. */
410+
await sql.unsafe(`CREATE FUNCTION bound_document_update() RETURNS trigger LANGUAGE plpgsql AS $$
411+
BEGIN
412+
IF (SELECT count(*) FROM changed_rows) > ${bound} THEN
413+
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
414+
END IF;
415+
RETURN NULL;
416+
END $$`)
417+
await sql`CREATE TRIGGER bound_document_update AFTER UPDATE ON document
418+
REFERENCING NEW TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_document_update()`
419+
/**
420+
* On a table this small the planner may answer any page with a sequential scan, which reads
421+
* every row whatever the window; production pages use the primary key, so the test does too.
422+
*/
423+
await sql`SET enable_seqscan = off`
424+
/** Document rows read by any scan, counted across committed and rolled-back pages alike. */
425+
async function documentReads() {
426+
await sql`SELECT pg_stat_force_next_flush()`
427+
await admin`SELECT pg_stat_clear_snapshot()`
428+
const [row] = await admin`SELECT (seq_tup_read + coalesce(idx_tup_fetch, 0))::int AS n
429+
FROM pg_stat_user_tables WHERE schemaname = ${schema} AND relname = 'document'`
430+
return row.n
431+
}
432+
try {
433+
const before = await documentReads()
434+
expect(await pass()).toBe(true)
435+
const reads = (await documentReads()) - before
436+
/**
437+
* About 120 pages retire the 3,001 documents 25 at a time once the limit has halved down to
438+
* its floor. A window of four IDs per row reads about 100 IDs and 25 update lookups per page,
439+
* roughly 5 reads per document, plus the chunk phase and completion rechecks. A fixed
440+
* 25,000-ID window re-reads the rest of the table on every attempt, over 100 per document.
441+
*/
442+
expect(reads).toBeLessThan(30 * docs)
443+
expect(
444+
(
445+
await sql`SELECT count(*)::int AS n FROM document
446+
WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
447+
)[0].n
448+
).toBe(0)
449+
} finally {
450+
await sql`RESET enable_seqscan`
451+
await sql`DROP TRIGGER IF EXISTS bound_document_update ON document`
452+
await sql`DROP FUNCTION bound_document_update()`
453+
}
454+
}, 120_000)
455+
456+
it('gives the completed-retirement recheck the completion timeout when it resumes', async () => {
457+
expect(await pass()).toBe(true)
458+
await sql`DELETE FROM script_migrations`
459+
/** Stands in for a recheck that outlasts the two-minute page timeout on a large target set. */
460+
await sql`CREATE FUNCTION require_recheck_timeout() RETURNS boolean LANGUAGE plpgsql AS $$
461+
BEGIN
462+
IF current_setting('statement_timeout') <> '30min' THEN
463+
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
464+
END IF;
465+
RETURN true;
466+
END $$`
467+
await sql`ALTER TABLE search_embedding_cleanup_targets RENAME TO captured_targets`
468+
await sql`CREATE VIEW search_embedding_cleanup_targets AS
469+
SELECT knowledge_base_id FROM captured_targets WHERE require_recheck_timeout()`
470+
try {
471+
expect(await sql`SELECT phase FROM search_embedding_cleanup_progress`).toEqual([
472+
{ phase: 'done' },
473+
])
474+
expect(await pass()).toBe(true)
475+
} finally {
476+
await sql`DROP VIEW IF EXISTS search_embedding_cleanup_targets`
477+
await sql`ALTER TABLE IF EXISTS captured_targets RENAME TO search_embedding_cleanup_targets`
478+
await sql`DROP FUNCTION require_recheck_timeout()`
479+
}
480+
})
481+
482+
it('scans past a run of already-retired documents longer than the row limit in one page', async () => {
483+
/** 6,000 retired Search documents exceed the initial 2,000-row limit but fit one 25,000-ID scan. */
484+
await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled)
485+
SELECT 'doc-' || lpad(i::text, 5, '0'), 'search', true, false FROM generate_series(1, 6000) i`
486+
await sql`INSERT INTO document (id, knowledge_base_id)
487+
SELECT 'doc-' || lpad(i::text, 5, '0'), 'search' FROM generate_series(6001, 6010) i`
488+
await sql`CREATE SEQUENCE document_page_statements`
489+
await sql`CREATE FUNCTION count_document_page() RETURNS trigger LANGUAGE plpgsql AS $$
490+
BEGIN
491+
PERFORM nextval('document_page_statements');
492+
RETURN NULL;
493+
END $$`
494+
/** Statement triggers fire even for zero rows, so this counts every documents-phase page. */
495+
await sql`CREATE TRIGGER count_document_page AFTER UPDATE ON document
496+
FOR EACH STATEMENT EXECUTE FUNCTION count_document_page()`
497+
try {
498+
expect(await pass()).toBe(true)
499+
/** One page retires the 11 unretired rows (10 bulk plus `search-doc`); one more finds the end. */
500+
expect((await sql`SELECT last_value::int AS n FROM document_page_statements`)[0].n).toBe(2)
501+
expect(
502+
(
503+
await sql`SELECT count(*)::int AS n FROM document
504+
WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
505+
)[0].n
506+
).toBe(0)
507+
} finally {
508+
await sql`DROP TRIGGER IF EXISTS count_document_page ON document`
509+
await sql`DROP FUNCTION count_document_page()`
510+
await sql`DROP SEQUENCE document_page_statements`
511+
}
512+
})
513+
514+
it('fails at once on a timeout outside the page mutation instead of shrinking the page', async () => {
515+
await sql`CREATE SEQUENCE completion_attempts`
516+
/** Times out the completion checkpoint, a statement no smaller row limit can speed up. */
517+
await sql`CREATE FUNCTION time_out_completion() RETURNS trigger LANGUAGE plpgsql AS $$
518+
BEGIN
519+
PERFORM nextval('completion_attempts');
520+
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
521+
END $$`
522+
await sql`CREATE TABLE search_embedding_cleanup_progress (
523+
id integer PRIMARY KEY CHECK (id = 1), knowledge_base_id text NOT NULL,
524+
phase text NOT NULL CHECK (phase IN ('documents', 'embeddings', 'done')),
525+
after_id text NOT NULL)`
526+
await sql`CREATE TRIGGER time_out_completion BEFORE UPDATE ON search_embedding_cleanup_progress
527+
FOR EACH ROW WHEN (NEW.phase = 'done') EXECUTE FUNCTION time_out_completion()`
528+
try {
529+
await expect(pass()).rejects.toMatchObject({ code: '57014' })
530+
/** The sequence is not transactional, so it counts rolled-back attempts too. */
531+
expect((await sql`SELECT last_value::int AS n FROM completion_attempts`)[0].n).toBe(1)
532+
expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0)
533+
expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501)
534+
await sql`DROP TRIGGER time_out_completion ON search_embedding_cleanup_progress`
535+
expect(await pass()).toBe(true)
536+
} finally {
537+
await sql`DROP TRIGGER IF EXISTS time_out_completion ON search_embedding_cleanup_progress`
538+
await sql`DROP FUNCTION time_out_completion()`
539+
await sql`DROP SEQUENCE completion_attempts`
540+
}
541+
})
542+
312543
it('rejects inconsistent document ownership before deleting any chunk in the page', async () => {
313544
await sql`UPDATE embedding SET document_id = 'ordinary-doc' WHERE id = '00002'`
314545
await expect(pass()).rejects.toThrow('Search content changed after retirement')

0 commit comments

Comments
 (0)