diff --git a/packages/db/drizzle.config.ts b/packages/db/drizzle.config.ts index aead0f23910..42d5b878d2e 100644 --- a/packages/db/drizzle.config.ts +++ b/packages/db/drizzle.config.ts @@ -16,5 +16,9 @@ export default { url: process.env.DATABASE_URL!, }, /** Runner-owned journals and resumable cleanup cursors must survive a development schema push. */ - tablesFilter: ['!script_migrations', '!search_embedding_cleanup_progress'], + tablesFilter: [ + '!script_migrations', + '!search_embedding_cleanup_progress', + '!search_embedding_cleanup_targets', + ], } satisfies Config 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 0ae423d754a..7e8f7f34269 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,6 @@ import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' import { maintainSearchRetirementMigration } from '@sim/db/script-migrations/0028_maintain_search_retirement' -import { runScriptMigrations } from '@sim/db/script-migrations/index' +import { runScriptMigrations, scriptMigrations } 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' @@ -18,6 +18,8 @@ describe('retiring dormant Search embeddings', () => { await admin.unsafe(`CREATE SCHEMA "${schema}"`) sql = postgres(readTestDatabaseUrl(), { max: 1, + /** Fixtures alternate legacy and upgraded table layouts on this connection. */ + prepare: false, connection: { search_path: schema }, onnotice: () => undefined, }) @@ -46,6 +48,7 @@ describe('retiring dormant Search embeddings', () => { beforeEach(async () => { 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_targets` 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)` @@ -68,20 +71,25 @@ describe('retiring dormant Search embeddings', () => { return receipts.length === 1 } - it('does nothing without Search data and fails an ambiguous target without changing data', async () => { + it('does nothing without Search 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` - 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 ).toBe(false) }) - it('cascades only Search chunks and preserves configuration and ordinary documents', async () => { + it('retires every Search KB while preserving configuration and ordinary documents', async () => { + await sql`INSERT INTO knowledge_base VALUES ('second-search', true), ('empty-search', true)` + await sql`INSERT INTO document (id, knowledge_base_id, processing_queue_token) + VALUES ('second-search-doc', 'second-search', 'old-dispatch')` + await sql`INSERT INTO embedding VALUES + ('00000', 'second-search', 'second-search-doc'), ('zz-last', 'second-search', 'second-search-doc')` + await sql`INSERT INTO embedding_search (id) VALUES ('00000'), ('zz-last')` + await sql`INSERT INTO embedding_keyword_search VALUES ('00000'), ('zz-last')` + await sql`INSERT INTO embedding_keyword_tin VALUES ('00000'), ('zz-last')` + await sql`INSERT INTO embedding_secret_provenance VALUES ('00000'), ('zz-last')` expect(await pass()).toBe(true) expect( ( @@ -103,7 +111,14 @@ describe('retiring dormant Search embeddings', () => { ]) { expect((await sql.unsafe(`SELECT count(*)::int AS n FROM ${table}`))[0].n).toBe(501) } - expect((await sql`SELECT count(*)::int AS n FROM knowledge_base`)[0].n).toBe(2) + expect((await sql`SELECT count(*)::int AS n FROM knowledge_base`)[0].n).toBe(4) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501) + expect( + ( + await sql`SELECT user_excluded, enabled, processing_queue_token FROM document + WHERE id = 'second-search-doc'` + )[0] + ).toEqual({ user_excluded: true, enabled: false, processing_queue_token: null }) expect( ( await sql`SELECT user_excluded, processing_queue_token FROM document WHERE id = 'ordinary-doc'` @@ -112,9 +127,134 @@ describe('retiring dormant Search embeddings', () => { expect(await pass()).toBe(true) }) - it('rolls back failed pages and resumes the frozen target, retiring documents inserted behind the cursor', async () => { + it.each([ + { phase: 'documents', legacy: 'present', otherSearch: true }, + { phase: 'embeddings', legacy: 'present', otherSearch: true }, + { phase: 'done', legacy: 'present', otherSearch: true }, + { phase: 'done', legacy: 'deleted', otherSearch: true }, + { phase: 'done', legacy: 'ordinary', otherSearch: true }, + { phase: 'done', legacy: 'deleted', otherSearch: false }, + ] as const)( + 'upgrades a legacy $phase checkpoint with a $legacy KB (other Search KBs: $otherSearch)', + async ({ phase, legacy, otherSearch }) => { + if (otherSearch) { + await sql`INSERT INTO knowledge_base VALUES ('second-search', true)` + await sql`INSERT INTO document (id, knowledge_base_id) + VALUES ('aaa-second-doc', 'second-search')` + await sql`INSERT INTO embedding VALUES ('00000', 'second-search', 'aaa-second-doc')` + await sql`INSERT INTO embedding_search (id) VALUES ('00000')` + } + await sql`UPDATE document SET user_excluded = true, enabled = false, + processing_queue_token = NULL WHERE knowledge_base_id = 'search'` + 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, reindexed_through text NOT NULL DEFAULT '', + vacuumed_tables integer NOT NULL DEFAULT 0)` + await sql`INSERT INTO search_embedding_cleanup_progress + VALUES (1, 'search', ${phase}, ${phase === 'documents' ? 'search-doc' : '00050'}, 'legacy_hnsw_idx', 6)` + await sql`CREATE INDEX legacy_hnsw_idx ON embedding_search USING hnsw (vector public.vector_l2_ops)` + const [original] = await sql`SELECT to_regclass('legacy_hnsw_idx')::oid AS oid` + await sql`CREATE TABLE script_migrations (name text PRIMARY KEY, applied_at timestamptz DEFAULT now())` + if (phase === 'embeddings') { + await sql`DELETE FROM embedding WHERE knowledge_base_id = 'search' AND id <= '00050'` + } + if (phase === 'done') { + await sql`DELETE FROM embedding WHERE knowledge_base_id = 'search'` + await sql`INSERT INTO script_migrations (name) + VALUES ('0027_retire_search_embeddings'), ('0028_maintain_search_retirement')` + } + if (legacy === 'deleted') { + await sql`DELETE FROM document WHERE knowledge_base_id = 'search'` + await sql`DELETE FROM knowledge_base WHERE id = 'search'` + } else if (legacy === 'ordinary') { + await sql`UPDATE knowledge_base SET is_search_index = false WHERE id = 'search'` + await sql`UPDATE document SET user_excluded = false, enabled = true, + processing_queue_token = 'keep-dispatch' WHERE id = 'search-doc'` + await sql`INSERT INTO embedding VALUES ('new-ordinary-chunk', 'search', 'search-doc')` + await sql`INSERT INTO embedding_search (id) VALUES ('new-ordinary-chunk')` + } + const migrations = scriptMigrations.filter((migration) => + [ + '0027_retire_search_embeddings', + '0028_maintain_search_retirement', + '0029_retire_all_search_embeddings', + ].includes(migration.name) + ) + try { + await runScriptMigrations(sql, migrations) + const preserved = legacy === 'ordinary' ? 502 : 501 + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(preserved) + expect((await sql`SELECT count(*)::int AS n FROM embedding_search`)[0].n).toBe(preserved) + if (otherSearch) { + expect( + (await sql`SELECT user_excluded, enabled FROM document WHERE id = 'aaa-second-doc'`)[0] + ).toEqual({ user_excluded: true, enabled: false }) + } + if (legacy === 'ordinary') { + expect( + ( + await sql`SELECT user_excluded, enabled, processing_queue_token FROM document WHERE id = 'search-doc'` + )[0] + ).toEqual({ + user_excluded: false, + enabled: true, + processing_queue_token: 'keep-dispatch', + }) + } + const [rebuilt] = await sql`SELECT to_regclass('legacy_hnsw_idx')::oid AS oid` + if (otherSearch) expect(rebuilt.oid).not.toBe(original.oid) + else expect(rebuilt.oid).toBe(original.oid) + expect( + await sql`SELECT name FROM script_migrations WHERE name = '0029_retire_all_search_embeddings'` + ).toHaveLength(1) + await runScriptMigrations(sql, migrations) + expect((await sql`SELECT to_regclass('legacy_hnsw_idx')::oid AS oid`)[0].oid).toBe( + rebuilt.oid + ) + } finally { + await sql`DROP INDEX legacy_hnsw_idx` + } + } + ) + + it.each(['embeddings', 'done'] as const)( + 'rejects changed markers outside the current page when resuming %s', + async (phase) => { + await sql`INSERT INTO knowledge_base VALUES ('empty-search', true)` + if (phase === 'done') { + await pass() + await sql`DELETE FROM script_migrations` + await sql`UPDATE knowledge_base SET is_search_index = false WHERE id = 'empty-search'` + } + await sql`CREATE FUNCTION change_empty_search_marker() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + UPDATE knowledge_base SET is_search_index = false WHERE id = 'empty-search'; + RETURN NULL; + END $$` + await sql`CREATE TRIGGER change_empty_search_marker AFTER DELETE ON embedding + FOR EACH STATEMENT EXECUTE FUNCTION change_empty_search_marker()` + try { + await expect(pass()).rejects.toThrow('no longer a Search knowledge base') + 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 change_empty_search_marker ON embedding` + await sql`UPDATE knowledge_base SET is_search_index = true WHERE id = 'empty-search'` + expect(await pass()).toBe(true) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501) + } finally { + await sql`DROP TRIGGER IF EXISTS change_empty_search_marker ON embedding` + await sql`DROP FUNCTION change_empty_search_marker()` + } + } + ) + + it('rolls back failed pages, rechecks every target marker, and resumes the frozen scope', async () => { await sql`INSERT INTO embedding SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 26002) i` + await sql`INSERT INTO knowledge_base VALUES ('second-search', true)` + await sql`INSERT INTO document (id, knowledge_base_id) VALUES ('second-search-doc', 'second-search')` + 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')` await expect(pass()).rejects.toThrow() @@ -123,6 +263,10 @@ describe('retiring dormant Search embeddings', () => { expect(remaining.n).toBeGreaterThan(501) expect(remaining.n).toBeLessThan(26002) expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0) + 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 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) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts index 99965cffcfc..1ebb2a837ec 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -2,7 +2,7 @@ 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 postgres, { type Sql } from 'postgres' +import postgres, { type Sql, type TransactionSql } from 'postgres' const logger = createLogger('RetireSearchEmbeddings') const BATCH_SIZE = 25_000 @@ -15,13 +15,13 @@ interface Progress { } /** - * Retires the sole legacy Search KB after the move to live Search. Each page commits with its + * Retires a frozen set of legacy Search KBs after the move to live Search. Each page commits with its * durable cursor. Other KBs are traversed without modification. The runner owns bookkeeping, as it owns `script_migrations`. */ export const retireSearchEmbeddingsMigration: ScriptMigration = { name: '0027_retire_search_embeddings', async up(sql) { - const knowledgeBaseId = await sql.begin(async (tx) => { + const hasTargets = await sql.begin('isolation level repeatable read', 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 ( @@ -29,30 +29,50 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { 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` + const [existing] = await tx` + SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress WHERE id = 1 FOR UPDATE` + const [snapshot] = + await tx`SELECT to_regclass('search_embedding_cleanup_targets') AS relation` + if (snapshot.relation) return Boolean(existing) + if (existing && existing.phase !== 'done') { + const [target] = await tx`SELECT id FROM knowledge_base + WHERE id = ${existing.knowledge_base_id} AND is_search_index FOR SHARE` + if (!target) throw new Error('Cleanup target is no longer a Search knowledge base') } - const [progress] = await tx< - Progress[] - >`SELECT * FROM search_embedding_cleanup_progress WHERE id = 1` - return progress.knowledge_base_id + + /** Creating the snapshot and resetting a legacy cursor commit atomically, once. */ + await tx`CREATE TABLE search_embedding_cleanup_targets (knowledge_base_id text PRIMARY KEY)` + let afterId = '' + for (;;) { + const [page] = await tx<{ after_id: string | null }[]>` + 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} + RETURNING knowledge_base_id + ) SELECT max(knowledge_base_id) AS after_id FROM targets` + if (page.after_id === null) break + afterId = page.after_id + } + const [first] = await tx<{ knowledge_base_id: string }[]>` + SELECT knowledge_base_id FROM search_embedding_cleanup_targets ORDER BY knowledge_base_id LIMIT 1` + if (!first) return false + await tx`ANALYZE search_embedding_cleanup_targets` + 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` + await tx`INSERT INTO search_embedding_cleanup_progress (id, knowledge_base_id, phase, after_id) + VALUES (1, ${first.knowledge_base_id}, 'documents', '') + ON CONFLICT (id) DO UPDATE SET knowledge_base_id = EXCLUDED.knowledge_base_id, + phase = 'documents', after_id = '', + reindexed_through = '', vacuumed_tables = 0` + return true }) - if (knowledgeBaseId === null) return + if (!hasTargets) return const startedAt = Date.now() let batches = 0 - while (!(await retirePage(sql, knowledgeBaseId))) { + while (!(await retirePage(sql))) { batches++ if (batches % 10 === 0) { logger.info('Search embedding retirement progress', { @@ -61,14 +81,14 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { }) } } - logger.info('Selected Search knowledge base embeddings retired', { + logger.info('Selected Search knowledge bases retired', { batches, elapsedMs: Date.now() - startedAt, }) }, } -async function retirePage(sql: Sql, knowledgeBaseId: string): Promise { +async function retirePage(sql: Sql): Promise { return retryOnLockTimeout( () => sql.begin(async (tx) => { @@ -76,26 +96,38 @@ async function retirePage(sql: Sql, knowledgeBaseId: string): Promise { 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') + if (progress.phase === 'done') { + await validateTargetMarkers(tx) + return true } - 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 [page] = await tx<{ after_id: string | null }[]>` + const [page] = await tx<{ after_id: string | null; invalid_target: boolean }[]>` WITH source_page AS MATERIALIZED ( - SELECT id FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE} + SELECT id, knowledge_base_id FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_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 + ), 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) + ORDER BY kb.id FOR SHARE OF kb + ), invalid_target AS MATERIALIZED ( + SELECT p.id FROM target_page p + LEFT JOIN locked_targets kb ON kb.id = p.knowledge_base_id + WHERE kb.id IS NULL OR NOT kb.is_search_index LIMIT 1 ), 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} + FROM target_page p + WHERE d.id = p.id AND d.knowledge_base_id = p.knowledge_base_id + 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 FROM source_page` + ) 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) { await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` return false @@ -105,37 +137,53 @@ async function retirePage(sql: Sql, knowledgeBaseId: string): Promise { } /** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */ - const [page] = await tx<{ after_id: string | null; unretired: boolean }[]>` + const [page] = await tx< + { after_id: string | null; unretired: boolean; invalid_target: boolean }[] + >` WITH source_page AS MATERIALIZED ( - SELECT 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 ${BATCH_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 + ), 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) + ORDER BY kb.id FOR SHARE OF kb + ), invalid_target AS MATERIALIZED ( + SELECT p.id FROM target_page p + LEFT JOIN locked_targets kb ON kb.id = p.knowledge_base_id + WHERE kb.id IS NULL OR NOT kb.is_search_index LIMIT 1 ), 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 + SELECT p.id FROM target_page p + JOIN document d ON d.id = p.document_id + WHERE NOT d.user_excluded OR d.knowledge_base_id <> p.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` + 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) + throw new Error('Cleanup target is no longer a Search knowledge base') 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 + 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 } - const [remaining] = - await tx`SELECT id FROM embedding WHERE knowledge_base_id = ${knowledgeBaseId} LIMIT 1` + 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 } + await validateTargetMarkers(tx) await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` return true } @@ -154,6 +202,27 @@ async function retirePage(sql: Sql, knowledgeBaseId: string): Promise { ) } +/** Validate even empty or fully scanned KBs, holding marker locks until completion commits. */ +async function validateTargetMarkers(tx: TransactionSql): Promise { + let afterId = '' + for (;;) { + 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} + ), 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) + ORDER BY kb.id FOR SHARE OF kb + ) SELECT max(p.knowledge_base_id) AS after_id, + coalesce(bool_or(kb.id IS NULL OR NOT kb.is_search_index), false) AS invalid_target + FROM target_page p LEFT JOIN locked_targets kb ON kb.id = p.knowledge_base_id` + if (page.invalid_target) throw new Error('Cleanup target is no longer a Search knowledge base') + if (page.after_id === null) return + afterId = page.after_id + } +} + /** The standalone entry resumes the deployment cursor and journals only a completed retirement. */ if (import.meta.main) { const url = resolveMigrationDatabaseUrl() @@ -161,13 +230,10 @@ if (import.meta.main) { const sql = postgres(url, { max: 1, max_lifetime: null, onnotice: () => undefined }) try { const { runScriptMigrations } = await import('@sim/db/script-migrations/index') - const { maintainSearchRetirementMigration } = await import( - '@sim/db/script-migrations/0028_maintain_search_retirement' + const { retireAllSearchEmbeddingsMigration } = await import( + '@sim/db/script-migrations/0029_retire_all_search_embeddings' ) - await runScriptMigrations(sql, [ - retireSearchEmbeddingsMigration, - maintainSearchRetirementMigration, - ]) + await runScriptMigrations(sql, [retireAllSearchEmbeddingsMigration]) } finally { await sql.end() } diff --git a/packages/db/script-migrations/0029_retire_all_search_embeddings.ts b/packages/db/script-migrations/0029_retire_all_search_embeddings.ts new file mode 100644 index 00000000000..5cfa5b6b582 --- /dev/null +++ b/packages/db/script-migrations/0029_retire_all_search_embeddings.ts @@ -0,0 +1,13 @@ +import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' +import { maintainSearchRetirementMigration } from '@sim/db/script-migrations/0028_maintain_search_retirement' +import type { ScriptMigration } from '@sim/db/script-migrations/types' + +/** Supersedes single-KB retirement receipts so every deployment receives the expanded cleanup. */ +export const retireAllSearchEmbeddingsMigration: ScriptMigration = { + name: '0029_retire_all_search_embeddings', + supersedes: [retireSearchEmbeddingsMigration.name, maintainSearchRetirementMigration.name], + async up(sql) { + await retireSearchEmbeddingsMigration.up(sql) + await maintainSearchRetirementMigration.up(sql) + }, +} diff --git a/packages/db/script-migrations/index.ts b/packages/db/script-migrations/index.ts index 9f98514f7bd..ec10731c9a3 100644 --- a/packages/db/script-migrations/index.ts +++ b/packages/db/script-migrations/index.ts @@ -10,8 +10,7 @@ import { projectionAclSkipUnfilledMigration } from '@sim/db/script-migrations/00 import { knowledgeProjectionAsyncMigration } from '@sim/db/script-migrations/0024_knowledge_projection_async' 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 { retireAllSearchEmbeddingsMigration } from '@sim/db/script-migrations/0029_retire_all_search_embeddings' import type { Sql } from 'postgres' import { backfillTableOrderKeys } from './0001_backfill_table_order_keys' import { backfillPausedBillingAttribution } from './0002_backfill_paused_billing_attribution' @@ -60,8 +59,8 @@ export const scriptMigrations: readonly ScriptMigration[] = [ scopeKeywordProjectionsMigration, /** 0026 installs the schema guard every table row write takes before it validates. */ userTableSchemaForWriteMigration, - retireSearchEmbeddingsMigration, - maintainSearchRetirementMigration, + /** 0029 expands single-KB retirement to every saved Search target and completes maintenance. */ + retireAllSearchEmbeddingsMigration, ] /** diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index 16c249b95ea..d55cb4594e2 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -1,11 +1,11 @@ -# Retiring one legacy Search index +# Retiring legacy Search indexes -`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 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. +`0029_retire_all_search_embeddings` runs through the existing script-migration registry without +cleanup flags. It supersedes the single-KB retirement and maintenance entries (`0027`/`0028`), +including databases that already recorded either receipt. It snapshots every knowledge base whose +persisted `is_search_index` marker is true. No Search KB is a completed no-op. Once saved, the +snapshot stays fixed across retries even if another Search KB is created. Ordinary KBs and the +selected KBs' live source/credential configuration, document metadata, and backing files are preserved. ## Deployment and execution @@ -18,11 +18,17 @@ also honor the indexed-search gate. 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. +`0029_retire_all_search_embeddings` and its superseded names in `script_migrations` only after all +selected KBs have no remaining chunks or unretired documents and index maintenance finishes. +On upgrading a legacy single-KB checkpoint, the snapshot and cursor reset commit atomically. The +scan starts at the beginning once so it includes other KBs behind the old cursor; previous deletes +remain committed. Maintenance checkpoints also reset once because the expanded cleanup creates new +dead entries. A completed legacy checkpoint does not require its former KB to still exist or remain +Search-marked; the new snapshot selects current Search KBs and preserves any KB now marked ordinary. +An unfinished legacy checkpoint still requires its target to remain Search-marked. Subsequent retries +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 @@ -35,7 +41,7 @@ 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 +retirement and maintenance through the successor migration and 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 @@ -43,37 +49,45 @@ migration runner already does for its session advisory lock and settings. `DATAB 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 -records only completion. `db:push` excludes the progress table from schema diffing. +The runner-owned `search_embedding_cleanup_targets` table stores the frozen KB set, populated in +bounded SQL pages within one repeatable-read transaction. The existing `search_embedding_cleanup_progress` +row stores the shared phase and ID cursor; its legacy `knowledge_base_id` remains an informational +anchor, not the full deletion scope. Page mutations and cursor advancement commit together. +The one-off migration journal records only completion. `db:push` excludes both bookkeeping tables +from schema diffing. The documents phase fences queued and in-flight processing by marking only target documents excluded/disabled and clearing their dispatch stamps. The embeddings phase deletes only target chunks, and refuses a page whose target document was not retired or has inconsistent ownership. The existing foreign keys cascade to vector/keyword projections and chunk provenance. Both phases -walk the primary key in bounded pages; unrelated rows are read only as IDs and are never updated. -This avoids sorting a whole KB or repeatedly scanning earlier pages when no suitable composite -cleanup index exists. Resumption continues the saved scan, including across pages containing only +walk the primary key once for the entire captured set in bounded pages; ordinary rows are never +updated. Each page locks and rechecks the Search markers for its target KBs before mutation, and +fails atomically if any target changed to an ordinary KB. This avoids a separate full-table scan +per KB or sorting a whole KB when no suitable composite cleanup index exists. Resumption continues the saved scan, including across pages containing only unrelated rows. Before completion, the cleanup checks for unretired documents and remaining chunks -behind either cursor and restarts the affected phase if needed. Keep target writers stopped and -do not change its marker during the pass. +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 +do not change their Search markers during the pass. Inspect progress with: ```sql SELECT * FROM search_embedding_cleanup_progress; SELECT name, applied_at FROM script_migrations -WHERE name IN ('0027_retire_search_embeddings', '0028_maintain_search_retirement'); +WHERE name IN ('0027_retire_search_embeddings', '0028_maintain_search_retirement', + '0029_retire_all_search_embeddings'); ``` -After completion, check the selected KB has no `embedding` rows, verify its live Search, and verify +After completion, check the selected KBs have no `embedding` rows, verify their live Search, and verify ordinary KB retrieval. A zero-row absence check may still scan index entries; use an appropriate timeout. The cleanup is destructive and not reversible by flipping the search flag. Re-enabling indexed Search requires deliberately restoring document eligibility and fully resyncing its sources. ## Storage maintenance -After deletion, `0028` runs `REINDEX INDEX CONCURRENTLY` on each HNSW index of `embedding_search`, +After deletion, `0029` invokes the existing maintenance implementation to run `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.