From 6a050c8bcd8f0cbda4bcc866519374066053c652 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Thu, 1 Oct 2026 09:47:07 -0700 Subject: [PATCH] improvement(search): run Search retirement as an operator command instead of a deploy step Deploys no longer register the 0027-0029 Search retirement. A long cleanup in the deploy migration generated heavy WAL and stalled application writes, and held the release until it finished. The retirement is optional storage reclamation once live Search is on, so it now runs as a resumable operator command, paced by --pause-ratio and --max-rows, with --maintenance running the index rebuilds and vacuum and journaling completion. --- ...016_backfill_search_vectors.integration.ts | 7 + ...27_retire_search_embeddings.integration.ts | 43 ++- .../0027_retire_search_embeddings.ts | 246 +++++++++++------- .../0029_retire_all_search_embeddings.ts | 28 +- packages/db/script-migrations/index.ts | 7 +- .../search-embedding-retirement.md | 149 +++++++---- 6 files changed, 307 insertions(+), 173 deletions(-) diff --git a/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts b/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts index e3c92cf8955..5f2d238637c 100644 --- a/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts +++ b/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts @@ -396,6 +396,13 @@ describe('search projection upgrade in PostgreSQL', () => { AND s.vector_512 = subvector(e.embedding, 1, 512)::halfvec(512) AND k.content_tsv = e.content_tsv` expect(complete).toBe(1001) + /** Search retirement is an operator command; a deploy's full registry run never starts it. */ + expect( + ( + await sql`SELECT to_regclass('search_embedding_cleanup_progress') AS progress, + to_regclass('search_embedding_cleanup_targets') AS targets` + )[0] + ).toEqual({ progress: null, targets: null }) await runScriptMigrations(sql) await sql`DELETE FROM embedding WHERE id LIKE 'upgrade-%'` }, 60_000) 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 b1e54720bd7..34738fd2c78 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,10 @@ -import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' +import { + retireSearchEmbeddings, + retireSearchEmbeddingsMigration, +} from '@sim/db/script-migrations/0027_retire_search_embeddings' import { maintainSearchRetirementMigration } from '@sim/db/script-migrations/0028_maintain_search_retirement' -import { runScriptMigrations, scriptMigrations } from '@sim/db/script-migrations/index' +import { retireAllSearchEmbeddings } from '@sim/db/script-migrations/0029_retire_all_search_embeddings' +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' @@ -174,13 +178,7 @@ describe('retiring dormant Search embeddings', () => { 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) - ) + const migrations = [retireAllSearchEmbeddings()] try { await runScriptMigrations(sql, migrations) const preserved = legacy === 'ordinary' ? 502 : 501 @@ -710,4 +708,31 @@ describe('retiring dormant Search embeddings', () => { await sql`DROP INDEX retirement_hnsw_idx` } }) + + it('pauses after each page for the pause ratio times the page, so a manual run leaves the primary idle', async () => { + /** Every delete page takes about 100 ms; with a ratio of 3 the next page starts 300 ms after it ends. */ + await sql`CREATE TABLE delete_page_started (at timestamptz NOT NULL DEFAULT clock_timestamp())` + await sql`CREATE FUNCTION slow_delete_page() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN INSERT INTO delete_page_started DEFAULT VALUES; PERFORM pg_sleep(0.1); RETURN NULL; END $$` + await sql`CREATE TRIGGER slow_delete_page BEFORE DELETE ON embedding + FOR EACH STATEMENT EXECUTE FUNCTION slow_delete_page()` + try { + await retireSearchEmbeddings(sql, { pauseRatio: 3, maxRows: 200 }) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0] + .n + ).toBe(0) + const starts = ( + await sql<{ at: Date }[]>`SELECT at FROM delete_page_started ORDER BY at` + ).map(({ at }) => at.getTime()) + expect(starts.length).toBeGreaterThanOrEqual(3) + for (let i = 1; i < starts.length; i++) { + expect(starts[i] - starts[i - 1]).toBeGreaterThanOrEqual(380) + } + } finally { + await sql`DROP TRIGGER slow_delete_page ON embedding` + await sql`DROP FUNCTION slow_delete_page()` + await sql`DROP TABLE delete_page_started` + } + }, 60_000) }) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts index 62554665582..f58f5b14b01 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -1,3 +1,4 @@ +import { parseArgs } from 'node:util' import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url' import type { ScriptMigration } from '@sim/db/script-migrations/types' import { retryOnLockTimeout } from '@sim/db/scripts/lock-timeout-retry' @@ -18,7 +19,8 @@ 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. + * not its scan, is what can outrun the statement timeout. `maxRows` lowers the starting limit and + * the ceiling it may grow to. */ const ROW_LIMIT = { initial: 2_000, min: 25, max: 8_000 } as const /** A page slower than this halves the row limit. */ @@ -28,12 +30,23 @@ const SLOW_PAGE_MS = 30_000 * to a size that timed out. */ const FAST_PAGE_MS = SLOW_PAGE_MS / 4 +/** The longest pause after one page, however slow the page was. */ +const MAX_PAGE_PAUSE_MS = 60_000 +const LOCK_RETRY_BUDGET_MS = 60_000 + /** - * Each page, committed or timed out, is followed by a pause as long as the page, up to this, to - * leave the primary headroom. + * How hard one run pushes the primary. Each page, committed or timed out, is followed by a pause of + * `pauseRatio` times the page's duration (up to a minute), so the run is busy at most + * `1 / (1 + pauseRatio)` of the time. A page is timed through its commit, so a slow synchronous + * replica or a checkpoint stall lengthens the pause by the same factor. */ -const MAX_PAGE_PAUSE_MS = 5_000 -const LOCK_RETRY_BUDGET_MS = 60_000 +export interface RetirementPacing { + pauseRatio: number + /** The most rows one page may update or delete, from 25 to 8,000. */ + maxRows: number +} + +export const DEFAULT_RETIREMENT_PACING: RetirementPacing = { pauseRatio: 2, maxRows: 2_000 } type Phase = 'documents' | 'embeddings' | 'done' @@ -87,116 +100,132 @@ interface Progress { */ export const retireSearchEmbeddingsMigration: ScriptMigration = { name: '0027_retire_search_embeddings', - async up(sql) { - 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 ( + up: (sql) => retireSearchEmbeddings(sql), +} + +/** Retires every captured target, resuming the saved cursor, paced by `pacing`. */ +export async function retireSearchEmbeddings( + sql: Sql, + pacing: RetirementPacing = DEFAULT_RETIREMENT_PACING +): Promise { + if ( + !(pacing.pauseRatio >= 0) || + !Number.isInteger(pacing.maxRows) || + pacing.maxRows < ROW_LIMIT.min || + pacing.maxRows > ROW_LIMIT.max + ) { + throw new Error( + `Search retirement pacing needs a pause ratio of at least 0 and ${ROW_LIMIT.min}-${ROW_LIMIT.max} max rows` + ) + } + const pause = (pageMs: number) => sleep(Math.min(pageMs * pacing.pauseRatio, MAX_PAGE_PAUSE_MS)) + 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 ( 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` + 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 + 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') - } + if (!target) throw new Error('Cleanup target is no longer a Search knowledge base') + } - /** 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 }[]>` + /** 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 ${SCAN_PAGE_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 }[]>` + 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 + 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) + 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 (!hasTargets) return + return true + }) + if (!hasTargets) return - const startedAt = Date.now() - let batches = 0 - let mutated = 0 - let rowLimit: number = ROW_LIMIT.initial - /** The largest limit the run may still try: half of the smallest limit that timed out. */ - let ceiling: number = ROW_LIMIT.max - for (;;) { - /** Timed around the whole call, so the synchronous-replication wait at commit counts. */ - const pageStartedAt = performance.now() - 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)) + const startedAt = Date.now() + let batches = 0 + let mutated = 0 + let rowLimit = Math.min(ROW_LIMIT.initial, pacing.maxRows) + /** The largest limit the run may still try: half of the smallest limit that timed out. */ + let ceiling = pacing.maxRows + 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 pause(performance.now() - pageStartedAt) + 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, + }) } - logger.info('Selected Search knowledge bases retired', { - batches, - mutated, - elapsedMs: Date.now() - startedAt, - }) - }, + await pause(pageMs) + } + logger.info('Selected Search knowledge bases retired', { + batches, + mutated, + elapsedMs: Date.now() - startedAt, + }) } /** @@ -380,17 +409,36 @@ async function validateTargetMarkers(tx: TransactionSql): Promise { } } -/** The standalone entry resumes the deployment cursor and journals only a completed retirement. */ +/** + * The operator entry: resumes the saved cursor and, with `--maintenance`, also rebuilds the indexes, + * vacuums, and journals the completed cleanup. `--pause-ratio` and `--max-rows` set the pacing. + */ if (import.meta.main) { + const { values } = parseArgs({ + options: { + maintenance: { type: 'boolean', default: false }, + 'pause-ratio': { type: 'string' }, + 'max-rows': { type: 'string' }, + }, + }) + const pacing: RetirementPacing = { + pauseRatio: Number(values['pause-ratio'] ?? DEFAULT_RETIREMENT_PACING.pauseRatio), + maxRows: Number(values['max-rows'] ?? DEFAULT_RETIREMENT_PACING.maxRows), + } 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') - const { retireAllSearchEmbeddingsMigration } = await import( - '@sim/db/script-migrations/0029_retire_all_search_embeddings' - ) - await runScriptMigrations(sql, [retireAllSearchEmbeddingsMigration]) + if (values.maintenance) { + const { runScriptMigrations } = await import('@sim/db/script-migrations/index') + const { retireAllSearchEmbeddings } = await import( + '@sim/db/script-migrations/0029_retire_all_search_embeddings' + ) + await runScriptMigrations(sql, [retireAllSearchEmbeddings(pacing)]) + } else { + await retireSearchEmbeddings(sql, pacing) + logger.info('Search retirement pass finished; run with --maintenance off-peak to complete it') + } } 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 index 5cfa5b6b582..1e6d5263256 100644 --- a/packages/db/script-migrations/0029_retire_all_search_embeddings.ts +++ b/packages/db/script-migrations/0029_retire_all_search_embeddings.ts @@ -1,13 +1,23 @@ -import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' +import { + type RetirementPacing, + retireSearchEmbeddings, + 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) - }, +/** + * The complete operator-run cleanup: every Search KB's retirement, then index maintenance. It is not + * in the deploy registry; the 0027 entry runs it with `--maintenance` and journals it on success. It + * supersedes single-KB retirement receipts so every database receives the expanded cleanup. + */ +export function retireAllSearchEmbeddings(pacing?: RetirementPacing): ScriptMigration { + return { + name: '0029_retire_all_search_embeddings', + supersedes: [retireSearchEmbeddingsMigration.name, maintainSearchRetirementMigration.name], + async up(sql) { + await retireSearchEmbeddings(sql, pacing) + await maintainSearchRetirementMigration.up(sql) + }, + } } diff --git a/packages/db/script-migrations/index.ts b/packages/db/script-migrations/index.ts index ec10731c9a3..7700a8f80a3 100644 --- a/packages/db/script-migrations/index.ts +++ b/packages/db/script-migrations/index.ts @@ -10,7 +10,6 @@ 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 { 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' @@ -59,8 +58,10 @@ export const scriptMigrations: readonly ScriptMigration[] = [ scopeKeywordProjectionsMigration, /** 0026 installs the schema guard every table row write takes before it validates. */ userTableSchemaForWriteMigration, - /** 0029 expands single-KB retirement to every saved Search target and completes maintenance. */ - retireAllSearchEmbeddingsMigration, + /** + * Search retirement (0027–0029) is an operator-run maintenance command, not a deploy step: + * see `search-embedding-retirement.md`. + */ ] /** diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index d9c531b50f3..86d48536e5b 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -1,48 +1,92 @@ # Retiring legacy Search indexes -`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 - -The app and workers must already use live Search, and older indexing jobs must be drained before -this cleanup ships: deployment migrations run before the new app switches over. `SIM_SEARCH_LIVE=true` -(the default) makes `isIndexedOrgSearchEnabled()` false. **`SIM_SEARCH_LIVE=false` enables indexed -Search again.** The cleanup does not inspect this flag. Live source setup may still create a Search -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 and continues until it is complete -in the same deployment. There is no page-count or one-minute deferral. A successful run records -`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. +The retirement is an **operator-run maintenance command, not a deploy step**. Deploy migrations no +longer register it: a long cleanup inside the deploy migration generated heavy WAL and stalled +application writes, and it held the release until it finished. It is optional storage reclamation +once live Search is on, so it runs separately, paced, at a time the operator chooses. Self-hosted +operators can run the same command. + +`0029_retire_all_search_embeddings` snapshots every knowledge base whose persisted `is_search_index` +marker is true and supersedes the single-KB retirement and maintenance receipts (`0027`/`0028`), +including databases that already recorded either. No Search KB is a completed no-op. Once saved, the +snapshot stays fixed across runs even if another Search KB is created. Ordinary KBs and the selected +KBs' live source/credential configuration, document metadata, and backing files are preserved. + +## Before running + +The app and workers must already use live Search, and older indexing jobs must be drained. +`SIM_SEARCH_LIVE=true` (the default) makes `isIndexedOrgSearchEnabled()` false. **`SIM_SEARCH_LIVE=false` +enables indexed Search again.** The cleanup does not inspect this flag. Live source setup may still +create a Search KB for configuration; it does not index content. Document uploads, dispatch and +queued processing also honor the indexed-search gate. + +## Running it + +From the repository root, with the migration role's writer DSN on a **direct or session-pooled** +connection (the run holds a session advisory lock and session settings; PgBouncer transaction pooling +is unsupported, and reserving a postgres.js client does not pin a backend through it): + +```sh +# Retire documents and delete their chunks, resuming the saved cursor. Safe to stop and rerun. +MIGRATION_DATABASE_URL= bun run packages/db/script-migrations/0027_retire_search_embeddings.ts + +# Off-peak: finish any remaining retirement, rebuild the HNSW indexes, vacuum, and record completion. +MIGRATION_DATABASE_URL= bun run packages/db/script-migrations/0027_retire_search_embeddings.ts --maintenance +``` + +| Flag | Default | Effect | +| --- | --- | --- | +| `--pause-ratio N` | `2` | After each page, pause N × the page's duration (at most one minute), so the run is busy at most `1 / (1 + N)` of the time. Raise it to go gentler. | +| `--max-rows N` | `2000` | The most rows one page may update or delete (25–8,000). Lower it to make each page lighter. | +| `--maintenance` | off | After retirement, run the index rebuilds and vacuums and journal `0029` with its superseded names. | + +Run it as the migration role: maintenance needs `pg_maintain`, which the application roles lack. Run +it outside peak traffic, and run `--maintenance` in the quietest window you have: concurrent HNSW +rebuilds are long and write a lot of WAL (GitLab, for example, schedules automatic reindexing for +weekends). Keep one run at a time. + +**Pausing.** Ctrl-C is safe at any point. The in-flight page rolls back with its cursor, and an +interrupted concurrent rebuild's leftover index is removed on the next run. Rerun the same command to +resume; completed pages stay committed. + +**Watching.** Every ten pages the run 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. In PostgreSQL, watch for `checkpoint starting: wal` in quick succession, slow checkpoint +sync times, `canceling wait for synchronous replication`, and replica lag. If they appear, stop the +run and resume later with a higher `--pause-ratio` or lower `--max-rows`. + +```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', + '0029_retire_all_search_embeddings'); +``` + +A plain run does not journal anything; only a `--maintenance` run that finishes records `0029` and its +superseded names. 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. + +## How a run paces itself 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 +sequentially, and each is followed by a pause of `--pause-ratio` times its duration, up to one +minute. Because a page is timed through its commit, a slow synchronous replica or a checkpoint stall +lengthens the following pause by the same factor. 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 +The row limit starts at 2,000 rows, or `--max-rows` if lower. 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 +30 seconds halves the limit. A fast page, one under 7.5 seconds, doubles it up to `--max-rows`, 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 @@ -55,23 +99,6 @@ speed it up. The completion rechecks, which walk every captured KB once, run wit 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 -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 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 -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_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 @@ -110,12 +137,11 @@ indexed Search requires deliberately restoring document eligibility and fully re ## Storage maintenance -After deletion, `0029` invokes the existing maintenance implementation to run `REINDEX INDEX CONCURRENTLY` on each HNSW index of `embedding_search`, +With `--maintenance`, 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. -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. +They wait for older transactions and can dominate total runtime. 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 @@ -134,3 +160,20 @@ bucket objects: their application hard-delete path also enqueues identity-bound and applies accounting. Its ordinary scoped mode excludes retired documents, so a follow-up must explicitly support these rows while preserving those side effects. Do not delete source accounts, integration policies or permission grants used by live Search. + +## Why it runs this way + +Long data changes belong outside deploy migrations, in batches, throttled on database health, and +resumable from a cursor: + +- [strong_migrations: Backfilling data](https://github.com/ankane/strong_migrations#backfilling-data) +- [GitLab batched background migrations](https://docs.gitlab.com/development/database/batched_background_migrations/) + and [automatic reindexing](https://docs.gitlab.com/omnibus/settings/database/) +- [Shopify maintenance_tasks](https://github.com/Shopify/maintenance_tasks) +- [gh-ost throttling](https://github.com/github/gh-ost/blob/master/doc/throttle.md) and + [pt-online-schema-change](https://docs.percona.com/percona-toolkit/pt-online-schema-change.html) +- [Stripe: online migrations at scale](https://stripe.com/blog/online-migrations) +- PostgreSQL 17: [WAL configuration](https://www.postgresql.org/docs/17/wal-configuration.html), + [synchronous replication](https://www.postgresql.org/docs/17/warm-standby.html#SYNCHRONOUS-REPLICATION), + [replication statistics](https://www.postgresql.org/docs/17/monitoring-stats.html) +- [PlanetScale: the only scalable delete](https://planetscale.com/blog/the-only-scalable-delete)