diff --git a/.github/workflows/search-retirement.yml b/.github/workflows/search-retirement.yml new file mode 100644 index 00000000000..d4f4495c4f7 --- /dev/null +++ b/.github/workflows/search-retirement.yml @@ -0,0 +1,87 @@ +name: Search Retirement Cleanup + +# Background runner for the `0029_retire_all_search_embeddings` script migration. A deploy only +# advances the cleanup for a minute and defers; this job advances it between deploys in budgeted, +# resumable slices throttled on the database's commit latency, WAL rate and replication lag, and +# journals it once retirement and index maintenance finish. Every slice resumes the saved cursor +# and maintenance checkpoints, so a cancelled or failed run loses at most its in-flight page. +# Runbook: packages/db/script-migrations/search-embedding-retirement.md +# +# Pause: `gh workflow disable search-retirement.yml`, then cancel any in-progress run. + +on: + schedule: + # Retirement slice every hour; a slice never starts index maintenance. + - cron: '23 * * * *' + # Off-peak maintenance window: concurrent HNSW rebuilds and vacuums may start. + - cron: '7 3 * * *' + workflow_dispatch: + inputs: + environment: + description: Target environment + required: true + type: choice + options: + - production + - staging + budget_minutes: + description: Minutes after which no page or maintenance operation starts + required: false + default: '50' + maintenance: + description: Allow concurrent index rebuilds and vacuums in this run + type: boolean + required: false + default: false + +permissions: + contents: read + +jobs: + cleanup: + name: Advance Search retirement (${{ matrix.environment }}) + if: github.repository == 'simstudioai/sim' + runs-on: ${{ (vars.CI_PROVIDER == '' || vars.CI_PROVIDER == 'blacksmith') && 'blacksmith-4vcpu-ubuntu-2404' || 'ubuntu-latest' }} + strategy: + fail-fast: false + matrix: + environment: ${{ fromJSON(github.event_name == 'workflow_dispatch' && format('["{0}"]', inputs.environment) || '["production","staging"]') }} + # One slice per database at a time; a queued slice waits rather than cancelling a running one. + concurrency: + group: search-retirement-${{ matrix.environment }} + cancel-in-progress: false + # A rebuild that starts inside the budget runs to completion, so the job outlives the budget. + timeout-minutes: 300 + + steps: + - name: Checkout code + uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6 + + - name: Setup Bun + uses: oven-sh/setup-bun@0c5077e51419868618aeaa5fe8019c62421857d6 # v2 + with: + bun-version: 1.4.2 + + - name: Install dependencies + run: bun install --frozen-lockfile --ignore-scripts + + # Same secret mapping as migrations.yml: the migration role owns the tables, which the + # concurrent rebuilds and vacuums require, and can read replication lag. + - name: Run a cleanup slice + working-directory: ./packages/db + env: + DATABASE_URL: ${{ matrix.environment == 'production' && secrets.DATABASE_URL || matrix.environment == 'staging' && secrets.STAGING_DATABASE_URL || '' }} + MIGRATION_DATABASE_URL: ${{ matrix.environment == 'production' && secrets.MIGRATION_DATABASE_URL || matrix.environment == 'staging' && secrets.STAGING_MIGRATION_DATABASE_URL || '' }} + BUDGET_MINUTES: ${{ github.event_name == 'workflow_dispatch' && inputs.budget_minutes || github.event.schedule == '7 3 * * *' && '120' || '50' }} + MAINTENANCE: ${{ (github.event_name == 'workflow_dispatch' && inputs.maintenance) || github.event.schedule == '7 3 * * *' }} + run: | + set -euo pipefail + if [ -z "$DATABASE_URL" ]; then + echo "ERROR: no database URL secret resolved" >&2 + exit 1 + fi + args=(--budget-minutes "$BUDGET_MINUTES") + if [ "$MAINTENANCE" = "true" ]; then + args+=(--maintenance) + fi + bun run ./script-migrations/0029_retire_all_search_embeddings.ts "${args[@]}" 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..747362fe00f 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -1,12 +1,43 @@ -import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' +import { + type RetirementThrottle, + retireSearchEmbeddings, +} from '@sim/db/script-migrations/0027_retire_search_embeddings' import { maintainSearchRetirementMigration } from '@sim/db/script-migrations/0028_maintain_search_retirement' +import { searchRetirementSlice } from '@sim/db/script-migrations/0029_retire_all_search_embeddings' import { runScriptMigrations, scriptMigrations } from '@sim/db/script-migrations/index' +import { type ScriptMigration, ScriptMigrationDeferred } from '@sim/db/script-migrations/types' 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' +/** + * Pages back to back, for the suites about scope and recovery rather than pacing. A local database's + * small `max_wal_size` would otherwise pace even these tiny fixtures for seconds per page. + */ +const UNTHROTTLED: Partial = { + dutyRatio: 0, + walBudgetShare: Number.POSITIVE_INFINITY, +} + +/** The retirement alone, run to completion under its own journal name. */ +const retireToCompletion: ScriptMigration = { + name: '0027_retire_search_embeddings', + async up(sql) { + if ((await retireSearchEmbeddings(sql, { throttle: UNTHROTTLED })) === 'deferred') { + throw new ScriptMigrationDeferred('Another Search retirement runner holds the lock') + } + }, +} + +/** The background runner's slice without a budget: retirement and maintenance to completion. */ +const backgroundRun = searchRetirementSlice({ + budgetMs: Number.POSITIVE_INFINITY, + maintenance: true, + throttle: UNTHROTTLED, +}) + /** Proves destructive scope, cascading integrity, and atomic restart against real PostgreSQL. */ describe('retiring dormant Search embeddings', () => { const schema = `search_retirement_${generateId().replaceAll('-', '')}` @@ -65,7 +96,7 @@ describe('retiring dormant Search embeddings', () => { }) async function pass() { - await runScriptMigrations(sql, [retireSearchEmbeddingsMigration]) + await runScriptMigrations(sql, [retireToCompletion]) const receipts = await sql`SELECT name FROM script_migrations WHERE name = '0027_retire_search_embeddings'` return receipts.length === 1 @@ -181,11 +212,21 @@ describe('retiring dormant Search embeddings', () => { '0029_retire_all_search_embeddings', ].includes(migration.name) ) + const journaled = async () => + ( + await sql`SELECT name FROM script_migrations WHERE name = '0029_retire_all_search_embeddings'` + ).length === 1 try { + /** The deploy slice retires within its budget but leaves rebuilds to the background runner. */ 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) + expect((await sql`SELECT to_regclass('legacy_hnsw_idx')::oid AS oid`)[0].oid).toBe( + original.oid + ) + expect(await journaled()).toBe(!otherSearch) + await runScriptMigrations(sql, [backgroundRun]) if (otherSearch) { expect( (await sql`SELECT user_excluded, enabled FROM document WHERE id = 'aaa-second-doc'`)[0] @@ -205,9 +246,7 @@ describe('retiring dormant Search embeddings', () => { 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) + expect(await journaled()).toBe(true) await runScriptMigrations(sql, migrations) expect((await sql`SELECT to_regclass('legacy_hnsw_idx')::oid AS oid`)[0].oid).toBe( rebuilt.oid @@ -302,16 +341,13 @@ describe('retiring dormant Search embeddings', () => { ).toEqual({ user_excluded: false, processing_queue_token: 'keep-dispatch' }) }) - it('finishes beyond the former page budget and journals completion in one invocation', async () => { + it('finishes in one unbounded background invocation and journals completion', 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` - await runScriptMigrations(sql, [ - retireSearchEmbeddingsMigration, - maintainSearchRetirementMigration, - ]) - expect(await sql`SELECT name FROM script_migrations`).toHaveLength(2) + await runScriptMigrations(sql, [backgroundRun]) + expect(await sql`SELECT name FROM script_migrations`).toHaveLength(3) expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501) expect( ( @@ -661,7 +697,7 @@ describe('retiring dormant Search embeddings', () => { await released }) const [{ pid }] = await sql`SELECT pg_backend_pid() AS pid` - const migrations = [retireSearchEmbeddingsMigration, maintainSearchRetirementMigration] + const migrations = [retireToCompletion, maintainSearchRetirementMigration] let outcome: Promise | undefined try { await locked @@ -710,4 +746,279 @@ describe('retiring dormant Search embeddings', () => { await sql`DROP INDEX retirement_hnsw_idx` } }) + + /** A cleanup progress row in its pre-maintenance layout, so a test can attach triggers first. */ + async function createLegacyProgressTable() { + 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)` + } + + /** Records each chunk delete statement's time, rows and smallest ID. */ + async function recordDeleteStatements() { + await sql`CREATE TABLE delete_statement (at timestamptz NOT NULL, rows integer NOT NULL, min_id text)` + await sql`CREATE FUNCTION record_delete_statement() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + INSERT INTO delete_statement + SELECT clock_timestamp(), count(*), min(id) FROM removed HAVING count(*) > 0; + RETURN NULL; + END $$` + await sql`CREATE TRIGGER record_delete_statement AFTER DELETE ON embedding + REFERENCING OLD TABLE AS removed FOR EACH STATEMENT EXECUTE FUNCTION record_delete_statement()` + } + + async function dropDeleteStatements() { + await sql`DROP TRIGGER IF EXISTS record_delete_statement ON embedding` + await sql`DROP FUNCTION IF EXISTS record_delete_statement()` + await sql`DROP TABLE IF EXISTS delete_statement` + } + + async function deleteStatements() { + return sql<{ at: Date; rows: number; min_id: string }[]>` + SELECT at, rows, min_id FROM delete_statement ORDER BY at` + } + + it('defers the deploy slice within its budget, then the background runner resumes the cursor and journals', async () => { + await sql`INSERT INTO embedding + SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 40002) i` + await createLegacyProgressTable() + /** Every page commit waits, as a synchronous-replication round trip would. */ + await sql`CREATE FUNCTION wait_at_commit() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + PERFORM pg_sleep(0.2); + RETURN NULL; + END $$` + await sql`CREATE CONSTRAINT TRIGGER wait_at_commit AFTER UPDATE ON search_embedding_cleanup_progress + DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION wait_at_commit()` + await recordDeleteStatements() + const lock = postgres(readTestDatabaseUrl(), { max: 1, onnotice: () => undefined }) + try { + const budgetMs = 1_000 + const deploy = searchRetirementSlice({ budgetMs, maintenance: false }) + const startedAt = performance.now() + await runScriptMigrations(sql, [deploy]) + /** No page starts after the budget, so only the page in flight at the deadline may overrun it. */ + expect(performance.now() - startedAt).toBeLessThan(budgetMs + 2_000) + expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0) + const [saved] = await sql`SELECT phase, after_id FROM search_embedding_cleanup_progress` + expect(saved.phase).toBe('embeddings') + expect(saved.after_id > '').toBe(true) + const [{ remaining }] = await sql`SELECT count(*)::int AS remaining FROM embedding + WHERE knowledge_base_id = 'search'` + expect(remaining).toBeGreaterThan(0) + expect(remaining).toBeLessThan(39_501) + + /** While another runner holds the retirement lock, a deploy slice defers without a page. */ + await lock`SELECT pg_advisory_lock(hashtextextended('sim:search-retirement', 0))` + await runScriptMigrations(sql, [deploy]) + await lock`SELECT pg_advisory_unlock(hashtextextended('sim:search-retirement', 0))` + expect(await sql`SELECT phase, after_id FROM search_embedding_cleanup_progress`).toEqual([ + saved, + ]) + expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0) + + await sql`TRUNCATE delete_statement` + await runScriptMigrations(sql, [backgroundRun]) + /** The background run picks up at the saved cursor rather than rescanning from the start. */ + const resumed = await deleteStatements() + expect(resumed.length).toBeGreaterThan(0) + for (const statement of resumed) expect(statement.min_id > saved.after_id).toBe(true) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0] + .n + ).toBe(0) + expect( + (await sql`SELECT name FROM script_migrations ORDER BY name`).map((row) => row.name) + ).toEqual([ + '0027_retire_search_embeddings', + '0028_maintain_search_retirement', + '0029_retire_all_search_embeddings', + ]) + } finally { + await lock.end() + await dropDeleteStatements() + await sql`DROP TRIGGER IF EXISTS wait_at_commit ON search_embedding_cleanup_progress` + await sql`DROP FUNCTION wait_at_commit()` + } + }, 60_000) + + it('halves the page and backs off exponentially while commits are slow', async () => { + await sql`INSERT INTO embedding + SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 12002) i` + await createLegacyProgressTable() + await sql`CREATE SEQUENCE slow_commits` + /** The first two chunk pages wait at commit, as a stalled synchronous standby would make them. */ + await sql`CREATE FUNCTION slow_commit() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + IF NEW.phase = 'embeddings' AND NEW.after_id <> '' THEN + IF nextval('slow_commits') <= 2 THEN PERFORM pg_sleep(0.3); END IF; + END IF; + RETURN NULL; + END $$` + await sql`CREATE CONSTRAINT TRIGGER slow_commit AFTER UPDATE ON search_embedding_cleanup_progress + DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION slow_commit()` + await recordDeleteStatements() + const backoffMs = 1_000 + try { + expect( + await retireSearchEmbeddings(sql, { + throttle: { ...UNTHROTTLED, slowCommitMs: 150, backoffMs }, + }) + ).toBe('complete') + const statements = await deleteStatements() + /** Halved after each slow commit, then growing again once commits recover. */ + const rows = statements.slice(0, 4).map((statement) => statement.rows) + expect(rows).toEqual([rows[0], rows[0] / 2, rows[0] / 4, rows[0] / 2]) + const gap = (index: number) => + statements[index].at.getTime() - statements[index - 1].at.getTime() + /** Each gap holds the slow commit plus a back-off that doubles while commits stay slow. */ + expect(gap(1)).toBeGreaterThanOrEqual(300 + backoffMs) + expect(gap(2)).toBeGreaterThanOrEqual(300 + 2 * backoffMs) + expect(gap(3)).toBeLessThan(backoffMs) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0] + .n + ).toBe(0) + } finally { + await dropDeleteStatements() + await sql`DROP TRIGGER IF EXISTS slow_commit ON search_embedding_cleanup_progress` + await sql`DROP FUNCTION slow_commit()` + await sql`DROP SEQUENCE slow_commits` + } + }, 60_000) + + it('paces pages so the database WAL rate stays within the budget', async () => { + await sql`INSERT INTO embedding + SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 40002) i` + await sql`INSERT INTO embedding_search (id) SELECT id FROM embedding WHERE id > '01002'` + /** Where each chunk delete statement starts and ends, in time and in WAL. */ + await sql`CREATE TABLE delete_bound (at timestamptz NOT NULL, lsn pg_lsn NOT NULL)` + await sql`CREATE FUNCTION record_delete_bound() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + INSERT INTO delete_bound VALUES (clock_timestamp(), pg_current_wal_insert_lsn()); + RETURN NULL; + END $$` + await sql`CREATE TRIGGER record_delete_start BEFORE DELETE ON embedding + FOR EACH STATEMENT EXECUTE FUNCTION record_delete_bound()` + await sql`CREATE TRIGGER record_delete_end AFTER DELETE ON embedding + FOR EACH STATEMENT EXECUTE FUNCTION record_delete_bound()` + const [settings] = await sql<{ max_wal_bytes: number; checkpoint_ms: number }[]>` + SELECT pg_size_bytes(current_setting('max_wal_size'))::float8 AS max_wal_bytes, + (extract(epoch FROM current_setting('checkpoint_timeout')::interval) * 1000)::float8 AS checkpoint_ms` + /** 1 MB/s, far below what back-to-back pages write locally. */ + const budgetBytesPerMs = 1_000 + const walBudgetShare = (budgetBytesPerMs * settings.checkpoint_ms) / settings.max_wal_bytes + try { + expect( + await retireSearchEmbeddings(sql, { throttle: { dutyRatio: 0, walBudgetShare } }) + ).toBe('complete') + const bounds = await sql<{ at: Date; wal: number }[]>` + SELECT at, pg_wal_lsn_diff(lsn, '0/0')::float8 AS wal FROM delete_bound ORDER BY at` + /** Pairs of start and end, one per page; the last page's bounds have no successor. */ + const pages = [] + for (let index = 0; index + 1 < bounds.length; index += 2) { + pages.push({ + startedAt: bounds[index].at.getTime(), + wal: bounds[index + 1].wal - bounds[index].wal, + }) + } + expect(pages.length).toBeGreaterThan(2) + /** + * Every page but the last was followed by its pause, so the WAL those pages wrote, over the + * time from the first page's start to the last page's start, is the paced rate. + */ + const paced = pages.slice(0, -1) + const wal = paced.reduce((total, page) => total + page.wal, 0) + const elapsedMs = pages[pages.length - 1].startedAt - pages[0].startedAt + expect(wal).toBeGreaterThan(0) + expect(wal / elapsedMs).toBeLessThanOrEqual(budgetBytesPerMs * 1.05) + } finally { + await sql`DROP TRIGGER IF EXISTS record_delete_start ON embedding` + await sql`DROP TRIGGER IF EXISTS record_delete_end ON embedding` + await sql`DROP FUNCTION record_delete_bound()` + await sql`DROP TABLE delete_bound` + } + }, 60_000) + + it('runs maintenance only in background slices, one checkpointed operation past the budget at most', async () => { + expect(await pass()).toBe(true) + await sql`CREATE INDEX aa_retirement_hnsw_idx ON embedding_search USING hnsw (vector public.vector_l2_ops)` + await sql`CREATE INDEX bb_retirement_hnsw_idx ON embedding_search USING hnsw (vector public.vector_l2_ops)` + const oids = async () => + ( + await sql`SELECT to_regclass('aa_retirement_hnsw_idx')::oid AS aa, + to_regclass('bb_retirement_hnsw_idx')::oid AS bb` + )[0] + const checkpoints = async () => + ( + await sql`SELECT reindexed_through, vacuumed_tables FROM search_embedding_cleanup_progress` + )[0] + const journaled = async () => + ( + await sql`SELECT name FROM script_migrations WHERE name = '0029_retire_all_search_embeddings'` + ).length === 1 + const original = await oids() + const blocker = postgres(readTestDatabaseUrl(), { + max: 1, + connection: { search_path: schema }, + onnotice: () => undefined, + }) + let release!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + let signalLocked!: () => void + const locked = new Promise((resolve) => { + signalLocked = resolve + }) + let holding: Promise | undefined + try { + await runScriptMigrations(sql, [ + searchRetirementSlice({ budgetMs: 60_000, maintenance: false }), + ]) + expect(await oids()).toEqual(original) + expect(await checkpoints()).toEqual({ reindexed_through: '', vacuumed_tables: 0 }) + expect(await journaled()).toBe(false) + + /** An open writer makes the first concurrent rebuild outlast the slice's budget. */ + holding = blocker.begin(async (tx) => { + await tx`UPDATE embedding_search SET vector = '[3,2,1]' WHERE id = '00001'` + signalLocked() + await released + }) + await locked + const slice = runScriptMigrations(sql, [ + searchRetirementSlice({ budgetMs: 300, maintenance: true, throttle: UNTHROTTLED }), + ]) + await sleep(1_000) + release() + await holding + await slice + const afterSlice = await oids() + expect(afterSlice.aa).not.toBe(original.aa) + expect(afterSlice.bb).toBe(original.bb) + expect(await checkpoints()).toEqual({ + reindexed_through: 'aa_retirement_hnsw_idx', + vacuumed_tables: 0, + }) + expect(await journaled()).toBe(false) + + await runScriptMigrations(sql, [backgroundRun]) + const finished = await oids() + expect(finished.aa).toBe(afterSlice.aa) + expect(finished.bb).not.toBe(original.bb) + expect(await checkpoints()).toEqual({ + reindexed_through: 'bb_retirement_hnsw_idx', + vacuumed_tables: 6, + }) + expect(await journaled()).toBe(true) + } finally { + release() + await holding + await blocker.end() + await sql`DROP INDEX IF EXISTS aa_retirement_hnsw_idx` + await sql`DROP INDEX IF EXISTS bb_retirement_hnsw_idx` + } + }, 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..1f3fb396d0c 100644 --- a/packages/db/script-migrations/0027_retire_search_embeddings.ts +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -1,10 +1,9 @@ -import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url' -import type { ScriptMigration } from '@sim/db/script-migrations/types' +import { type ScriptMigration, ScriptMigrationDeferred } from '@sim/db/script-migrations/types' import { retryOnLockTimeout } from '@sim/db/scripts/lock-timeout-retry' import { createLogger } from '@sim/logger' import { getPostgresCancellationReason } from '@sim/utils/errors' import { sleep } from '@sim/utils/helpers' -import postgres, { type Sql, type TransactionSql } from 'postgres' +import type { Sql, TransactionSql } from 'postgres' const logger = createLogger('RetireSearchEmbeddings') /** Most IDs one page reads in primary-key order; reading is cheap next to the mutation. */ @@ -28,12 +27,52 @@ const SLOW_PAGE_MS = 30_000 * to a size that timed out. */ const FAST_PAGE_MS = SLOW_PAGE_MS / 4 +const LOCK_RETRY_BUDGET_MS = 60_000 +/** One runner at a time: the deploy slice defers while the background runner holds it, and vice versa. */ +const RETIREMENT_LOCK = 'sim:search-retirement' + /** - * 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 the cleanup yields to the database's own load. Every page is followed by the longest of three + * pauses, capped at `maxPauseMs`: a duty-cycle pause proportional to the page, a WAL pause that + * keeps the database's WAL rate inside the budget, and an exponential back-off while commits are + * slow or a replica lags. */ -const MAX_PAGE_PAUSE_MS = 5_000 -const LOCK_RETRY_BUDGET_MS = 60_000 +export interface RetirementThrottle { + /** Pause after each page as a multiple of its duration; 1 keeps the cleanup busy at most half the time. */ + dutyRatio: number + /** + * Share of `max_wal_size` the database may write per `checkpoint_timeout` while the cleanup runs, + * counting every writer, so checkpoints stay time-triggered instead of WAL-triggered. A burst of + * WAL-triggered checkpoints re-logs a full image of each page first touched after every + * checkpoint, which is what multiplies WAL during a bulk delete. + */ + walBudgetShare: number + /** A page commit slower than this backs off. It includes the synchronous-replication wait. */ + slowCommitMs: number + /** The first back-off after a slow commit or a lagging replica; it doubles while either persists. */ + backoffMs: number + maxPauseMs: number + /** Replication lag (write, flush or replay) above this pauses before the next page. */ + maxReplicaLagMs: number +} + +export const DEFAULT_RETIREMENT_THROTTLE: RetirementThrottle = { + dutyRatio: 1, + walBudgetShare: 0.1, + slowCommitMs: 1_000, + backoffMs: 15_000, + maxPauseMs: 5 * 60_000, + maxReplicaLagMs: 10_000, +} + +export interface RetirementRunOptions { + /** `performance.now()` after which no page starts; a page already running finishes. */ + deadline?: number + throttle?: Partial +} + +/** `complete` once every captured KB is retired; `deferred` when the budget ran out or another runner holds the lock. */ +export type RetirementOutcome = 'complete' | 'deferred' type Phase = 'documents' | 'embeddings' | 'done' @@ -47,6 +86,8 @@ interface PageResult { mutated: number /** A phase change the page committed. */ transition?: 'embeddings' | 'documents_rescan' | 'embeddings_rescan' + /** Time from the page's last statement to its acknowledged commit. */ + commitMs: number } /** @@ -83,120 +124,254 @@ interface Progress { /** * 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`. + * durable cursor. Other KBs are traversed without modification. The runner owns bookkeeping, as it + * owns `script_migrations`. Runs to completion; the registered successor runs it in budgeted slices. */ 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 ( - 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` - 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') - } + if ((await retireSearchEmbeddings(sql)) === 'deferred') { + throw new ScriptMigrationDeferred('Another Search retirement runner holds the lock') + } + }, +} - /** 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 }[]>` - 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 (!hasTargets) return +/** + * Advances the retirement until it completes or `deadline` passes, throttled on the database's + * commit latency, WAL rate and replication lag. Resumes the saved cursor; safe to interrupt at any + * point, since each page commits its mutation together with the cursor. + */ +export async function retireSearchEmbeddings( + sql: Sql, + options: RetirementRunOptions = {} +): Promise { + const deadline = options.deadline ?? Number.POSITIVE_INFINITY + const throttle = { ...DEFAULT_RETIREMENT_THROTTLE, ...options.throttle } + const [{ locked }] = + await sql`SELECT pg_try_advisory_lock(hashtextextended(${RETIREMENT_LOCK}, 0)) AS locked` + if (!locked) { + logger.info('Another Search retirement runner is active; deferring') + return 'deferred' + } + try { + if (!(await prepareTargets(sql))) return 'complete' + return await retireTargets(sql, deadline, throttle) + } finally { + await sql`SELECT pg_advisory_unlock(hashtextextended(${RETIREMENT_LOCK}, 0))` + } +} - 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 +/** Captures the target snapshot once and reports whether any target remains. */ +async function prepareTargets(sql: Sql): Promise { + return 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` + 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') + } + + /** 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 (;;) { - /** 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) { + 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 }[]>` + 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 + }) +} + +async function retireTargets( + sql: Sql, + deadline: number, + throttle: RetirementThrottle +): Promise { + const walBytesPerMs = await walBudgetBytesPerMs(sql, throttle.walBudgetShare) + 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 + /** Consecutive slow commits or lagging-replica checks; each doubles the back-off. */ + let strain = 0 + const backoffMs = () => Math.min(throttle.maxPauseMs, throttle.backoffMs * 2 ** (strain - 1)) + const pause = (ms: number) => sleep(Math.max(0, Math.min(ms, deadline - performance.now()))) + + for (;;) { + if (performance.now() >= deadline) { + logger.info('Search retirement budget spent; deferring', { batches, mutated, rowLimit }) + return 'deferred' + } + const health = await readHealth(sql) + if (health.replicaLagMs !== null && health.replicaLagMs > throttle.maxReplicaLagMs) { + strain++ + const pauseMs = backoffMs() + logger.warn('Search retirement batch', { + throttle: 'replica_lag', + replicaLagMs: Math.round(health.replicaLagMs), + pauseMs, + rowLimit, + }) + await pause(pauseMs) + continue + } + + /** 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 + const pauseMs = Math.min( + throttle.maxPauseMs, + (performance.now() - pageStartedAt) * throttle.dutyRatio + ) + logger.warn('Search retirement page timed out; retrying with fewer rows', { + rowLimit, + pauseMs: Math.round(pauseMs), + }) + await pause(pauseMs) + continue + } + if (page.done) break + const pageMs = performance.now() - pageStartedAt + const walBytes = await walSince(sql, health.lsn) + batches++ + mutated += page.mutated + + let state: 'steady' | 'slow_commit' | 'slow_page' | 'phase_change' = 'steady' + if (page.transition) { + /** A phase change may include a full recheck, which says nothing about page cost. */ + state = 'phase_change' + logger.info('Search retirement phase changed', { + transition: page.transition, + batches, + mutated, + }) + } else if (page.commitMs > throttle.slowCommitMs) { + state = 'slow_commit' + strain++ + rowLimit = halve(rowLimit) + } else { + strain = 0 + if (pageMs > SLOW_PAGE_MS) { + state = 'slow_page' rowLimit = halve(rowLimit) - logger.warn('Search retirement page was slow; halving the row limit', { - pageMs: Math.round(pageMs), - rowLimit, - }) } else if (pageMs < FAST_PAGE_MS) { rowLimit = Math.min(ceiling, rowLimit * 2) } - if (batches % 10 === 0) { - logger.info('Search embedding retirement progress', { - batches, - phase: page.phase, - afterId: page.afterId, - mutated, - rowLimit, - elapsedMs: Date.now() - startedAt, - }) - } - await sleep(Math.min(pageMs, MAX_PAGE_PAUSE_MS)) } - logger.info('Selected Search knowledge bases retired', { - batches, - mutated, - elapsedMs: Date.now() - startedAt, + + const dutyPauseMs = pageMs * throttle.dutyRatio + /** Long enough that the WAL written during the page, spread over page and pause, fits the budget. */ + const walPauseMs = walBytes / walBytesPerMs - pageMs + const strainPauseMs = state === 'slow_commit' ? backoffMs() : 0 + const pauseMs = Math.max( + 0, + Math.min(throttle.maxPauseMs, Math.max(dutyPauseMs, walPauseMs, strainPauseMs)) + ) + const limitedBy = + strainPauseMs >= Math.max(dutyPauseMs, walPauseMs) && strainPauseMs > 0 + ? 'backoff' + : walPauseMs > dutyPauseMs + ? 'wal' + : 'duty' + logger.info('Search retirement batch', { + phase: page.phase, + rows: page.mutated, + pageMs: Math.round(pageMs), + commitMs: Math.round(page.commitMs), + walBytes, + pauseMs: Math.round(pauseMs), + throttle: state === 'steady' ? limitedBy : state, + rowLimit, }) - }, + if (batches % 10 === 0) { + logger.info('Search embedding retirement progress', { + batches, + phase: page.phase, + afterId: page.afterId, + mutated, + rowLimit, + elapsedMs: Date.now() - startedAt, + }) + } + await pause(pauseMs) + } + logger.info('Selected Search knowledge bases retired', { + batches, + mutated, + elapsedMs: Date.now() - startedAt, + }) + return 'complete' +} + +/** + * The WAL rate the budget allows, in bytes per millisecond: a share of the WAL that may accumulate + * within one checkpoint interval before `max_wal_size` forces an early checkpoint. + */ +async function walBudgetBytesPerMs(sql: Sql, share: number): Promise { + const [settings] = await sql<{ max_wal_bytes: number; checkpoint_ms: number }[]>` + SELECT pg_size_bytes(current_setting('max_wal_size'))::float8 AS max_wal_bytes, + (extract(epoch FROM current_setting('checkpoint_timeout')::interval) * 1000)::float8 AS checkpoint_ms` + return (share * settings.max_wal_bytes) / settings.checkpoint_ms +} + +interface Health { + lsn: string + /** The worst replica's lag, or null without replicas or without `pg_read_all_stats`. */ + replicaLagMs: number | null +} + +async function readHealth(sql: Sql): Promise { + const [health] = await sql<{ lsn: string; replica_lag_ms: number | null }[]>` + SELECT pg_current_wal_lsn()::text AS lsn, + (SELECT extract(epoch FROM max(greatest(write_lag, flush_lag, replay_lag))) * 1000 + FROM pg_stat_replication)::float8 AS replica_lag_ms` + return { lsn: health.lsn, replicaLagMs: health.replica_lag_ms } +} + +/** WAL the whole database wrote since `lsn`, every writer included. */ +async function walSince(sql: Sql, lsn: string): Promise { + const [row] = await sql<{ bytes: number }[]>` + SELECT pg_wal_lsn_diff(pg_current_wal_lsn(), ${lsn}::pg_lsn)::float8 AS bytes` + return row.bytes } /** @@ -207,32 +382,57 @@ export const retireSearchEmbeddingsMigration: ScriptMigration = { */ async function retirePage(sql: Sql, rowLimit: number): Promise { return retryOnLockTimeout( - () => - sql.begin(async (tx) => { - const scanLimit = Math.min(SCAN_PAGE_SIZE, rowLimit * SCAN_ROWS_PER_MUTATION) - await tx`SET LOCAL statement_timeout = '120s'` - await tx`SET LOCAL lock_timeout = '1s'` - const [progress] = await tx` + async () => { + let statementsDoneAt = 0 + const page = await sql.begin(async (tx) => { + const result = await retirePageStatements(tx, rowLimit) + statementsDoneAt = performance.now() + return result + }) + return { ...page, commitMs: performance.now() - statementsDoneAt } + }, + { + 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), + }), + } + ) +} + +async function retirePageStatements( + tx: TransactionSql, + rowLimit: number +): Promise> { + const scanLimit = Math.min(SCAN_PAGE_SIZE, rowLimit * SCAN_ROWS_PER_MUTATION) + await tx`SET LOCAL statement_timeout = '120s'` + await tx`SET LOCAL lock_timeout = '1s'` + const [progress] = await tx` SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress WHERE id = 1 FOR UPDATE` - const result = (page: Partial = {}): PageResult => ({ - done: false, - phase: progress.phase, - afterId: progress.after_id, - mutated: 0, - ...page, - }) - if (progress.phase === 'done') { - /** A retry after maintenance failed rechecks every captured KB, as completion did. */ - await tx`SET LOCAL statement_timeout = '30min'` - await validateTargetMarkers(tx) - return result({ done: true }) - } + const result = ( + page: Partial> = {} + ): Omit => ({ + done: false, + phase: progress.phase, + afterId: progress.after_id, + mutated: 0, + ...page, + }) + if (progress.phase === 'done') { + /** A retry after maintenance failed rechecks every captured KB, as completion did. */ + await tx`SET LOCAL statement_timeout = '30min'` + await validateTargetMarkers(tx) + return result({ done: true }) + } - if (progress.phase === 'documents') { - /** Already-retired documents are skipped so they never spend the row limit. */ - const [page] = await pageMutation(tx< - { after_id: string; mutated: number; invalid_target: boolean }[] - >` + if (progress.phase === 'documents') { + /** Already-retired documents are skipped so they never spend the row limit. */ + const [page] = await pageMutation(tx< + { after_id: string; mutated: number; invalid_target: boolean }[] + >` WITH source_page AS MATERIALIZED ( SELECT id, knowledge_base_id, (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL @@ -266,25 +466,24 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { (SELECT count(*) FROM retired)::int AS mutated, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target`) - if (!page) { - await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` - return result({ afterId: '', transition: 'embeddings' }) - } - if (page.invalid_target) - throw new Error('Cleanup target is no longer a Search knowledge base') - await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return result({ afterId: page.after_id, mutated: page.mutated }) - } + if (!page) { + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` + return result({ afterId: '', transition: 'embeddings' }) + } + if (page.invalid_target) throw new Error('Cleanup target is no longer a Search knowledge base') + await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` + return result({ afterId: page.after_id, mutated: page.mutated }) + } - /** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */ - const [page] = await pageMutation(tx< - { - after_id: string - mutated: number - unretired: boolean - invalid_target: boolean - }[] - >` + /** Keep page IDs in PostgreSQL; foreign keys cascade projection and provenance deletes. */ + const [page] = await pageMutation(tx< + { + after_id: string + mutated: number + unretired: boolean + invalid_target: boolean + }[] + >` WITH source_page AS MATERIALIZED ( SELECT id, knowledge_base_id, document_id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${scanLimit} ), target_page AS MATERIALIZED ( @@ -315,48 +514,36 @@ async function retirePage(sql: Sql, rowLimit: number): Promise { EXISTS (SELECT 1 FROM unretired) AS unretired, EXISTS (SELECT 1 FROM invalid_target) AS invalid_target FROM source_page p CROSS JOIN page_end e GROUP BY e.limited, e.last_target`) - if (page?.invalid_target) - throw new Error('Cleanup target is no longer a Search knowledge base') - if (page?.unretired) - throw new Error('Search content changed after retirement; stop writers before resuming') - if (!page) { - /** - * A late insert may sort behind either UUID cursor; completion must recheck the target. - * These rechecks walk every captured KB once, which no single page does. - */ - await tx`SET LOCAL statement_timeout = '30min'` - logger.info('Rechecking captured Search knowledge bases before completion') - const [unretired] = await tx`SELECT d.id FROM document d + 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) { + /** + * A late insert may sort behind either UUID cursor; completion must recheck the target. + * These rechecks walk every captured KB once, which no single page does. + */ + await tx`SET LOCAL statement_timeout = '30min'` + logger.info('Rechecking captured Search knowledge bases before completion') + const [unretired] = await tx`SELECT d.id FROM document d JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = d.knowledge_base_id WHERE (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL) LIMIT 1` - if (unretired) { - await tx`UPDATE search_embedding_cleanup_progress SET phase = 'documents', after_id = '' WHERE id = 1` - return result({ afterId: '', transition: 'documents_rescan' }) - } - const [remaining] = await tx`SELECT e.id FROM embedding e + if (unretired) { + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'documents', after_id = '' WHERE id = 1` + return result({ afterId: '', transition: 'documents_rescan' }) + } + const [remaining] = await tx`SELECT e.id FROM embedding e JOIN search_embedding_cleanup_targets t ON t.knowledge_base_id = e.knowledge_base_id LIMIT 1` - if (remaining) { - await tx`UPDATE search_embedding_cleanup_progress SET after_id = '' WHERE id = 1` - return result({ afterId: '', transition: 'embeddings_rescan' }) - } - await validateTargetMarkers(tx) - await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` - return result({ done: true }) - } - await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` - return result({ afterId: page.after_id, mutated: page.mutated }) - }), - { - 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), - }), + if (remaining) { + await tx`UPDATE search_embedding_cleanup_progress SET after_id = '' WHERE id = 1` + return result({ afterId: '', transition: 'embeddings_rescan' }) } - ) + await validateTargetMarkers(tx) + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` + return result({ done: true }) + } + await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${page.after_id} WHERE id = 1` + return result({ afterId: page.after_id, mutated: page.mutated }) } /** Validate even empty or fully scanned KBs, holding marker locks until completion commits. */ @@ -379,19 +566,3 @@ async function validateTargetMarkers(tx: TransactionSql): Promise { afterId = page.after_id } } - -/** 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') - const { retireAllSearchEmbeddingsMigration } = await import( - '@sim/db/script-migrations/0029_retire_all_search_embeddings' - ) - await runScriptMigrations(sql, [retireAllSearchEmbeddingsMigration]) - } 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 index 747d08612c0..3337152bc60 100644 --- a/packages/db/script-migrations/0028_maintain_search_retirement.ts +++ b/packages/db/script-migrations/0028_maintain_search_retirement.ts @@ -16,36 +16,75 @@ const VACUUM_TABLES = [ /** * Runs after retirement, including on databases that already journaled 0027. Concurrent rebuilds - * and vacuum must run outside transactions; checkpoints follow each successful operation. + * and vacuum must run outside transactions; checkpoints follow each successful operation. Runs to + * completion; the registered successor runs it within a budget instead. */ 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') + if ((await maintainSearchRetirement(sql)) === 'deferred') { + throw new Error('Search retirement maintenance is already running') + } + }, +} - 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 +export interface MaintenanceRunOptions { + /** + * `performance.now()` after which no rebuild or vacuum starts. One that already started runs to + * completion: an interrupted concurrent rebuild restarts from scratch, so cancelling it at a + * deadline could keep a large index from ever finishing. + */ + deadline?: number +} + +/** + * Rebuilds the HNSW indexes and vacuums the retired tables, one checkpointed operation at a time. + * Returns `deferred` when the deadline passes with work left or another worker holds the lock. + */ +export async function maintainSearchRetirement( + sql: Sql, + options: MaintenanceRunOptions = {} +): Promise<'complete' | 'deferred'> { + const deadline = options.deadline ?? Number.POSITIVE_INFINITY + const outOfBudget = (operation: Record) => { + if (performance.now() < deadline) return false + logger.info('Search retirement maintenance deferred to a later slice', operation) + return true + } + const [table] = await sql`SELECT to_regclass('search_embedding_cleanup_progress') AS relation` + if (!table.relation) return 'complete' + const [retirement] = await sql`SELECT phase FROM search_embedding_cleanup_progress WHERE id = 1` + if (!retirement) return 'complete' + 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) { + logger.info('Search retirement maintenance is already running; deferring') + return 'deferred' + } + const [settings] = await sql`SELECT current_setting('statement_timeout') AS statement_timeout, + current_setting('lock_timeout') AS lock_timeout, + current_setting('vacuum_cost_delay') AS vacuum_cost_delay` + 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 }[]>` + }) + /** Concurrent maintenance waits for old snapshots without blocking ordinary table writes. */ + await sql`SET statement_timeout = 0` + await sql`SET lock_timeout = 0` + /** + * A manual VACUUM is unthrottled by default; this gives it autovacuum's default cost delay so + * its reads, dirtied pages and WAL are spread out instead of competing with live traffic. + */ + await sql`SET vacuum_cost_delay = '2ms'` + 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 @@ -53,48 +92,51 @@ export const maintainSearchRetirementMigration: ScriptMigration = { 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, - }) - } + if (!index) break + if (outOfBudget({ index: index.name })) return 'deferred' + 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 }[]>` + 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 }[]>` + for (let step = progress.vacuumed_tables; step < VACUUM_TABLES.length; step++) { + const tableName = VACUUM_TABLES[step] + if (outOfBudget({ table: tableName })) return 'deferred' + 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, - }) - } + 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), + set_config('vacuum_cost_delay', ${settings.vacuum_cost_delay}, false)` } 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))` - } + await sql`SELECT pg_advisory_unlock(hashtextextended(${MAINTENANCE_LOCK}, 0))` } - }, + } + return 'complete' } /** PostgreSQL leaves invalid _ccnew/_ccold siblings if a concurrent rebuild is interrupted. */ 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..6288f0d801f 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,90 @@ -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' +import { parseArgs } from 'node:util' +import { + type RetirementThrottle, + retireSearchEmbeddings, + retireSearchEmbeddingsMigration, +} from '@sim/db/script-migrations/0027_retire_search_embeddings' +import { + maintainSearchRetirement, + maintainSearchRetirementMigration, +} from '@sim/db/script-migrations/0028_maintain_search_retirement' +import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url' +import { type ScriptMigration, ScriptMigrationDeferred } from '@sim/db/script-migrations/types' +import postgres from 'postgres' + +/** + * How long a deploy advances the retirement before deferring to the background runner. The deploy + * still makes progress on databases without that runner, and stays off the release's critical path. + */ +export const DEPLOY_RETIREMENT_BUDGET_MS = 60_000 + +export interface SearchRetirementSliceOptions { + /** Time after which no retirement page or maintenance operation starts. */ + budgetMs: number + /** Whether this slice may start index rebuilds and vacuums; the deploy never does. */ + maintenance: boolean + throttle?: Partial +} + +/** + * One budgeted slice of the Search retirement under the `0029` journal name. A slice that finishes + * retirement and maintenance is journaled with its superseded names; any other slice throws + * `ScriptMigrationDeferred`, so the runner leaves it unrecorded and the next slice resumes the saved + * cursor and maintenance checkpoints. + */ +export function searchRetirementSlice(options: SearchRetirementSliceOptions): ScriptMigration { + return { + name: '0029_retire_all_search_embeddings', + supersedes: [retireSearchEmbeddingsMigration.name, maintainSearchRetirementMigration.name], + async up(sql) { + const deadline = performance.now() + options.budgetMs + const retirement = await retireSearchEmbeddings(sql, { + deadline, + throttle: options.throttle, + }) + if (retirement === 'deferred') { + throw new ScriptMigrationDeferred('Search retirement continues in the background runner') + } + const maintenance = await maintainSearchRetirement(sql, { + deadline: options.maintenance ? deadline : Number.NEGATIVE_INFINITY, + }) + if (maintenance === 'deferred') { + throw new ScriptMigrationDeferred('Search retirement maintenance waits for its window') + } + }, + } +} /** 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) - }, +export const retireAllSearchEmbeddingsMigration = searchRetirementSlice({ + budgetMs: DEPLOY_RETIREMENT_BUDGET_MS, + maintenance: false, +}) + +/** + * The background runner: `--budget-minutes` bounds the slice (default: run to completion) and + * `--maintenance` lets it start index rebuilds and vacuums. Journals only a completed cleanup. + */ +if (import.meta.main) { + const { values } = parseArgs({ + options: { + 'budget-minutes': { type: 'string' }, + maintenance: { type: 'boolean', default: false }, + }, + }) + const budgetMinutes = values['budget-minutes'] + const budgetMs = + budgetMinutes === undefined ? Number.POSITIVE_INFINITY : Number(budgetMinutes) * 60_000 + if (!(budgetMs > 0)) throw new Error('--budget-minutes must be a positive number') + 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, [ + searchRetirementSlice({ budgetMs, maintenance: values.maintenance }), + ]) + } finally { + await sql.end() + } } diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md index d9c531b50f3..5c4e7db799e 100644 --- a/packages/db/script-migrations/search-embedding-retirement.md +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -14,61 +14,115 @@ this cleanup ships: deployment migrations run before the new app switches over. (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. +also honor the indexed-search gate. The app does not depend on the cleanup finishing: once documents +are retired, what remains is storage reclamation. + +The cleanup is a background data migration, in the pattern of GitLab's +[batched background migrations](https://docs.gitlab.com/development/database/batched_background_migrations/): +the deploy starts it, a background runner advances it between deploys, throttled on database +health, and the journal records it only once it is finished. + +- **Deploy slice.** The ordinary migration runner applies `0029_retire_all_search_embeddings` with a + one-minute budget. It captures the target snapshot on first run, advances the retirement, and + throws `ScriptMigrationDeferred`, so the runner leaves it unrecorded and the release proceeds. It + never starts index maintenance. If the background runner holds the retirement lock, the deploy + slice defers at once without running a page. +- **Background runner.** `.github/workflows/search-retirement.yml` runs a slice every hour against + staging and production with the same secrets, and therefore the same migration role, as + `migrations.yml`. Each hourly slice may start pages for 50 minutes. A daily slice in an off-peak + window (03:07 UTC) may also start index rebuilds and vacuums for two hours. The first slice that + finishes retirement and maintenance records `0029` and its superseded names in + `script_migrations`; later slices find nothing pending. The migration role is used because it + owns these tables, which the rebuilds and vacuums require, and can read replication statistics. + GitHub runs scheduled workflows from the default branch, so the schedule starts once this ships + to `main`. Remove the workflow once `0029` is journaled in every environment. +- **Manual slice.** `bun run packages/db/script-migrations/0029_retire_all_search_embeddings.ts + [--budget-minutes N] [--maintenance]`, with the writer supplied through + `MIGRATION_DATABASE_URL`/`DATABASE_URL`, runs one slice the same way. Without `--budget-minutes` it + runs until the retirement finishes; without `--maintenance` it never starts a rebuild or vacuum. + Self-hosted deployments without the workflow use this command to finish the cleanup. + +A session advisory lock lets only one runner retire pages at a time, and another serializes +maintenance. No page starts after a slice's budget; the page in flight when the budget runs out +finishes or rolls back with its cursor, so a slice can overrun its budget by at most one page. +Cancelling a slice at any point loses at most its in-flight page. + +### Pacing Each page mutates at most a row limit of target rows and reads at most four IDs per row of that -limit, never more than 25,000 IDs. Pages execute -sequentially, and each is followed by a pause as long as the page took, up to five seconds, to -leave the primary headroom. Retiring a document is a non-HOT update that writes every index on -`document`, and deleting a chunk cascades into its projections, so a page's cost follows the target -rows it mutates, not the IDs it reads. A page that reaches the row limit advances the cursor only to -its last mutated row; the rest of its scan is read again by the next page. Tying the scan window to -the limit keeps that re-reading proportional to the work, even after the limit shrinks. Documents -that are already retired never count against the limit. - -The row limit starts at 2,000 rows. A page is timed from the start of its transaction through its -commit, including the synchronous-replication wait and any lock-timeout retries. A page slower than -30 seconds halves the limit. A fast page, one under 7.5 seconds, doubles it up to 8,000, which -also widens the scan window, so sparse stretches are not crawled in small windows. The limit never drops below 25 rows. Phase changes do not adjust it. +limit, never more than 25,000 IDs. Retiring a document is a non-HOT update that writes every index +on `document`, and deleting a chunk cascades into its projections, so a page's cost follows the +target rows it mutates, not the IDs it reads. A page that reaches the row limit advances the cursor +only to its last mutated row; the rest of its scan is read again by the next page. Documents that are +already retired never count against the limit. + +The row limit starts at 2,000 rows, never drops below 25 and never exceeds 8,000. A page is timed +from the start of its transaction through its commit, including any lock-timeout retries. A page +slower than 30 seconds halves the limit. A page faster than 7.5 seconds doubles it, which also +widens the scan window. Phase changes do not adjust it. + +After every page the cleanup pauses for the longest of: + +- **Duty cycle**: as long as the page took, so the cleanup is busy at most half the time. +- **WAL budget**: long enough that the WAL the whole database wrote during the page, spread over + the page and the pause, stays within 10% of `max_wal_size` per `checkpoint_timeout`. The rate is + read from `pg_current_wal_lsn()` before and after each page and counts every writer, so the cleanup + yields when the application itself writes heavily. Holding WAL well below `max_wal_size` keeps + checkpoints time-triggered. When WAL outruns `max_wal_size`, checkpoints start back to back. The + first change to each page after a checkpoint logs the whole page, so WAL grows further, and + synchronous commits slow down for every writer. +- **Commit back-off**: a page whose commit took longer than one second halves the row limit and + waits 15 seconds, doubling while commits stay slow. The commit is timed separately from the page's + statements, so it isolates the synchronous-replication wait and checkpoint fsync pressure that + every other writer is also paying. + +Before each page the cleanup also reads `pg_stat_replication`; a replica whose write, flush or +replay lag exceeds 10 seconds pauses the cleanup with the same doubling back-off. Without replicas, +or without `pg_read_all_stats`, the lag reads as unknown and only the other signals apply. No pause +exceeds five minutes. Each pause is cut short at the slice's budget. Materialized SQL pages keep the IDs inside PostgreSQL; the migration process receives only a cursor and a validation result. Each page uses a two-minute statement timeout and a one-second lock timeout. If a page's mutating statement exceeds the statement timeout, the page rolls back with its -cursor and is retried with half the row limit after the usual pause. From then on, fast pages grow -the limit only up to that halved size, so a size that timed out is never tried again. A page that -still times out at 25 rows fails the migration. Any other statement timeout fails the migration at once, because a smaller page cannot -speed it up. The completion rechecks, which walk every captured KB once, run with a 30-minute -timeout. Brief lock timeouts retry the rolled-back page with bounded backoff for up to one minute. -Other errors, or exhausted lock retries, fail the migration without a completion receipt. - -Every ten pages the migration logs the phase, cursor, rows mutated so far and current row limit. It -also logs each halving after a slow page, each phase change and the start of the completion -recheck. The deployment job retains its five-hour overall timeout; it is not a runtime estimate. A -large cleanup can need more than one job run, and each run resumes from the saved cursor. - -After interruption or failure, rerun the migration job, or run -`bun run packages/db/script-migrations/0027_retire_search_embeddings.ts` with the writer supplied -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 +cursor and is retried with half the row limit after a duty-cycle pause. From then on, fast pages grow +the limit only up to that halved size. A page that still times out at 25 rows fails the slice. Any +other statement timeout fails the slice at once, because a smaller page cannot speed it up. The +completion rechecks, which walk every captured KB once, run with a 30-minute timeout. Brief lock +timeouts retry the rolled-back page with bounded backoff for up to one minute. Other errors, or +exhausted lock retries, fail the slice without a completion receipt; the next slice resumes. + +### Observing + +Every page logs one `Search retirement batch` line with its phase, rows mutated, page and commit +duration, WAL bytes, the pause it chose and the signal that set it (`duty`, `wal`, `backoff`, +`slow_commit`, `slow_page` or `phase_change`), and the current row limit. A replica-lag pause logs a +`replica_lag` batch line without rows. Every ten pages the cleanup also logs its cursor and running +totals, and it logs each phase change, the completion recheck and each deferral. In production, +also watch primary commit latency, checkpoint frequency (`checkpoint starting: wal` in the server +log means WAL is forcing them), replica lag and free disk. + +### Pausing and resuming + +Pause the background runner with `gh workflow disable search-retirement.yml` and cancel any +in-progress run; the in-flight page rolls back with its cursor. Deploys keep advancing the cleanup +for one minute each while the workflow is disabled. Resume with `gh workflow enable` or a manual +`workflow_dispatch`. Nothing is reset on resume: every slice continues from the saved scope, phase, +cursor and maintenance checkpoints. + +### Legacy checkpoints + +`0029` supersedes the single-KB retirement and maintenance entries (`0027`/`0028`), including +databases that already recorded either receipt. 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 slices resume the saved scope, phase, cursor +and maintenance checkpoints. + +Every slice requires a direct or session-pooled PostgreSQL connection, as the deployment +migration runner already does for its session advisory locks 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. @@ -110,12 +164,19 @@ 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`, -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. +Once retirement completes, a slice allowed to run maintenance 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. Deploy slices and hourly slices +never start these; only the off-peak daily slice or a manual `--maintenance` slice does. Operations +run sequentially outside transactions, and none starts after the slice's budget. One that already +started runs to completion, because an interrupted concurrent rebuild starts over and a large index +cancelled at every deadline would never finish. The job timeout leaves room for that. + +Rebuilds keep ordinary reads and writes available, but each writes its whole new index to WAL and +needs temporary index space; they wait for older transactions. A manual `VACUUM` is not cost-limited +by default, so maintenance runs it with autovacuum's default `vacuum_cost_delay` of 2 ms to spread +its I/O and WAL. 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