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 80ef1608a5c..0ae423d754a 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -1,6 +1,8 @@ import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' -import { ScriptMigrationDeferred } from '@sim/db/script-migrations/types' +import { maintainSearchRetirementMigration } from '@sim/db/script-migrations/0028_maintain_search_retirement' +import { runScriptMigrations } from '@sim/db/script-migrations/index' import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure' +import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' import postgres, { type Sql } from 'postgres' import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest' @@ -27,7 +29,9 @@ describe('retiring dormant Search embeddings', () => { await sql`CREATE TABLE embedding ( id text PRIMARY KEY, knowledge_base_id text REFERENCES knowledge_base(id), document_id text REFERENCES document(id))` - await sql`CREATE TABLE embedding_search (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` + await sql`CREATE TABLE embedding_search ( + id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE, + vector public.vector(3) NOT NULL DEFAULT '[1,2,3]')` await sql`CREATE TABLE embedding_keyword_search (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` await sql`CREATE TABLE embedding_keyword_tin (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` await sql`CREATE TABLE embedding_secret_provenance (embedding_id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` @@ -43,6 +47,7 @@ describe('retiring dormant Search embeddings', () => { await sql`TRUNCATE knowledge_base, document, embedding, embedding_search, embedding_keyword_search, embedding_keyword_tin, embedding_secret_provenance` await sql`DROP TABLE IF EXISTS search_embedding_cleanup_progress` + await sql`DROP TABLE IF EXISTS script_migrations` await sql`INSERT INTO knowledge_base VALUES ('search', true), ('ordinary', false)` await sql`INSERT INTO document (id, knowledge_base_id, processing_queue_token) VALUES ('search-doc', 'search', 'old-dispatch'), ('ordinary-doc', 'ordinary', 'keep-dispatch')` @@ -50,27 +55,26 @@ describe('retiring dormant Search embeddings', () => { SELECT lpad(i::text, 5, '0'), CASE WHEN i % 2 = 0 THEN 'search' ELSE 'ordinary' END, CASE WHEN i % 2 = 0 THEN 'search-doc' ELSE 'ordinary-doc' END FROM generate_series(1, 1002) i` - await sql`INSERT INTO embedding_search SELECT id FROM embedding` + await sql`INSERT INTO embedding_search (id) SELECT id FROM embedding` await sql`INSERT INTO embedding_keyword_search SELECT id FROM embedding` await sql`INSERT INTO embedding_keyword_tin SELECT id FROM embedding` await sql`INSERT INTO embedding_secret_provenance SELECT id FROM embedding` }) async function pass() { - try { - await retireSearchEmbeddingsMigration.up(sql) - return true - } catch (error) { - if (error instanceof ScriptMigrationDeferred) return false - throw error - } + await runScriptMigrations(sql, [retireSearchEmbeddingsMigration]) + const receipts = await sql`SELECT name FROM script_migrations + WHERE name = '0027_retire_search_embeddings'` + return receipts.length === 1 } - it('does nothing without Search data and defers an ambiguous target without changing data', async () => { + it('does nothing without Search data and fails an ambiguous target without changing data', async () => { await sql`UPDATE knowledge_base SET is_search_index = false` expect(await pass()).toBe(true) + await sql`DELETE FROM script_migrations` await sql`UPDATE knowledge_base SET is_search_index = true` - expect(await pass()).toBe(false) + await expect(pass()).rejects.toThrow('ambiguous') + expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0) expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(1002) expect( (await sql`SELECT user_excluded FROM document WHERE id = 'search-doc'`)[0].user_excluded @@ -109,15 +113,22 @@ describe('retiring dormant Search embeddings', () => { }) it('rolls back failed pages and resumes the frozen target, retiring documents inserted behind the cursor', async () => { + await sql`INSERT INTO embedding + SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 26002) i` await sql`CREATE TABLE deletion_blocker (id text REFERENCES embedding(id))` - await sql`INSERT INTO deletion_blocker VALUES ('00002')` + await sql`INSERT INTO deletion_blocker VALUES ('26002')` 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 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(1002) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(remaining.n) 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')` await sql`INSERT INTO knowledge_base VALUES ('other-search', true)` await sql`INSERT INTO document (id, knowledge_base_id, processing_queue_token) VALUES ('other-search-doc', 'other-search', 'keep-dispatch')` @@ -135,15 +146,193 @@ describe('retiring dormant Search embeddings', () => { ).toEqual({ user_excluded: false, processing_queue_token: 'keep-dispatch' }) }) - it('defers at the page budget and resumes without skipping remaining chunks', async () => { + it('finishes beyond the former page budget and journals completion in one invocation', async () => { + await sql`INSERT INTO document (id, knowledge_base_id) + SELECT 'bulk-doc-' || i::text, 'search' FROM generate_series(1, 51002) i` await sql`INSERT INTO embedding SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 51002) i` - expect(await pass()).toBe(false) - const [remaining] = - await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'` - expect(remaining.n).toBeGreaterThan(0) - expect(remaining.n).toBeLessThan(50501) - expect(await pass()).toBe(true) + await runScriptMigrations(sql, [ + retireSearchEmbeddingsMigration, + maintainSearchRetirementMigration, + ]) + expect(await sql`SELECT name FROM script_migrations`).toHaveLength(2) expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501) + 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) }, 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') + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(1002) + expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0) + await sql`UPDATE embedding SET document_id = 'search-doc' WHERE id = '00002'` + expect(await pass()).toBe(true) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501) + }) + + it('retries a briefly locked document page and completes in the same invocation', async () => { + const blocker = postgres(readTestDatabaseUrl(), { + max: 1, + connection: { search_path: schema }, + onnotice: () => undefined, + }) + let signalLocked!: () => void + const locked = new Promise((resolve) => { + signalLocked = resolve + }) + let release!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + const holding = blocker.begin(async (tx) => { + await tx`SELECT id FROM document WHERE id = 'search-doc' FOR UPDATE` + signalLocked() + await released + }) + try { + await locked + const result = pass().then( + (complete) => ({ complete, error: undefined }), + (error: unknown) => ({ complete: false, error }) + ) + await sleep(1_500) + release() + await holding + expect(await result).toEqual({ complete: true, error: undefined }) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0] + .n + ).toBe(0) + } finally { + release() + await holding + await blocker.end() + } + }) + + it('maintains an already-retired index and resumes failed vacuum bookkeeping without rebuilding it again', async () => { + await pass() + await sql`CREATE INDEX retirement_hnsw_idx ON embedding_search USING hnsw (vector public.vector_l2_ops)` + const [original] = await sql`SELECT to_regclass('retirement_hnsw_idx')::oid AS oid` + await sql`CREATE FUNCTION interrupt_maintenance_checkpoint() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + IF NEW.vacuumed_tables = 2 THEN RAISE EXCEPTION 'interrupt maintenance checkpoint'; END IF; + RETURN NEW; + END $$` + await sql`CREATE TRIGGER interrupt_maintenance_checkpoint BEFORE UPDATE ON search_embedding_cleanup_progress + FOR EACH ROW EXECUTE FUNCTION interrupt_maintenance_checkpoint()` + try { + await expect(runScriptMigrations(sql, [maintainSearchRetirementMigration])).rejects.toThrow( + 'interrupt maintenance checkpoint' + ) + const [rebuilt] = + await sql`SELECT c.oid, i.indisvalid FROM pg_class c JOIN pg_index i ON i.indexrelid = c.oid + WHERE c.oid = to_regclass('retirement_hnsw_idx')` + expect(rebuilt.oid).not.toBe(original.oid) + expect(rebuilt.indisvalid).toBe(true) + expect( + await sql`SELECT reindexed_through, vacuumed_tables FROM search_embedding_cleanup_progress` + ).toEqual([{ reindexed_through: 'retirement_hnsw_idx', vacuumed_tables: 1 }]) + expect( + await sql`SELECT name FROM script_migrations WHERE name = '0028_maintain_search_retirement'` + ).toHaveLength(0) + await sql`DROP TRIGGER interrupt_maintenance_checkpoint ON search_embedding_cleanup_progress` + await runScriptMigrations(sql, [maintainSearchRetirementMigration]) + expect((await sql`SELECT to_regclass('retirement_hnsw_idx')::oid AS oid`)[0].oid).toBe( + rebuilt.oid + ) + expect( + await sql`SELECT name FROM script_migrations WHERE name = '0028_maintain_search_retirement'` + ).toHaveLength(1) + expect( + ( + await sql`SELECT reltuples::int AS n FROM pg_class WHERE oid = to_regclass('embedding_search')` + )[0].n + ).toBe(501) + await runScriptMigrations(sql, [maintainSearchRetirementMigration]) + expect((await sql`SELECT to_regclass('retirement_hnsw_idx')::oid AS oid`)[0].oid).toBe( + rebuilt.oid + ) + } finally { + await sql`DROP TRIGGER IF EXISTS interrupt_maintenance_checkpoint ON search_embedding_cleanup_progress` + await sql`DROP FUNCTION interrupt_maintenance_checkpoint()` + await sql`DROP INDEX retirement_hnsw_idx` + } + }) + + it('recovers a canceled concurrent rebuild and completes cleanup plus maintenance in one rerun', async () => { + await sql`CREATE INDEX retirement_hnsw_idx ON embedding_search USING hnsw (vector public.vector_l2_ops)` + const blocker = postgres(readTestDatabaseUrl(), { + max: 1, + connection: { search_path: schema }, + onnotice: () => undefined, + }) + let signalLocked!: () => void + const locked = new Promise((resolve) => { + signalLocked = resolve + }) + let release!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + const holding = blocker.begin(async (tx) => { + await tx`UPDATE embedding_search SET vector = '[3,2,1]' WHERE id = '00001'` + signalLocked() + await released + }) + const [{ pid }] = await sql`SELECT pg_backend_pid() AS pid` + const migrations = [retireSearchEmbeddingsMigration, maintainSearchRetirementMigration] + let outcome: Promise | undefined + try { + await locked + outcome = runScriptMigrations(sql, migrations).then( + () => undefined, + (error: unknown) => error + ) + let waiting = false + for (let attempt = 0; attempt < 200; attempt++) { + const [progress] = + await admin`SELECT phase FROM pg_stat_progress_create_index WHERE pid = ${pid}` + if (progress?.phase === 'waiting for writers before build') { + waiting = true + break + } + await sleep(25) + } + expect(waiting).toBe(true) + await admin`SELECT pg_cancel_backend(${pid})` + expect(await outcome).toMatchObject({ code: '57014' }) + release() + await holding + const leftovers = + await sql`SELECT c.relname FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid + WHERE i.indrelid = to_regclass('embedding_search') AND NOT i.indisvalid` + expect(leftovers.length).toBeGreaterThan(0) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0] + .n + ).toBe(0) + await runScriptMigrations(sql, migrations) + expect( + await sql`SELECT c.relname FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid + WHERE i.indrelid = to_regclass('embedding_search') AND NOT i.indisvalid` + ).toHaveLength(0) + expect((await sql`SELECT count(*)::int AS n FROM embedding_search`)[0].n).toBe(501) + expect( + await sql`SELECT name FROM script_migrations WHERE name = '0028_maintain_search_retirement'` + ).toHaveLength(1) + } finally { + release() + await holding + await admin`SELECT pg_cancel_backend(${pid})` + await outcome + await blocker.end() + await sql`DROP INDEX retirement_hnsw_idx` + } + }) }) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts index 5afec02dc11..99965cffcfc 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -1,13 +1,12 @@ import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url' -import { type ScriptMigration, ScriptMigrationDeferred } from '@sim/db/script-migrations/types' +import type { ScriptMigration } from '@sim/db/script-migrations/types' +import { retryOnLockTimeout } from '@sim/db/scripts/lock-timeout-retry' import { createLogger } from '@sim/logger' -import { sleep } from '@sim/utils/helpers' import postgres, { type Sql } from 'postgres' const logger = createLogger('RetireSearchEmbeddings') -const BATCH_SIZE = 500 -const MAX_BATCHES = 100 -const RUN_BUDGET_MS = 60_000 +const BATCH_SIZE = 25_000 +const LOCK_RETRY_BUDGET_MS = 60_000 interface Progress { knowledge_base_id: string @@ -22,119 +21,153 @@ interface Progress { export const retireSearchEmbeddingsMigration: ScriptMigration = { name: '0027_retire_search_embeddings', async up(sql) { - await sql`CREATE TABLE IF NOT EXISTS 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 - )` - const existing = await sql< - Progress[] - >`SELECT * FROM search_embedding_cleanup_progress WHERE id = 1` - if (existing.length === 0) { - const targets = await sql<{ id: string }[]>` - SELECT id FROM knowledge_base WHERE is_search_index LIMIT 2` - if (targets.length === 0) return - if (targets.length > 1) { - throw new ScriptMigrationDeferred( - 'Multiple Search knowledge bases found; cleanup target is ambiguous' - ) + const knowledgeBaseId = await sql.begin(async (tx) => { + await tx`SET LOCAL statement_timeout = '120s'` + await tx`SET LOCAL lock_timeout = '1s'` + await tx`CREATE TABLE IF NOT EXISTS 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 + )` + const existing = await tx< + Progress[] + >`SELECT * FROM search_embedding_cleanup_progress WHERE id = 1` + if (existing.length === 0) { + const targets = await tx<{ id: string }[]>` + SELECT id FROM knowledge_base WHERE is_search_index LIMIT 2` + if (targets.length === 0) return null + if (targets.length > 1) { + throw new Error('Multiple Search knowledge bases found; cleanup target is ambiguous') + } + await tx`INSERT INTO search_embedding_cleanup_progress (id, knowledge_base_id, phase, after_id) + VALUES (1, ${targets[0].id}, 'documents', '') + ON CONFLICT (id) DO NOTHING` } - await sql`INSERT INTO search_embedding_cleanup_progress VALUES (1, ${targets[0].id}, 'documents', '') - ON CONFLICT (id) DO NOTHING` - } - const [progress] = await sql< - Progress[] - >`SELECT * FROM search_embedding_cleanup_progress WHERE id = 1` - const knowledgeBaseId = progress.knowledge_base_id + const [progress] = await tx< + Progress[] + >`SELECT * FROM search_embedding_cleanup_progress WHERE id = 1` + return progress.knowledge_base_id + }) + if (knowledgeBaseId === null) return - const deadline = Date.now() + RUN_BUDGET_MS - for (let batch = 0; batch < MAX_BATCHES && Date.now() < deadline; batch++) { - if (await retirePage(sql, knowledgeBaseId)) { - logger.info('Selected Search knowledge base embeddings retired') - return + const startedAt = Date.now() + let batches = 0 + while (!(await retirePage(sql, knowledgeBaseId))) { + batches++ + if (batches % 10 === 0) { + logger.info('Search embedding retirement progress', { + batches, + elapsedMs: Date.now() - startedAt, + }) } - await sleep(100) } - throw new ScriptMigrationDeferred( - 'Cleanup page budget reached; rerun to resume the saved cursor' - ) + logger.info('Selected Search knowledge base embeddings retired', { + batches, + elapsedMs: Date.now() - startedAt, + }) }, } async function retirePage(sql: Sql, knowledgeBaseId: string): Promise { - return sql.begin(async (tx) => { - await tx`SET LOCAL statement_timeout = '15s'` - 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.knowledge_base_id !== knowledgeBaseId) { - throw new Error('Cannot change the selected Search KB during a cleanup') - } - const [target] = await tx`SELECT id FROM knowledge_base - WHERE id = ${knowledgeBaseId} AND is_search_index FOR SHARE` - if (!target) throw new Error('Cleanup target is no longer a Search knowledge base') - if (progress.phase === 'done') return true + return retryOnLockTimeout( + () => + sql.begin(async (tx) => { + 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.knowledge_base_id !== knowledgeBaseId) { + throw new Error('Cannot change the selected Search KB during a cleanup') + } + const [target] = await tx`SELECT id FROM knowledge_base + WHERE id = ${knowledgeBaseId} AND is_search_index FOR SHARE` + if (!target) throw new Error('Cleanup target is no longer a Search knowledge base') + if (progress.phase === 'done') return true - if (progress.phase === 'documents') { - const rows = await tx<{ id: string }[]>` - SELECT id FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE}` - if (rows.length === 0) { - await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` - return false - } - await tx`UPDATE document - SET user_excluded = true, enabled = false, processing_queue_token = NULL, - processing_queued_at = NULL, processing_deferred_until = NULL - WHERE id IN ${tx(rows.map((row) => row.id))} AND knowledge_base_id = ${knowledgeBaseId} - AND (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)` - await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${rows[rows.length - 1].id} WHERE id = 1` - return false - } + if (progress.phase === 'documents') { + const [page] = await tx<{ after_id: string | null }[]>` + WITH source_page AS MATERIALIZED ( + SELECT id FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE} + ), retired AS ( + UPDATE document d + SET user_excluded = true, enabled = false, processing_queue_token = NULL, + processing_queued_at = NULL, processing_deferred_until = NULL + FROM source_page p WHERE d.id = p.id AND d.knowledge_base_id = ${knowledgeBaseId} + 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 FROM source_page` + if (page.after_id === null) { + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` + return false + } + await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` + return false + } - const rows = await tx<{ id: string }[]>` - SELECT id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE}` - if (rows.length === 0) { - /** A late insert may sort behind either UUID cursor; completion must recheck the target. */ - const [unretired] = - await tx`SELECT id FROM document WHERE knowledge_base_id = ${knowledgeBaseId} - AND (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 - } - const [remaining] = - await tx`SELECT id FROM embedding WHERE knowledge_base_id = ${knowledgeBaseId} LIMIT 1` - if (remaining) { - await tx`UPDATE search_embedding_cleanup_progress SET after_id = '' WHERE id = 1` + /** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */ + const [page] = await tx<{ after_id: string | null; unretired: boolean }[]>` + WITH source_page AS MATERIALIZED ( + SELECT id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE} + ), unretired AS MATERIALIZED ( + SELECT e.id FROM source_page p JOIN embedding e ON e.id = p.id + JOIN document d ON d.id = e.document_id + WHERE e.knowledge_base_id = ${knowledgeBaseId} + AND (NOT d.user_excluded OR d.knowledge_base_id <> e.knowledge_base_id) LIMIT 1 + ), deleted AS ( + DELETE FROM embedding e USING source_page p + WHERE e.id = p.id AND e.knowledge_base_id = ${knowledgeBaseId} + AND NOT EXISTS (SELECT 1 FROM unretired) + ) SELECT max(id) AS after_id, EXISTS (SELECT 1 FROM unretired) AS unretired FROM source_page` + 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. */ + const [unretired] = + await tx`SELECT id FROM document WHERE knowledge_base_id = ${knowledgeBaseId} + AND (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 + } + const [remaining] = + await tx`SELECT id FROM embedding WHERE knowledge_base_id = ${knowledgeBaseId} LIMIT 1` + if (remaining) { + await tx`UPDATE search_embedding_cleanup_progress SET after_id = '' WHERE id = 1` + return false + } + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` + return true + } + await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` return false - } - await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` - return true + }), + { + budgetMs: LOCK_RETRY_BUDGET_MS, + backoff: { baseMs: 1_000, maxMs: 5_000 }, + onRetry: ({ attempt, delayMs }) => + logger.warn('Search retirement page locked; retrying', { + attempt, + retryInMs: Math.round(delayMs), + }), } - const ids = rows.map((row) => row.id) - const [unretired] = await tx` - SELECT e.id FROM embedding e JOIN document d ON d.id = e.document_id - WHERE e.id IN ${tx(ids)} AND e.knowledge_base_id = ${knowledgeBaseId} - AND (NOT d.user_excluded OR d.knowledge_base_id <> e.knowledge_base_id) LIMIT 1` - if (unretired) - throw new Error('Search content changed after retirement; stop writers before resuming') - /** Foreign keys cascade to vector/keyword projections and private chunk provenance. */ - await tx`DELETE FROM embedding WHERE id IN ${tx(ids)} AND knowledge_base_id = ${knowledgeBaseId}` - await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${rows[rows.length - 1].id} WHERE id = 1` - return false - }) + ) } -/** The standalone entry uses the deployment journal and leaves budget-deferred work unrecorded. */ +/** The standalone entry resumes the deployment cursor and journals only a completed retirement. */ if (import.meta.main) { const url = resolveMigrationDatabaseUrl() if (!url) throw new Error('DATABASE_URL is required for Search retirement') const sql = postgres(url, { max: 1, max_lifetime: null, onnotice: () => undefined }) try { const { runScriptMigrations } = await import('@sim/db/script-migrations/index') - await runScriptMigrations(sql, [retireSearchEmbeddingsMigration]) + const { maintainSearchRetirementMigration } = await import( + '@sim/db/script-migrations/0028_maintain_search_retirement' + ) + await runScriptMigrations(sql, [ + retireSearchEmbeddingsMigration, + maintainSearchRetirementMigration, + ]) } finally { await sql.end() } diff --git a/packages/db/script-migrations/0028_maintain_search_retirement.ts b/packages/db/script-migrations/0028_maintain_search_retirement.ts new file mode 100644 index 00000000000..747d08612c0 --- /dev/null +++ b/packages/db/script-migrations/0028_maintain_search_retirement.ts @@ -0,0 +1,113 @@ +import type { ScriptMigration } from '@sim/db/script-migrations/types' +import { createLogger } from '@sim/logger' +import { escapeRegExp } from '@sim/utils/string' +import type { Sql } from 'postgres' + +const logger = createLogger('SearchRetirementMaintenance') +const MAINTENANCE_LOCK = 'sim:search-retirement-maintenance' +const VACUUM_TABLES = [ + 'embedding_search', + 'embedding_keyword_search', + 'embedding_keyword_tin', + 'embedding_secret_provenance', + 'embedding', + 'document', +] as const + +/** + * Runs after retirement, including on databases that already journaled 0027. Concurrent rebuilds + * and vacuum must run outside transactions; checkpoints follow each successful operation. + */ +export const maintainSearchRetirementMigration: ScriptMigration = { + name: '0028_maintain_search_retirement', + async up(sql) { + const [table] = await sql`SELECT to_regclass('search_embedding_cleanup_progress') AS relation` + if (!table.relation) return + const [retirement] = await sql`SELECT phase FROM search_embedding_cleanup_progress WHERE id = 1` + if (!retirement) return + if (retirement.phase !== 'done') + throw new Error('Search retirement must finish before maintenance') + + const [{ locked }] = + await sql`SELECT pg_try_advisory_lock(hashtextextended(${MAINTENANCE_LOCK}, 0)) AS locked` + if (!locked) throw new Error('Search retirement maintenance is already running') + const [settings] = await sql`SELECT current_setting('statement_timeout') AS statement_timeout, + current_setting('lock_timeout') AS lock_timeout` + try { + await sql.begin(async (tx) => { + await tx`SET LOCAL lock_timeout = '1s'` + await tx`SET LOCAL statement_timeout = '120s'` + await tx`ALTER TABLE search_embedding_cleanup_progress + ADD COLUMN IF NOT EXISTS reindexed_through text NOT NULL DEFAULT '', + ADD COLUMN IF NOT EXISTS vacuumed_tables integer NOT NULL DEFAULT 0` + }) + /** Concurrent maintenance waits for old snapshots without blocking ordinary table writes. */ + await sql`SET statement_timeout = 0` + await sql`SET lock_timeout = 0` + for (;;) { + const [index] = await sql<{ name: string; qualified_name: string }[]>` + SELECT c.relname::text AS name, format('%I.%I', n.nspname, c.relname) AS qualified_name + FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid + JOIN pg_namespace n ON n.oid = c.relnamespace JOIN pg_am am ON am.oid = c.relam + WHERE i.indrelid = to_regclass('embedding_search') AND am.amname = 'hnsw' + AND c.relname::text > (SELECT reindexed_through FROM search_embedding_cleanup_progress WHERE id = 1) + AND c.relname !~ '_cc(new|old)[0-9]*$' + ORDER BY c.relname::text LIMIT 1` + if (!index) break + await removeInterruptedRebuilds(sql, index.name) + const startedAt = Date.now() + logger.info('Rebuilding retired Search vector index', { index: index.name }) + await sql.unsafe(`REINDEX INDEX CONCURRENTLY ${index.qualified_name}`) + await sql`UPDATE search_embedding_cleanup_progress SET reindexed_through = ${index.name} WHERE id = 1` + logger.info('Search vector index rebuilt', { + index: index.name, + elapsedMs: Date.now() - startedAt, + }) + } + + const [progress] = await sql<{ vacuumed_tables: number }[]>` + SELECT vacuumed_tables FROM search_embedding_cleanup_progress WHERE id = 1` + for (let step = progress.vacuumed_tables; step < VACUUM_TABLES.length; step++) { + const tableName = VACUUM_TABLES[step] + const [relation] = await sql<{ qualified_name: string; can_maintain: boolean }[]>` + SELECT format('%I.%I', n.nspname, c.relname) AS qualified_name, + CASE WHEN current_setting('server_version_num')::int >= 170000 + THEN has_table_privilege(c.oid, 'MAINTAIN') + ELSE pg_has_role(c.relowner, 'USAGE') END AS can_maintain + FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE c.oid = to_regclass(${tableName})` + if (!relation?.can_maintain) throw new Error(`Cannot vacuum retirement table ${tableName}`) + const startedAt = Date.now() + logger.info('Vacuuming retired Search storage', { table: tableName }) + await sql.unsafe(`VACUUM (ANALYZE, TRUNCATE FALSE) ${relation.qualified_name}`) + await sql`UPDATE search_embedding_cleanup_progress SET vacuumed_tables = ${step + 1} WHERE id = 1` + logger.info('Search storage vacuumed', { + table: tableName, + elapsedMs: Date.now() - startedAt, + }) + } + } finally { + try { + await sql`SELECT set_config('statement_timeout', ${settings.statement_timeout}, false), + set_config('lock_timeout', ${settings.lock_timeout}, false)` + } finally { + await sql`SELECT pg_advisory_unlock(hashtextextended(${MAINTENANCE_LOCK}, 0))` + } + } + }, +} + +/** PostgreSQL leaves invalid _ccnew/_ccold siblings if a concurrent rebuild is interrupted. */ +async function removeInterruptedRebuilds(sql: Sql, indexName: string): Promise { + for (;;) { + const [leftover] = await sql<{ qualified_name: string }[]>` + SELECT format('%I.%I', n.nspname, c.relname) AS qualified_name + FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid + JOIN pg_namespace n ON n.oid = c.relnamespace JOIN pg_am am ON am.oid = c.relam + WHERE i.indrelid = to_regclass('embedding_search') AND am.amname = 'hnsw' + AND NOT i.indisvalid AND c.relname ~ ${`^${escapeRegExp(indexName)}_cc(new|old)[0-9]*$`} + ORDER BY c.relname LIMIT 1` + if (!leftover) return + await sql.unsafe(`DROP INDEX CONCURRENTLY ${leftover.qualified_name}`) + } +} diff --git a/packages/db/script-migrations/index.ts b/packages/db/script-migrations/index.ts index 77c9cc5869d..9f98514f7bd 100644 --- a/packages/db/script-migrations/index.ts +++ b/packages/db/script-migrations/index.ts @@ -11,6 +11,7 @@ import { knowledgeProjectionAsyncMigration } from '@sim/db/script-migrations/002 import { scopeKeywordProjectionsMigration } from '@sim/db/script-migrations/0025_scope_keyword_projections' import { userTableSchemaForWriteMigration } from '@sim/db/script-migrations/0026_user_table_schema_for_write' import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' +import { maintainSearchRetirementMigration } from '@sim/db/script-migrations/0028_maintain_search_retirement' import type { Sql } from 'postgres' import { backfillTableOrderKeys } from './0001_backfill_table_order_keys' import { backfillPausedBillingAttribution } from './0002_backfill_paused_billing_attribution' @@ -60,6 +61,7 @@ export const scriptMigrations: readonly ScriptMigration[] = [ /** 0026 installs the schema guard every table row write takes before it validates. */ userTableSchemaForWriteMigration, retireSearchEmbeddingsMigration, + maintainSearchRetirementMigration, ] /** diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index e012eab8286..16c249b95ea 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -3,7 +3,7 @@ `0027_retire_search_embeddings` runs through the existing script-migration registry without cleanup flags. On its first invocation it discovers the sole knowledge base whose persisted `is_search_index` marker is true and saves that target. No Search KB is a completed no-op; -multiple Search KBs defer without changing content because the target is ambiguous. Once saved, +multiple Search KBs fail without changing content because the target is ambiguous. Once saved, the target stays fixed even if another Search KB is created. Ordinary KBs and the target's live source/credential configuration, document metadata, and backing files are preserved. @@ -16,16 +16,32 @@ Search again.** The cleanup does not inspect this flag. Live source setup may st KB for configuration; it does not index content. Document uploads, dispatch and queued processing also honor the indexed-search gate. -The ordinary migration runner starts the cleanup automatically. For subsequent maintenance passes, -run `bun run packages/db/script-migrations/0027_retire_search_embeddings.ts` with the writer supplied -through the normal `MIGRATION_DATABASE_URL`/`DATABASE_URL` configuration. Repeat until -`script_migrations` records `0027_retire_search_embeddings`. A budget deferral exits normally without -recording completion; it is **not** a completed purge. +The ordinary migration runner starts the cleanup automatically and continues until it is complete +in the same deployment. There is no page-count or one-minute deferral. A successful run records +`0027_retire_search_embeddings` in `script_migrations` only after the target has no remaining chunks +or unretired documents. Previously deferred runs resume their existing saved target, phase and cursor. +The next registered migration, `0028_maintain_search_retirement`, rebuilds the HNSW indexes and vacuums +the affected tables before deployment continues. Its separate receipt also makes maintenance run +where the original retirement was already completed before this upgrade. -Each invocation handles at most 100 pages of 500 IDs, with a 60-second budget between pages, -a 15-second statement timeout, a one-second lock timeout and 100 ms pacing. One in-flight statement -can finish after the run budget. Keep one maintenance worker and monitor primary latency, WAL, -replica lag and available disk. Lock/query failures stop the invocation; rerun after resolving them. +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. + +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 +through the normal `MIGRATION_DATABASE_URL`/`DATABASE_URL` configuration. Completed pages remain +committed and the failed page is retried from its saved cursor. The standalone command runs both +retirement and maintenance through the same journal. Keep one maintenance worker and monitor primary +latency, WAL, replica lag and available disk. + +Both entry points require a direct or session-pooled PostgreSQL connection, as the deployment +migration runner already does for its session advisory lock and settings. `DATABASE_URL` is a valid +fallback only when it provides that session affinity. PgBouncer transaction pooling is unsupported; +reserving a postgres.js client connection does not pin a backend through a transaction pooler. The runner-owned `search_embedding_cleanup_progress` table stores the selected KB, phase and ID cursor. Page mutations and cursor advancement commit together. The one-off migration journal @@ -45,8 +61,9 @@ do not change its marker during the pass. Inspect progress with: ```sql -SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress; -SELECT name, applied_at FROM script_migrations WHERE name = '0027_retire_search_embeddings'; +SELECT * FROM search_embedding_cleanup_progress; +SELECT name, applied_at FROM script_migrations +WHERE name IN ('0027_retire_search_embeddings', '0028_maintain_search_retirement'); ``` After completion, check the selected KB has no `embedding` rows, verify its live Search, and verify @@ -56,13 +73,23 @@ indexed Search requires deliberately restoring document eligibility and fully re ## Storage maintenance -No vacuum, reindex or table rewrite runs inside this backfill. Plan maintenance separately on the -writer, outside transactions. Ordinary `VACUUM (ANALYZE, TRUNCATE FALSE)` on affected source, -projection and document tables makes dead space reusable and refreshes planner statistics. It -generally does not shrink allocated files. `VACUUM FULL` takes an exclusive lock and requires extra -space; use a separately planned maintenance window or a provider-supported online rewrite for -physical shrinkage. HNSW may benefit from `REINDEX INDEX CONCURRENTLY` before vacuum, with enough -temporary disk and WAL capacity. See the [PostgreSQL vacuum documentation](https://www.postgresql.org/docs/17/sql-vacuum.html) +After deletion, `0028` runs `REINDEX INDEX CONCURRENTLY` on each HNSW index of `embedding_search`, +then `VACUUM (ANALYZE, TRUNCATE FALSE)` on the vector and keyword projections, chunk provenance, +embeddings, and documents. These operations execute sequentially outside transactions. Rebuilds +keep ordinary reads and writes available and require temporary index space and WAL capacity. +They wait for older transactions and can dominate total runtime; the five-hour job deadline still +applies. PostgreSQL's `pg_stat_progress_create_index` and `pg_stat_progress_vacuum` expose progress. + +The existing progress row gains `reindexed_through` and `vacuumed_tables` checkpoints. Completed +indexes and tables are skipped on retry; interruption between an operation and its checkpoint may +repeat that one operation. Invalid `_ccnew`/`_ccold` siblings from an interrupted concurrent rebuild +are removed concurrently before retrying their original index. One advisory lock serializes the +maintenance worker. Do not run other index maintenance on these tables at the same time. + +Ordinary vacuum makes dead space reusable and refreshes planner statistics; it generally does not +shrink table files. This migration does not run `VACUUM FULL` or rewrite tables. See the +[PostgreSQL vacuum documentation](https://www.postgresql.org/docs/17/sql-vacuum.html), +[concurrent reindex recovery](https://www.postgresql.org/docs/17/sql-reindex.html#SQL-REINDEX-CONCURRENTLY), and [pgvector maintenance guidance](https://github.com/pgvector/pgvector#vacuuming). Document counts are historical until a later document cleanup. Do not raw-delete documents or