From c3b2ba9900ed2468bf103bfedcda6a2024e06d59 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 02:48:58 -0700 Subject: [PATCH 1/5] 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. --- ...27_retire_search_embeddings.integration.ts | 79 +++++++++++- .../0027_retire_search_embeddings.ts | 121 ++++++++++++++---- .../search-embedding-retirement.md | 19 ++- 3 files changed, 183 insertions(+), 36 deletions(-) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts index 7e8f7f34269..08ca9672371 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -257,19 +257,27 @@ describe('retiring dormant Search embeddings', () => { await sql`INSERT INTO embedding VALUES ('26003', 'second-search', 'second-search-doc')` await sql`CREATE TABLE deletion_blocker (id text REFERENCES embedding(id))` await sql`INSERT INTO deletion_blocker VALUES ('26002')` + /** Every committed delete sits at or behind the cursor; a failed page leaves all rows past it. */ + async function expectRolledBackPastCursor() { + const [progress] = await sql`SELECT phase, after_id FROM search_embedding_cleanup_progress` + expect(progress.phase).toBe('embeddings') + const [beyond] = + await sql`SELECT count(*)::int AS n FROM embedding WHERE id > ${progress.after_id}` + expect(beyond.n).toBe(26003 - Number(progress.after_id)) + expect(progress.after_id < '26002').toBe(true) + } await expect(pass()).rejects.toThrow() - const before = await sql`SELECT * FROM search_embedding_cleanup_progress` const [remaining] = await sql`SELECT count(*)::int AS n FROM embedding` expect(remaining.n).toBeGreaterThan(501) expect(remaining.n).toBeLessThan(26002) expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0) + await expectRolledBackPastCursor() await sql`UPDATE knowledge_base SET is_search_index = false WHERE id = 'second-search'` await expect(pass()).rejects.toThrow('no longer a Search knowledge base') - expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(remaining.n) + await expectRolledBackPastCursor() await sql`UPDATE knowledge_base SET is_search_index = true WHERE id = 'second-search'` await expect(pass()).rejects.toThrow() - expect(await sql`SELECT * FROM search_embedding_cleanup_progress`).toEqual(before) - expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(remaining.n) + await expectRolledBackPastCursor() await sql`DROP TABLE deletion_blocker` await sql`INSERT INTO document (id, knowledge_base_id) VALUES ('aaa-late-document', 'search')` await sql`INSERT INTO embedding VALUES ('00000', 'search', 'aaa-late-document')` @@ -309,6 +317,69 @@ describe('retiring dormant Search embeddings', () => { ).toBe(0) }, 60_000) + it('splits pages whose writes would outrun the statement timeout and resumes at the split', async () => { + const bound = 300 + await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled) + SELECT 'doc-' || lpad(i::text, 5, '0'), CASE WHEN i % 3 = 0 THEN 'ordinary' ELSE 'search' END, + i % 7 = 0, i % 7 <> 0 + FROM generate_series(1, 4000) i` + await sql`INSERT INTO embedding + SELECT lpad(i::text, 5, '0'), CASE WHEN i % 3 = 0 THEN 'ordinary' ELSE 'search' END, + CASE WHEN i % 3 = 0 THEN 'ordinary-doc' ELSE 'search-doc' END + FROM generate_series(1003, 5002) i` + await sql`CREATE TABLE committed_statement (rows integer NOT NULL)` + /** Stands in for write cost: a statement touching more than `bound` rows times out and rolls back. */ + await sql.unsafe(`CREATE FUNCTION bound_statement_rows() RETURNS trigger LANGUAGE plpgsql AS $$ + DECLARE touched integer; + BEGIN + SELECT count(*) INTO touched FROM changed_rows; + IF touched > ${bound} THEN + RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled'; + END IF; + INSERT INTO committed_statement VALUES (touched); + RETURN NULL; + END $$`) + await sql`CREATE TRIGGER bound_document_update AFTER UPDATE ON document + REFERENCING NEW TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_statement_rows()` + await sql`CREATE TRIGGER bound_embedding_delete AFTER DELETE ON embedding + REFERENCING OLD TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_statement_rows()` + try { + expect(await pass()).toBe(true) + const [{ largest, total }] = await sql`SELECT max(rows)::int AS largest, + sum(rows)::int AS total FROM committed_statement` + expect(largest).toBeLessThanOrEqual(bound) + expect(largest).toBeGreaterThan(0) + /** 2,286 unretired bulk documents plus `search-doc`, then 3,168 chunks; nothing twice. */ + expect(total).toBe(2287 + 3168) + expect( + ( + await sql`SELECT count(*)::int AS n FROM document + WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)` + )[0].n + ).toBe(0) + expect( + ( + await sql`SELECT count(*)::int AS n FROM document + WHERE knowledge_base_id = 'ordinary' AND NOT user_excluded AND enabled` + )[0].n + ).toBe(1 + 1333 - 190) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0] + .n + ).toBe(0) + expect( + ( + await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'ordinary'` + )[0].n + ).toBe(501 + 1333) + } finally { + await sql`DROP TRIGGER IF EXISTS bound_document_update ON document` + await sql`DROP TRIGGER IF EXISTS bound_embedding_delete ON embedding` + await sql`DROP FUNCTION bound_statement_rows()` + await sql`DROP TABLE committed_statement` + } + }, 60_000) + it('rejects inconsistent document ownership before deleting any chunk in the page', async () => { await sql`UPDATE embedding SET document_id = 'ordinary-doc' WHERE id = '00002'` await expect(pass()).rejects.toThrow('Search content changed after retirement') diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts index 1ebb2a837ec..347f100c7a2 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -2,12 +2,29 @@ import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database- import type { ScriptMigration } from '@sim/db/script-migrations/types' import { retryOnLockTimeout } from '@sim/db/scripts/lock-timeout-retry' import { createLogger } from '@sim/logger' +import { getPostgresCancellationReason } from '@sim/utils/errors' import postgres, { type Sql, type TransactionSql } from 'postgres' const logger = createLogger('RetireSearchEmbeddings') -const BATCH_SIZE = 25_000 +/** IDs one page reads in primary-key order; reading is cheap next to the mutation. */ +const SCAN_PAGE_SIZE = 25_000 +/** + * Rows one page may update or delete. Every retired document is a non-HOT update touching each of + * its indexes, and every deleted chunk cascades into its projections, so the write cost of a page, + * not its scan, is what can outrun the statement timeout. + */ +const ROW_LIMIT = { initial: 2_000, min: 25, max: SCAN_PAGE_SIZE } as const +/** A page slower than this halves the row limit; one well under it with a full limit doubles it. */ +const TARGET_PAGE_MS = 30_000 const LOCK_RETRY_BUDGET_MS = 60_000 +interface PageResult { + done: boolean + /** The page stopped at the row limit rather than the end of its scan. */ + limited: boolean + elapsedMs: number +} + interface Progress { knowledge_base_id: string phase: 'documents' | 'embeddings' | 'done' @@ -48,7 +65,7 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { WITH targets AS ( INSERT INTO search_embedding_cleanup_targets (knowledge_base_id) SELECT id FROM knowledge_base WHERE is_search_index AND id > ${afterId} - ORDER BY id LIMIT ${BATCH_SIZE} + ORDER BY id LIMIT ${SCAN_PAGE_SIZE} RETURNING knowledge_base_id ) SELECT max(knowledge_base_id) AS after_id FROM targets` if (page.after_id === null) break @@ -72,11 +89,30 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { const startedAt = Date.now() let batches = 0 - while (!(await retirePage(sql))) { + let rowLimit: number = ROW_LIMIT.initial + for (;;) { + let page: PageResult + try { + page = await retirePage(sql, rowLimit) + } catch (error) { + /** The timed-out page rolled back with its cursor, so it is retried with fewer rows. */ + if (getPostgresCancellationReason(error) !== 'statement_timeout') throw error + if (rowLimit <= ROW_LIMIT.min) throw error + rowLimit = Math.max(ROW_LIMIT.min, Math.floor(rowLimit / 2)) + logger.warn('Search retirement page timed out; retrying with fewer rows', { rowLimit }) + continue + } + if (page.done) break batches++ + if (page.elapsedMs > TARGET_PAGE_MS) { + rowLimit = Math.max(ROW_LIMIT.min, Math.floor(rowLimit / 2)) + } else if (page.limited && page.elapsedMs < TARGET_PAGE_MS / 4) { + rowLimit = Math.min(ROW_LIMIT.max, rowLimit * 2) + } if (batches % 10 === 0) { logger.info('Search embedding retirement progress', { batches, + rowLimit, elapsedMs: Date.now() - startedAt, }) } @@ -88,26 +124,46 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { }, } -async function retirePage(sql: Sql): Promise { +/** + * Retires one page: the target rows among the next `SCAN_PAGE_SIZE` IDs, capped at `rowLimit`. + * A capped page advances the cursor only to its last mutated row, so the rest of the scan is + * read again by the next page; an uncapped page advances past its whole scan. + */ +async function retirePage(sql: Sql, rowLimit: number): Promise { return retryOnLockTimeout( () => sql.begin(async (tx) => { + const startedAt = performance.now() + const result = (done: boolean, limited = false): PageResult => ({ + done, + limited, + elapsedMs: performance.now() - startedAt, + }) await tx`SET LOCAL statement_timeout = '120s'` await tx`SET LOCAL lock_timeout = '1s'` const [progress] = await tx` SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress WHERE id = 1 FOR UPDATE` if (progress.phase === 'done') { await validateTargetMarkers(tx) - return true + return result(true) } if (progress.phase === 'documents') { - const [page] = await tx<{ after_id: string | null; invalid_target: boolean }[]>` + /** Already-retired documents are skipped so they never spend the row limit. */ + const [page] = await tx< + { after_id: string; limited: boolean; invalid_target: boolean }[] + >` WITH source_page AS MATERIALIZED ( - SELECT id, knowledge_base_id FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE} + SELECT id, knowledge_base_id, + (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL + OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL) AS unretired + FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${SCAN_PAGE_SIZE} ), target_page AS MATERIALIZED ( - SELECT p.* FROM source_page p + SELECT p.id, p.knowledge_base_id FROM source_page p JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = p.knowledge_base_id + WHERE p.unretired ORDER BY p.id LIMIT ${rowLimit} + ), page_end AS MATERIALIZED ( + SELECT count(*) >= ${rowLimit} AS limited, max(id) AS last_target FROM target_page ), locked_targets AS MATERIALIZED ( SELECT kb.id, kb.is_search_index FROM knowledge_base kb WHERE kb.id IN (SELECT knowledge_base_id FROM target_page) @@ -125,26 +181,31 @@ async function retirePage(sql: Sql): Promise { AND NOT EXISTS (SELECT 1 FROM invalid_target) AND (NOT d.user_excluded OR d.enabled OR d.processing_queue_token IS NOT NULL OR d.processing_queued_at IS NOT NULL OR d.processing_deferred_until IS NOT NULL) - ) SELECT max(id) AS after_id, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target FROM source_page` - if (page.invalid_target) - throw new Error('Cleanup target is no longer a Search knowledge base') - if (page.after_id === null) { + ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, + e.limited, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target + FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target` + if (!page) { await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` - return false + return result(false) } + if (page.invalid_target) + throw new Error('Cleanup target is no longer a Search knowledge base') await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return false + return result(false, page.limited) } /** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */ const [page] = await tx< - { after_id: string | null; unretired: boolean; invalid_target: boolean }[] + { after_id: string; limited: boolean; unretired: boolean; invalid_target: boolean }[] >` WITH source_page AS MATERIALIZED ( - SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE} + SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${SCAN_PAGE_SIZE} ), target_page AS MATERIALIZED ( SELECT p.* FROM source_page p JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = p.knowledge_base_id + ORDER BY p.id LIMIT ${rowLimit} + ), page_end AS MATERIALIZED ( + SELECT count(*) >= ${rowLimit} AS limited, max(id) AS last_target FROM target_page ), locked_targets AS MATERIALIZED ( SELECT kb.id, kb.is_search_index FROM knowledge_base kb WHERE kb.id IN (SELECT knowledge_base_id FROM target_page) @@ -161,34 +222,40 @@ async function retirePage(sql: Sql): Promise { DELETE FROM embedding e USING target_page p WHERE e.id = p.id AND e.knowledge_base_id = p.knowledge_base_id AND NOT EXISTS (SELECT 1 FROM unretired) AND NOT EXISTS (SELECT 1 FROM invalid_target) - ) SELECT max(id) AS after_id, EXISTS (SELECT 1 FROM unretired) AS unretired, - EXISTS (SELECT 1 FROM invalid_target) AS invalid_target FROM source_page` - if (page.invalid_target) + ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, + e.limited, EXISTS (SELECT 1 FROM unretired) AS unretired, + EXISTS (SELECT 1 FROM invalid_target) AS invalid_target + FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target` + if (page?.invalid_target) throw new Error('Cleanup target is no longer a Search knowledge base') - if (page.unretired) + if (page?.unretired) throw new Error('Search content changed after retirement; stop writers before resuming') - if (page.after_id === null) { - /** A late insert may sort behind either UUID cursor; completion must recheck the target. */ + if (!page) { + /** + * A late insert may sort behind either UUID cursor; completion must recheck the target. + * These rechecks walk every captured KB once, which no single page does. + */ + await tx`SET LOCAL statement_timeout = '30min'` const [unretired] = await tx`SELECT d.id FROM document d JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = d.knowledge_base_id WHERE (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL) LIMIT 1` if (unretired) { await tx`UPDATE search_embedding_cleanup_progress SET phase = 'documents', after_id = '' WHERE id = 1` - return false + return result(false) } const [remaining] = await tx`SELECT e.id FROM embedding e JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = e.knowledge_base_id LIMIT 1` if (remaining) { await tx`UPDATE search_embedding_cleanup_progress SET after_id = '' WHERE id = 1` - return false + return result(false) } await validateTargetMarkers(tx) await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` - return true + return result(true) } await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return false + return result(false, page.limited) }), { budgetMs: LOCK_RETRY_BUDGET_MS, @@ -209,7 +276,7 @@ async function validateTargetMarkers(tx: TransactionSql): Promise { const [page] = await tx<{ after_id: string | null; invalid_target: boolean }[]>` WITH target_page AS MATERIALIZED ( SELECT knowledge_base_id FROM search_embedding_cleanup_targets - WHERE knowledge_base_id > ${afterId} ORDER BY knowledge_base_id LIMIT ${BATCH_SIZE} + WHERE knowledge_base_id > ${afterId} ORDER BY knowledge_base_id LIMIT ${SCAN_PAGE_SIZE} ), locked_targets AS MATERIALIZED ( SELECT kb.id, kb.is_search_index FROM knowledge_base kb WHERE kb.id IN (SELECT knowledge_base_id FROM target_page) diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index d55cb4594e2..3a1c0908323 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -30,11 +30,20 @@ resume the saved scope, phase, cursor and maintenance checkpoints. The existing maintenance implementation rebuilds HNSW indexes and vacuums affected tables before deployment continues. -Pages contain at most 25,000 IDs and execute sequentially without a pacing delay. Materialized SQL -pages keep the IDs inside PostgreSQL; the migration process receives only a cursor and a validation -result. Each page uses a two-minute statement timeout and a one-second lock timeout. Brief lock -timeouts retry the rolled-back page with bounded backoff for up to one minute. Other errors, or -exhausted lock retries, fail the migration without a completion receipt. Progress is logged every +Each page reads at most 25,000 IDs and mutates at most a row limit of them, and pages execute +sequentially without a pacing delay. Retiring a document is a non-HOT update that writes every +index on `document`, and deleting a chunk cascades into its projections, so a page's cost follows +the target rows it mutates, not the IDs it reads. A page that reaches the row limit advances the +cursor only to its last mutated row; the rest of its scan is read again by the next page. Documents +that are already retired never count against the limit. The limit starts at 2,000 rows, halves after +a page slower than 30 seconds, and doubles (up to 25,000) after a fast page that reached it. +Materialized SQL pages keep the IDs inside PostgreSQL; the migration process receives only a cursor +and a validation result. Each page uses a two-minute statement timeout and a one-second lock timeout. +A page that exceeds the statement timeout rolls back with its cursor and is retried with half the +row limit; one that still times out at 25 rows fails the migration. The completion rechecks, which +walk every captured KB once, run with a 30-minute timeout. Brief lock timeouts retry the rolled-back +page with bounded backoff for up to one minute. Other errors, or exhausted lock retries, fail the +migration without a completion receipt. Progress is logged every ten pages. The deployment job retains its five-hour overall timeout; it is not a runtime estimate. After interruption or failure, rerun the migration job, or run From a921532a6f15f2e202e7c1fcf12443d143a65434 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 02:56:18 -0700 Subject: [PATCH 2/5] 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. --- ...27_retire_search_embeddings.integration.ts | 43 +++++- .../0027_retire_search_embeddings.ts | 138 ++++++++++++++---- .../search-embedding-retirement.md | 39 +++-- 3 files changed, 172 insertions(+), 48 deletions(-) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts index 08ca9672371..22f9d945643 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -257,7 +257,11 @@ describe('retiring dormant Search embeddings', () => { await sql`INSERT INTO embedding VALUES ('26003', 'second-search', 'second-search-doc')` await sql`CREATE TABLE deletion_blocker (id text REFERENCES embedding(id))` await sql`INSERT INTO deletion_blocker VALUES ('26002')` - /** Every committed delete sits at or behind the cursor; a failed page leaves all rows past it. */ + /** + * Pages split by mutated rows under an adaptive limit, so how far each failed run gets depends + * on page geometry. The geometry-free invariant: every committed delete sits at or behind the + * cursor, and a failed page leaves every row past it (IDs 00001-26003 are contiguous). + */ async function expectRolledBackPastCursor() { const [progress] = await sql`SELECT phase, after_id FROM search_embedding_cleanup_progress` expect(progress.phase).toBe('embeddings') @@ -349,7 +353,11 @@ describe('retiring dormant Search embeddings', () => { sum(rows)::int AS total FROM committed_statement` expect(largest).toBeLessThanOrEqual(bound) expect(largest).toBeGreaterThan(0) - /** 2,286 unretired bulk documents plus `search-doc`, then 3,168 chunks; nothing twice. */ + /** + * Documents: i % 3 <> 0 gives 2,667 Search rows, of which i % 7 = 0 leaves 381 retired, so + * 2,286 are updated, plus `search-doc`. Chunks: 501 Search rows from the fixture plus the + * 2,667 with i % 3 <> 0 among 1003-5002. Equality proves no row was mutated twice. + */ expect(total).toBe(2287 + 3168) expect( ( @@ -362,6 +370,7 @@ describe('retiring dormant Search embeddings', () => { await sql`SELECT count(*)::int AS n FROM document WHERE knowledge_base_id = 'ordinary' AND NOT user_excluded AND enabled` )[0].n + /** `ordinary-doc` plus the 1,333 i % 3 = 0 rows, less the 190 of them seeded retired. */ ).toBe(1 + 1333 - 190) expect( (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0] @@ -371,6 +380,7 @@ describe('retiring dormant Search embeddings', () => { ( await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'ordinary'` )[0].n + /** 501 fixture chunks plus the 1,333 i % 3 = 0 rows among 1003-5002. */ ).toBe(501 + 1333) } finally { await sql`DROP TRIGGER IF EXISTS bound_document_update ON document` @@ -380,6 +390,35 @@ describe('retiring dormant Search embeddings', () => { } }, 60_000) + it('fails at once on a timeout outside the page mutation instead of shrinking the page', async () => { + await sql`CREATE SEQUENCE completion_attempts` + /** Times out the completion checkpoint, a statement no smaller row limit can speed up. */ + await sql`CREATE FUNCTION time_out_completion() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + PERFORM nextval('completion_attempts'); + RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled'; + END $$` + await sql`CREATE TABLE search_embedding_cleanup_progress ( + id integer PRIMARY KEY CHECK (id = 1), knowledge_base_id text NOT NULL, + phase text NOT NULL CHECK (phase IN ('documents', 'embeddings', 'done')), + after_id text NOT NULL)` + await sql`CREATE TRIGGER time_out_completion BEFORE UPDATE ON search_embedding_cleanup_progress + FOR EACH ROW WHEN (NEW.phase = 'done') EXECUTE FUNCTION time_out_completion()` + try { + await expect(pass()).rejects.toMatchObject({ code: '57014' }) + /** The sequence is not transactional, so it counts rolled-back attempts too. */ + expect((await sql`SELECT last_value::int AS n FROM completion_attempts`)[0].n).toBe(1) + expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501) + await sql`DROP TRIGGER time_out_completion ON search_embedding_cleanup_progress` + expect(await pass()).toBe(true) + } finally { + await sql`DROP TRIGGER IF EXISTS time_out_completion ON search_embedding_cleanup_progress` + await sql`DROP FUNCTION time_out_completion()` + await sql`DROP SEQUENCE completion_attempts` + } + }) + it('rejects inconsistent document ownership before deleting any chunk in the page', async () => { await sql`UPDATE embedding SET document_id = 'ordinary-doc' WHERE id = '00002'` await expect(pass()).rejects.toThrow('Search content changed after retirement') diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts index 347f100c7a2..29658ad0ddc 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -3,6 +3,7 @@ import type { ScriptMigration } from '@sim/db/script-migrations/types' import { retryOnLockTimeout } from '@sim/db/scripts/lock-timeout-retry' import { createLogger } from '@sim/logger' import { getPostgresCancellationReason } from '@sim/utils/errors' +import { sleep } from '@sim/utils/helpers' import postgres, { type Sql, type TransactionSql } from 'postgres' const logger = createLogger('RetireSearchEmbeddings') @@ -13,21 +14,60 @@ const SCAN_PAGE_SIZE = 25_000 * its indexes, and every deleted chunk cascades into its projections, so the write cost of a page, * not its scan, is what can outrun the statement timeout. */ -const ROW_LIMIT = { initial: 2_000, min: 25, max: SCAN_PAGE_SIZE } as const -/** A page slower than this halves the row limit; one well under it with a full limit doubles it. */ -const TARGET_PAGE_MS = 30_000 +const ROW_LIMIT = { initial: 2_000, min: 25, max: 8_000 } as const +/** A page slower than this halves the row limit. */ +const SLOW_PAGE_MS = 30_000 +/** A page faster than this that reached the row limit doubles it. */ +const FAST_PAGE_MS = SLOW_PAGE_MS / 4 +/** Each page is followed by a pause as long as the page, up to this, to leave the primary headroom. */ +const MAX_PAGE_PAUSE_MS = 5_000 const LOCK_RETRY_BUDGET_MS = 60_000 +type Phase = 'documents' | 'embeddings' | 'done' + interface PageResult { done: boolean + /** The phase the page ran in. */ + phase: Phase + /** The cursor the page committed. */ + afterId: string + /** Rows the page updated or deleted. */ + mutated: number /** The page stopped at the row limit rather than the end of its scan. */ limited: boolean - elapsedMs: number + /** A phase change the page committed. */ + transition?: 'embeddings' | 'documents_rescan' | 'embeddings_rescan' +} + +/** + * A statement timeout from a page's mutating statement, the only statement a smaller page speeds + * up. Any other timeout, such as a completion recheck, propagates unchanged and fails the run. + */ +class PageMutationTimeout extends Error { + override name = 'PageMutationTimeout' + constructor(readonly timeout: unknown) { + super('Search retirement page mutation timed out', { cause: timeout }) + } +} + +async function pageMutation(statement: PromiseLike): Promise { + try { + return await statement + } catch (error) { + if (getPostgresCancellationReason(error) === 'statement_timeout') { + throw new PageMutationTimeout(error) + } + throw error + } +} + +function halve(rowLimit: number): number { + return Math.max(ROW_LIMIT.min, Math.floor(rowLimit / 2)) } interface Progress { knowledge_base_id: string - phase: 'documents' | 'embeddings' | 'done' + phase: Phase after_id: string } @@ -89,36 +129,57 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { const startedAt = Date.now() let batches = 0 + let mutated = 0 let rowLimit: number = ROW_LIMIT.initial for (;;) { + /** Timed around the whole call, so the synchronous-replication wait at commit counts. */ + const pageStartedAt = performance.now() let page: PageResult try { page = await retirePage(sql, rowLimit) } catch (error) { + if (!(error instanceof PageMutationTimeout)) throw error /** The timed-out page rolled back with its cursor, so it is retried with fewer rows. */ - if (getPostgresCancellationReason(error) !== 'statement_timeout') throw error - if (rowLimit <= ROW_LIMIT.min) throw error - rowLimit = Math.max(ROW_LIMIT.min, Math.floor(rowLimit / 2)) + if (rowLimit <= ROW_LIMIT.min) throw error.timeout + rowLimit = halve(rowLimit) logger.warn('Search retirement page timed out; retrying with fewer rows', { rowLimit }) continue } if (page.done) break + const pageMs = performance.now() - pageStartedAt batches++ - if (page.elapsedMs > TARGET_PAGE_MS) { - rowLimit = Math.max(ROW_LIMIT.min, Math.floor(rowLimit / 2)) - } else if (page.limited && page.elapsedMs < TARGET_PAGE_MS / 4) { + mutated += page.mutated + if (page.transition) { + /** A phase change may include a full recheck, which says nothing about page cost. */ + logger.info('Search retirement phase changed', { + transition: page.transition, + batches, + mutated, + }) + } else if (pageMs > SLOW_PAGE_MS) { + rowLimit = halve(rowLimit) + logger.warn('Search retirement page was slow; halving the row limit', { + pageMs: Math.round(pageMs), + rowLimit, + }) + } else if (page.limited && pageMs < FAST_PAGE_MS) { rowLimit = Math.min(ROW_LIMIT.max, rowLimit * 2) } if (batches % 10 === 0) { logger.info('Search embedding retirement progress', { batches, + phase: page.phase, + afterId: page.afterId, + mutated, rowLimit, elapsedMs: Date.now() - startedAt, }) } + await sleep(Math.min(pageMs, MAX_PAGE_PAUSE_MS)) } logger.info('Selected Search knowledge bases retired', { batches, + mutated, elapsedMs: Date.now() - startedAt, }) }, @@ -133,25 +194,27 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { return retryOnLockTimeout( () => sql.begin(async (tx) => { - const startedAt = performance.now() - const result = (done: boolean, limited = false): PageResult => ({ - done, - limited, - elapsedMs: performance.now() - startedAt, - }) await tx`SET LOCAL statement_timeout = '120s'` await tx`SET LOCAL lock_timeout = '1s'` const [progress] = await tx` SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress WHERE id = 1 FOR UPDATE` + const result = (page: Partial = {}): PageResult => ({ + done: false, + phase: progress.phase, + afterId: progress.after_id, + mutated: 0, + limited: false, + ...page, + }) if (progress.phase === 'done') { await validateTargetMarkers(tx) - return result(true) + return result({ done: true }) } if (progress.phase === 'documents') { /** Already-retired documents are skipped so they never spend the row limit. */ - const [page] = await tx< - { after_id: string; limited: boolean; invalid_target: boolean }[] + const [page] = await pageMutation(tx< + { after_id: string; limited: boolean; mutated: number; invalid_target: boolean }[] >` WITH source_page AS MATERIALIZED ( SELECT id, knowledge_base_id, @@ -181,22 +244,30 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { AND NOT EXISTS (SELECT 1 FROM invalid_target) AND (NOT d.user_excluded OR d.enabled OR d.processing_queue_token IS NOT NULL OR d.processing_queued_at IS NOT NULL OR d.processing_deferred_until IS NOT NULL) + RETURNING d.id ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, - e.limited, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target - FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target` + e.limited, (SELECT count(*) FROM retired)::int AS mutated, + EXISTS (SELECT 1 FROM invalid_target) AS invalid_target + FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target`) if (!page) { await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` - return result(false) + return result({ afterId: '', transition: 'embeddings' }) } if (page.invalid_target) throw new Error('Cleanup target is no longer a Search knowledge base') await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return result(false, page.limited) + return result({ afterId: page.after_id, mutated: page.mutated, limited: page.limited }) } /** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */ - const [page] = await tx< - { after_id: string; limited: boolean; unretired: boolean; invalid_target: boolean }[] + const [page] = await pageMutation(tx< + { + after_id: string + limited: boolean + mutated: number + unretired: boolean + invalid_target: boolean + }[] >` WITH source_page AS MATERIALIZED ( SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${SCAN_PAGE_SIZE} @@ -222,10 +293,12 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { DELETE FROM embedding e USING target_page p WHERE e.id = p.id AND e.knowledge_base_id = p.knowledge_base_id AND NOT EXISTS (SELECT 1 FROM unretired) AND NOT EXISTS (SELECT 1 FROM invalid_target) + RETURNING e.id ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, - e.limited, EXISTS (SELECT 1 FROM unretired) AS unretired, + e.limited, (SELECT count(*) FROM deleted)::int AS mutated, + EXISTS (SELECT 1 FROM unretired) AS unretired, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target - FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target` + FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target`) if (page?.invalid_target) throw new Error('Cleanup target is no longer a Search knowledge base') if (page?.unretired) @@ -236,26 +309,27 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { * These rechecks walk every captured KB once, which no single page does. */ await tx`SET LOCAL statement_timeout = '30min'` + logger.info('Rechecking captured Search knowledge bases before completion') const [unretired] = await tx`SELECT d.id FROM document d JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = d.knowledge_base_id WHERE (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL) LIMIT 1` if (unretired) { await tx`UPDATE search_embedding_cleanup_progress SET phase = 'documents', after_id = '' WHERE id = 1` - return result(false) + return result({ afterId: '', transition: 'documents_rescan' }) } const [remaining] = await tx`SELECT e.id FROM embedding e JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = e.knowledge_base_id LIMIT 1` if (remaining) { await tx`UPDATE search_embedding_cleanup_progress SET after_id = '' WHERE id = 1` - return result(false) + return result({ afterId: '', transition: 'embeddings_rescan' }) } await validateTargetMarkers(tx) await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` - return result(true) + return result({ done: true }) } await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return result(false, page.limited) + return result({ afterId: page.after_id, mutated: page.mutated, limited: page.limited }) }), { budgetMs: LOCK_RETRY_BUDGET_MS, diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index 3a1c0908323..e84683c8bbd 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -30,21 +30,32 @@ resume the saved scope, phase, cursor and maintenance checkpoints. The existing maintenance implementation rebuilds HNSW indexes and vacuums affected tables before deployment continues. -Each page reads at most 25,000 IDs and mutates at most a row limit of them, and pages execute -sequentially without a pacing delay. Retiring a document is a non-HOT update that writes every -index on `document`, and deleting a chunk cascades into its projections, so a page's cost follows -the target rows it mutates, not the IDs it reads. A page that reaches the row limit advances the -cursor only to its last mutated row; the rest of its scan is read again by the next page. Documents -that are already retired never count against the limit. The limit starts at 2,000 rows, halves after -a page slower than 30 seconds, and doubles (up to 25,000) after a fast page that reached it. +Each page reads at most 25,000 IDs and mutates at most a row limit of them. Pages execute +sequentially, and each is followed by a pause as long as the page took, up to five seconds, to +leave the primary headroom. Retiring a document is a non-HOT update that writes every index on +`document`, and deleting a chunk cascades into its projections, so a page's cost follows the target +rows it mutates, not the IDs it reads. A page that reaches the row limit advances the cursor only to +its last mutated row; the rest of its scan is read again by the next page. Documents that are +already retired never count against the limit. + +The row limit starts at 2,000 rows. A page is timed from the start of its transaction through its +commit, including the synchronous-replication wait and any lock-timeout retries. A page slower than +30 seconds halves the limit. A fast page, one under 7.5 seconds, that reached the limit doubles it, +up to 8,000. The limit never drops below 25 rows. Phase changes do not adjust it. + Materialized SQL pages keep the IDs inside PostgreSQL; the migration process receives only a cursor -and a validation result. Each page uses a two-minute statement timeout and a one-second lock timeout. -A page that exceeds the statement timeout rolls back with its cursor and is retried with half the -row limit; one that still times out at 25 rows fails the migration. The completion rechecks, which -walk every captured KB once, run with a 30-minute timeout. Brief lock timeouts retry the rolled-back -page with bounded backoff for up to one minute. Other errors, or exhausted lock retries, fail the -migration without a completion receipt. Progress is logged every -ten pages. The deployment job retains its five-hour overall timeout; it is not a runtime estimate. +and a validation result. Each page uses a two-minute statement timeout and a one-second lock +timeout. If a page's mutating statement exceeds the statement timeout, the page rolls back with its +cursor and is retried with half the row limit; one that still times out at 25 rows fails the +migration. Any other statement timeout fails the migration at once, because a smaller page cannot +speed it up. The completion rechecks, which walk every captured KB once, run with a 30-minute +timeout. Brief lock timeouts retry the rolled-back page with bounded backoff for up to one minute. +Other errors, or exhausted lock retries, fail the migration without a completion receipt. + +Every ten pages the migration logs the phase, cursor, rows mutated so far and current row limit. It +also logs each halving after a slow page, each phase change and the start of the completion +recheck. The deployment job retains its five-hour overall timeout; it is not a runtime estimate. A +large cleanup can need more than one job run, and each run resumes from the saved cursor. After interruption or failure, rerun the migration job, or run `bun run packages/db/script-migrations/0027_retire_search_embeddings.ts` with the writer supplied From 404e365ca5c81448dc02df37deaecb164d090cdb Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 03:02:01 -0700 Subject: [PATCH 3/5] 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. --- ...27_retire_search_embeddings.integration.ts | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts index 22f9d945643..d3b3573919a 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -390,6 +390,38 @@ describe('retiring dormant Search embeddings', () => { } }, 60_000) + it('scans past a run of already-retired documents longer than the row limit in one page', async () => { + /** 6,000 retired Search documents exceed the initial 2,000-row limit but fit one 25,000-ID scan. */ + await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled) + SELECT 'doc-' || lpad(i::text, 5, '0'), 'search', true, false FROM generate_series(1, 6000) i` + await sql`INSERT INTO document (id, knowledge_base_id) + SELECT 'doc-' || lpad(i::text, 5, '0'), 'search' FROM generate_series(6001, 6010) i` + await sql`CREATE SEQUENCE document_page_statements` + await sql`CREATE FUNCTION count_document_page() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + PERFORM nextval('document_page_statements'); + RETURN NULL; + END $$` + /** Statement triggers fire even for zero rows, so this counts every documents-phase page. */ + await sql`CREATE TRIGGER count_document_page AFTER UPDATE ON document + FOR EACH STATEMENT EXECUTE FUNCTION count_document_page()` + try { + expect(await pass()).toBe(true) + /** One page retires the 11 unretired rows (10 bulk plus `search-doc`); one more finds the end. */ + expect((await sql`SELECT last_value::int AS n FROM document_page_statements`)[0].n).toBe(2) + expect( + ( + await sql`SELECT count(*)::int AS n FROM document + WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)` + )[0].n + ).toBe(0) + } finally { + await sql`DROP TRIGGER IF EXISTS count_document_page ON document` + await sql`DROP FUNCTION count_document_page()` + await sql`DROP SEQUENCE document_page_statements` + } + }) + it('fails at once on a timeout outside the page mutation instead of shrinking the page', async () => { await sql`CREATE SEQUENCE completion_attempts` /** Times out the completion checkpoint, a statement no smaller row limit can speed up. */ From ee1d868f292cd362d4b8f8a2bd0c45af6099d798 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 03:30:45 -0700 Subject: [PATCH 4/5] 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. --- ...27_retire_search_embeddings.integration.ts | 78 +++++++++++++++++++ .../0027_retire_search_embeddings.ts | 36 +++++---- .../search-embedding-retirement.md | 14 ++-- 3 files changed, 107 insertions(+), 21 deletions(-) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts index d3b3573919a..a0785173012 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -390,6 +390,84 @@ describe('retiring dormant Search embeddings', () => { } }, 60_000) + it('bounds the IDs each page reads by the row limit once the limit shrinks', async () => { + const docs = 3000 + const bound = 25 + await sql`INSERT INTO document (id, knowledge_base_id) + SELECT 'doc-' || lpad(i::text, 5, '0'), 'search' FROM generate_series(1, ${docs}) i` + /** Statements over `bound` rows time out, which pins the row limit at its 25-row floor. */ + await sql.unsafe(`CREATE FUNCTION bound_document_update() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + IF (SELECT count(*) FROM changed_rows) > ${bound} THEN + RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled'; + END IF; + RETURN NULL; + END $$`) + await sql`CREATE TRIGGER bound_document_update AFTER UPDATE ON document + REFERENCING NEW TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_document_update()` + /** + * On a table this small the planner may answer any page with a sequential scan, which reads + * every row whatever the window; production pages use the primary key, so the test does too. + */ + await sql`SET enable_seqscan = off` + /** Document rows read by any scan, counted across committed and rolled-back pages alike. */ + async function documentReads() { + await sql`SELECT pg_stat_force_next_flush()` + await admin`SELECT pg_stat_clear_snapshot()` + const [row] = await admin`SELECT (seq_tup_read + coalesce(idx_tup_fetch, 0))::int AS n + FROM pg_stat_user_tables WHERE schemaname = ${schema} AND relname = 'document'` + return row.n + } + try { + const before = await documentReads() + expect(await pass()).toBe(true) + const reads = (await documentReads()) - before + /** + * About 120 pages retire the 3,001 documents 25 at a time, each after one rolled-back attempt + * at 50 rows. A window of four IDs per row reads about 300 IDs and 75 update lookups per page, + * roughly 15 reads per document, plus the chunk phase and completion rechecks. A fixed + * 25,000-ID window re-reads the rest of the table on every attempt, over 100 per document. + */ + expect(reads).toBeLessThan(30 * docs) + expect( + ( + await sql`SELECT count(*)::int AS n FROM document + WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)` + )[0].n + ).toBe(0) + } finally { + await sql`RESET enable_seqscan` + await sql`DROP TRIGGER IF EXISTS bound_document_update ON document` + await sql`DROP FUNCTION bound_document_update()` + } + }, 120_000) + + it('gives the completed-retirement recheck the completion timeout when it resumes', async () => { + expect(await pass()).toBe(true) + await sql`DELETE FROM script_migrations` + /** Stands in for a recheck that outlasts the two-minute page timeout on a large target set. */ + await sql`CREATE FUNCTION require_recheck_timeout() RETURNS boolean LANGUAGE plpgsql AS $$ + BEGIN + IF current_setting('statement_timeout') <> '30min' THEN + RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled'; + END IF; + RETURN true; + END $$` + await sql`ALTER TABLE search_embedding_cleanup_targets RENAME TO captured_targets` + await sql`CREATE VIEW search_embedding_cleanup_targets AS + SELECT knowledge_base_id FROM captured_targets WHERE require_recheck_timeout()` + try { + expect(await sql`SELECT phase FROM search_embedding_cleanup_progress`).toEqual([ + { phase: 'done' }, + ]) + expect(await pass()).toBe(true) + } finally { + await sql`DROP VIEW IF EXISTS search_embedding_cleanup_targets` + await sql`ALTER TABLE IF EXISTS captured_targets RENAME TO search_embedding_cleanup_targets` + await sql`DROP FUNCTION require_recheck_timeout()` + } + }) + it('scans past a run of already-retired documents longer than the row limit in one page', async () => { /** 6,000 retired Search documents exceed the initial 2,000-row limit but fit one 25,000-ID scan. */ await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts index 29658ad0ddc..ee35fe96710 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -7,8 +7,14 @@ import { sleep } from '@sim/utils/helpers' import postgres, { type Sql, type TransactionSql } from 'postgres' const logger = createLogger('RetireSearchEmbeddings') -/** IDs one page reads in primary-key order; reading is cheap next to the mutation. */ +/** Most IDs one page reads in primary-key order; reading is cheap next to the mutation. */ const SCAN_PAGE_SIZE = 25_000 +/** + * A page reads at most this many IDs per row it may mutate. A page that reaches the row limit is + * re-read from its last mutated row, so a window far wider than the limit would be read again and + * again as the limit shrinks, exactly when the database is already slow. + */ +const SCAN_ROWS_PER_MUTATION = 4 /** * Rows one page may update or delete. Every retired document is a non-HOT update touching each of * 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 const ROW_LIMIT = { initial: 2_000, min: 25, max: 8_000 } as const /** A page slower than this halves the row limit. */ const SLOW_PAGE_MS = 30_000 -/** A page faster than this that reached the row limit doubles it. */ +/** A page faster than this doubles the row limit, widening its scan window with it. */ const FAST_PAGE_MS = SLOW_PAGE_MS / 4 /** Each page is followed by a pause as long as the page, up to this, to leave the primary headroom. */ const MAX_PAGE_PAUSE_MS = 5_000 @@ -33,8 +39,6 @@ interface PageResult { afterId: string /** Rows the page updated or deleted. */ mutated: number - /** The page stopped at the row limit rather than the end of its scan. */ - limited: boolean /** A phase change the page committed. */ transition?: 'embeddings' | 'documents_rescan' | 'embeddings_rescan' } @@ -162,7 +166,7 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { pageMs: Math.round(pageMs), rowLimit, }) - } else if (page.limited && pageMs < FAST_PAGE_MS) { + } else if (pageMs < FAST_PAGE_MS) { rowLimit = Math.min(ROW_LIMIT.max, rowLimit * 2) } if (batches % 10 === 0) { @@ -186,7 +190,8 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { } /** - * Retires one page: the target rows among the next `SCAN_PAGE_SIZE` IDs, capped at `rowLimit`. + * Retires one page: the target rows among the next `rowLimit * SCAN_ROWS_PER_MUTATION` IDs (at + * most `SCAN_PAGE_SIZE`), capped at `rowLimit`. * A capped page advances the cursor only to its last mutated row, so the rest of the scan is * read again by the next page; an uncapped page advances past its whole scan. */ @@ -194,6 +199,7 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { return retryOnLockTimeout( () => sql.begin(async (tx) => { + const scanLimit = Math.min(SCAN_PAGE_SIZE, rowLimit * SCAN_ROWS_PER_MUTATION) await tx`SET LOCAL statement_timeout = '120s'` await tx`SET LOCAL lock_timeout = '1s'` const [progress] = await tx` @@ -203,10 +209,11 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { phase: progress.phase, afterId: progress.after_id, mutated: 0, - limited: false, ...page, }) if (progress.phase === 'done') { + /** A retry after maintenance failed rechecks every captured KB, as completion did. */ + await tx`SET LOCAL statement_timeout = '30min'` await validateTargetMarkers(tx) return result({ done: true }) } @@ -214,13 +221,13 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { if (progress.phase === 'documents') { /** Already-retired documents are skipped so they never spend the row limit. */ const [page] = await pageMutation(tx< - { after_id: string; limited: boolean; mutated: number; invalid_target: boolean }[] + { after_id: string; mutated: number; invalid_target: boolean }[] >` WITH source_page AS MATERIALIZED ( SELECT id, knowledge_base_id, (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL) AS unretired - FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${SCAN_PAGE_SIZE} + FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${scanLimit} ), target_page AS MATERIALIZED ( SELECT p.id, p.knowledge_base_id FROM source_page p 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 { OR d.processing_queued_at IS NOT NULL OR d.processing_deferred_until IS NOT NULL) RETURNING d.id ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, - e.limited, (SELECT count(*) FROM retired)::int AS mutated, + (SELECT count(*) FROM retired)::int AS mutated, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target`) if (!page) { @@ -256,21 +263,20 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { if (page.invalid_target) throw new Error('Cleanup target is no longer a Search knowledge base') await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return result({ afterId: page.after_id, mutated: page.mutated, limited: page.limited }) + return result({ afterId: page.after_id, mutated: page.mutated }) } /** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */ const [page] = await pageMutation(tx< { after_id: string - limited: boolean mutated: number unretired: boolean invalid_target: boolean }[] >` WITH source_page AS MATERIALIZED ( - SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${SCAN_PAGE_SIZE} + SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${scanLimit} ), target_page AS MATERIALIZED ( SELECT p.* FROM source_page p 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 { AND NOT EXISTS (SELECT 1 FROM unretired) AND NOT EXISTS (SELECT 1 FROM invalid_target) RETURNING e.id ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, - e.limited, (SELECT count(*) FROM deleted)::int AS mutated, + (SELECT count(*) FROM deleted)::int AS mutated, EXISTS (SELECT 1 FROM unretired) AS unretired, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target 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 { return result({ done: true }) } await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return result({ afterId: page.after_id, mutated: page.mutated, limited: page.limited }) + return result({ afterId: page.after_id, mutated: page.mutated }) }), { budgetMs: LOCK_RETRY_BUDGET_MS, diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index e84683c8bbd..0b966d14874 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -30,18 +30,20 @@ resume the saved scope, phase, cursor and maintenance checkpoints. The existing maintenance implementation rebuilds HNSW indexes and vacuums affected tables before deployment continues. -Each page reads at most 25,000 IDs and mutates at most a row limit of them. Pages execute +Each page mutates at most a row limit of target rows and reads at most four IDs per row of that +limit, never more than 25,000 IDs. Pages execute sequentially, and each is followed by a pause as long as the page took, up to five seconds, to leave the primary headroom. Retiring a document is a non-HOT update that writes every index on `document`, and deleting a chunk cascades into its projections, so a page's cost follows the target rows it mutates, not the IDs it reads. A page that reaches the row limit advances the cursor only to -its last mutated row; the rest of its scan is read again by the next page. Documents that are -already retired never count against the limit. +its last mutated row; the rest of its scan is read again by the next page. Tying the scan window to +the limit keeps that re-reading proportional to the work, even after the limit shrinks. Documents +that are already retired never count against the limit. The row limit starts at 2,000 rows. A page is timed from the start of its transaction through its commit, including the synchronous-replication wait and any lock-timeout retries. A page slower than -30 seconds halves the limit. A fast page, one under 7.5 seconds, that reached the limit doubles it, -up to 8,000. The limit never drops below 25 rows. Phase changes do not adjust it. +30 seconds halves the limit. A fast page, one under 7.5 seconds, doubles it up to 8,000, which +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. Materialized SQL pages keep the IDs inside PostgreSQL; the migration process receives only a cursor 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 behind either cursor and restarts the affected phase if needed. A final bounded pass validates all captured KB markers, including empty KBs and KBs whose rows were already scanned, holding shared marker locks until the completion checkpoint commits. Resuming a completed cleanup before maintenance -also revalidates the captured set. Keep target writers stopped and +also revalidates the captured set, with the same 30-minute timeout as the completion recheck. Keep target writers stopped and do not change their Search markers during the pass. Inspect progress with: From 1e7d6fd546ec3fcb705ebd0b95761540499415de Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 10:16:13 -0700 Subject: [PATCH 5/5] 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. --- ...027_retire_search_embeddings.integration.ts | 17 ++++++++++++++--- .../0027_retire_search_embeddings.ts | 18 ++++++++++++++---- .../search-embedding-retirement.md | 5 +++-- 3 files changed, 31 insertions(+), 9 deletions(-) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts index a0785173012..b1e54720bd7 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -332,12 +332,15 @@ describe('retiring dormant Search embeddings', () => { CASE WHEN i % 3 = 0 THEN 'ordinary-doc' ELSE 'search-doc' END FROM generate_series(1003, 5002) i` await sql`CREATE TABLE committed_statement (rows integer NOT NULL)` + /** Sequences are not transactional, so this counts the timed-out statements that rolled back. */ + await sql`CREATE SEQUENCE timed_out_statement` /** Stands in for write cost: a statement touching more than `bound` rows times out and rolls back. */ await sql.unsafe(`CREATE FUNCTION bound_statement_rows() RETURNS trigger LANGUAGE plpgsql AS $$ DECLARE touched integer; BEGIN SELECT count(*) INTO touched FROM changed_rows; IF touched > ${bound} THEN + PERFORM nextval('timed_out_statement'); RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled'; END IF; INSERT INTO committed_statement VALUES (touched); @@ -353,6 +356,13 @@ describe('retiring dormant Search embeddings', () => { sum(rows)::int AS total FROM committed_statement` expect(largest).toBeLessThanOrEqual(bound) expect(largest).toBeGreaterThan(0) + /** + * The limit halves 2,000 → 1,000 → 500 → 250 on the first document page and never grows back + * to a size that timed out, in either phase: exactly three rolled-back statements. A page + * that read past its row limit in either phase would time out again. + */ + const [{ timeouts }] = await sql`SELECT last_value::int AS timeouts FROM timed_out_statement` + expect(timeouts).toBe(3) /** * Documents: i % 3 <> 0 gives 2,667 Search rows, of which i % 7 = 0 leaves 381 retired, so * 2,286 are updated, plus `search-doc`. Chunks: 501 Search rows from the fixture plus the @@ -387,6 +397,7 @@ describe('retiring dormant Search embeddings', () => { await sql`DROP TRIGGER IF EXISTS bound_embedding_delete ON embedding` await sql`DROP FUNCTION bound_statement_rows()` await sql`DROP TABLE committed_statement` + await sql`DROP SEQUENCE timed_out_statement` } }, 60_000) @@ -423,9 +434,9 @@ describe('retiring dormant Search embeddings', () => { expect(await pass()).toBe(true) const reads = (await documentReads()) - before /** - * About 120 pages retire the 3,001 documents 25 at a time, each after one rolled-back attempt - * at 50 rows. A window of four IDs per row reads about 300 IDs and 75 update lookups per page, - * roughly 15 reads per document, plus the chunk phase and completion rechecks. A fixed + * About 120 pages retire the 3,001 documents 25 at a time once the limit has halved down to + * its floor. A window of four IDs per row reads about 100 IDs and 25 update lookups per page, + * roughly 5 reads per document, plus the chunk phase and completion rechecks. A fixed * 25,000-ID window re-reads the rest of the table on every attempt, over 100 per document. */ expect(reads).toBeLessThan(30 * docs) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts index ee35fe96710..62554665582 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -23,9 +23,15 @@ const SCAN_ROWS_PER_MUTATION = 4 const ROW_LIMIT = { initial: 2_000, min: 25, max: 8_000 } as const /** A page slower than this halves the row limit. */ const SLOW_PAGE_MS = 30_000 -/** A page faster than this doubles the row limit, widening its scan window with it. */ +/** + * A page faster than this doubles the row limit, widening its scan window with it, but never back + * to a size that timed out. + */ const FAST_PAGE_MS = SLOW_PAGE_MS / 4 -/** Each page is followed by a pause as long as the page, up to this, to leave the primary headroom. */ +/** + * Each page, committed or timed out, is followed by a pause as long as the page, up to this, to + * leave the primary headroom. + */ const MAX_PAGE_PAUSE_MS = 5_000 const LOCK_RETRY_BUDGET_MS = 60_000 @@ -135,6 +141,8 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { let batches = 0 let mutated = 0 let rowLimit: number = ROW_LIMIT.initial + /** The largest limit the run may still try: half of the smallest limit that timed out. */ + let ceiling: number = ROW_LIMIT.max for (;;) { /** Timed around the whole call, so the synchronous-replication wait at commit counts. */ const pageStartedAt = performance.now() @@ -145,8 +153,10 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { if (!(error instanceof PageMutationTimeout)) throw error /** The timed-out page rolled back with its cursor, so it is retried with fewer rows. */ if (rowLimit <= ROW_LIMIT.min) throw error.timeout - rowLimit = halve(rowLimit) + ceiling = halve(rowLimit) + rowLimit = ceiling logger.warn('Search retirement page timed out; retrying with fewer rows', { rowLimit }) + await sleep(Math.min(performance.now() - pageStartedAt, MAX_PAGE_PAUSE_MS)) continue } if (page.done) break @@ -167,7 +177,7 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { rowLimit, }) } else if (pageMs < FAST_PAGE_MS) { - rowLimit = Math.min(ROW_LIMIT.max, rowLimit * 2) + rowLimit = Math.min(ceiling, rowLimit * 2) } if (batches % 10 === 0) { logger.info('Search embedding retirement progress', { diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index 0b966d14874..d9c531b50f3 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -48,8 +48,9 @@ also widens the scan window, so sparse stretches are not crawled in small window Materialized SQL pages keep the IDs inside PostgreSQL; the migration process receives only a cursor and a validation result. Each page uses a two-minute statement timeout and a one-second lock timeout. If a page's mutating statement exceeds the statement timeout, the page rolls back with its -cursor and is retried with half the row limit; one that still times out at 25 rows fails the -migration. Any other statement timeout fails the migration at once, because a smaller page cannot +cursor and is retried with half the row limit after the usual pause. From then on, fast pages grow +the limit only up to that halved size, so a size that timed out is never tried again. A page that +still times out at 25 rows fails the migration. Any other statement timeout fails the migration at once, because a smaller page cannot speed it up. The completion rechecks, which walk every captured KB once, run with a 30-minute timeout. Brief lock timeouts retry the rolled-back page with bounded backoff for up to one minute. Other errors, or exhausted lock retries, fail the migration without a completion receipt.