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..b1e54720bd7 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,31 @@ 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')` + /** + * 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') + 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 +321,225 @@ 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)` + /** 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); + 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) + /** + * 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 + * 2,667 with i % 3 <> 0 among 1003-5002. Equality proves no row was mutated 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 + /** `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] + .n + ).toBe(0) + expect( + ( + 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` + 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) + + 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 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) + 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) + 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. */ + 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 1ebb2a837ec..62554665582 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -2,15 +2,82 @@ 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 { sleep } from '@sim/utils/helpers' import postgres, { type Sql, type TransactionSql } from 'postgres' const logger = createLogger('RetireSearchEmbeddings') -const BATCH_SIZE = 25_000 +/** 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, + * not its scan, is what can outrun the statement timeout. + */ +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, but never back + * to a size that timed out. + */ +const FAST_PAGE_MS = SLOW_PAGE_MS / 4 +/** + * 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 +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 + /** 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 } @@ -48,7 +115,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,42 +139,111 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { const startedAt = Date.now() let batches = 0 - while (!(await retirePage(sql))) { + 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() + 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 (rowLimit <= ROW_LIMIT.min) throw error.timeout + 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 + const pageMs = performance.now() - pageStartedAt batches++ + 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 (pageMs < FAST_PAGE_MS) { + rowLimit = Math.min(ceiling, 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, }) }, } -async function retirePage(sql: Sql): Promise { +/** + * 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. + */ +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` 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, + ...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 true + return result({ done: 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 pageMutation(tx< + { after_id: string; mutated: number; 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 ${scanLimit} ), 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 +261,38 @@ 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) { + RETURNING d.id + ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, + (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 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 false + return result({ afterId: page.after_id, mutated: page.mutated }) } /** 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 }[] + const [page] = await pageMutation(tx< + { + after_id: string + 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 ${BATCH_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 + 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 +309,43 @@ 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) + RETURNING e.id + ) SELECT CASE WHEN e.limited THEN e.last_target ELSE max(p.id) END AS after_id, + (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`) + 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'` + 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 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 false + return result({ afterId: '', transition: 'embeddings_rescan' }) } await validateTargetMarkers(tx) await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` - return true + return result({ done: true }) } await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return false + return result({ afterId: page.after_id, mutated: page.mutated }) }), { budgetMs: LOCK_RETRY_BUDGET_MS, @@ -209,7 +366,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..d9c531b50f3 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -30,12 +30,35 @@ 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 -ten pages. The deployment job retains its five-hour overall timeout; it is not a runtime estimate. +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. 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, 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 +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 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. + +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 @@ -68,7 +91,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: