Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
@@ -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'
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
})
246 changes: 147 additions & 99 deletions packages/db/script-migrations/0027_retire_search_embeddings.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand All @@ -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. */
Expand All @@ -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'

Expand Down Expand Up @@ -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<void> {
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<Progress[]>`
const [existing] = await tx<Progress[]>`
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,
})
}

/**
Expand Down Expand Up @@ -380,17 +409,36 @@ async function validateTargetMarkers(tx: TransactionSql): Promise<void> {
}
}

/** 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()
}
Expand Down
Loading
Loading