1+ import { parseArgs } from 'node:util'
12import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url'
23import type { ScriptMigration } from '@sim/db/script-migrations/types'
34import { retryOnLockTimeout } from '@sim/db/scripts/lock-timeout-retry'
@@ -18,7 +19,8 @@ const SCAN_ROWS_PER_MUTATION = 4
1819/**
1920 * Rows one page may update or delete. Every retired document is a non-HOT update touching each of
2021 * its indexes, and every deleted chunk cascades into its projections, so the write cost of a page,
21- * not its scan, is what can outrun the statement timeout.
22+ * not its scan, is what can outrun the statement timeout. `maxRows` lowers the starting limit and
23+ * the ceiling it may grow to.
2224 */
2325const ROW_LIMIT = { initial : 2_000 , min : 25 , max : 8_000 } as const
2426/** A page slower than this halves the row limit. */
@@ -28,12 +30,23 @@ const SLOW_PAGE_MS = 30_000
2830 * to a size that timed out.
2931 */
3032const FAST_PAGE_MS = SLOW_PAGE_MS / 4
33+ /** The longest pause after one page, however slow the page was. */
34+ const MAX_PAGE_PAUSE_MS = 60_000
35+ const LOCK_RETRY_BUDGET_MS = 60_000
36+
3137/**
32- * Each page, committed or timed out, is followed by a pause as long as the page, up to this, to
33- * leave the primary headroom.
38+ * How hard one run pushes the primary. Each page, committed or timed out, is followed by a pause of
39+ * `pauseRatio` times the page's duration (up to a minute), so the run is busy at most
40+ * `1 / (1 + pauseRatio)` of the time. A page is timed through its commit, so a slow synchronous
41+ * replica or a checkpoint stall lengthens the pause by the same factor.
3442 */
35- const MAX_PAGE_PAUSE_MS = 5_000
36- const LOCK_RETRY_BUDGET_MS = 60_000
43+ export interface RetirementPacing {
44+ pauseRatio : number
45+ /** The most rows one page may update or delete, from 25 to 8,000. */
46+ maxRows : number
47+ }
48+
49+ export const DEFAULT_RETIREMENT_PACING : RetirementPacing = { pauseRatio : 2 , maxRows : 2_000 }
3750
3851type Phase = 'documents' | 'embeddings' | 'done'
3952
@@ -87,116 +100,132 @@ interface Progress {
87100 */
88101export const retireSearchEmbeddingsMigration : ScriptMigration = {
89102 name : '0027_retire_search_embeddings' ,
90- async up ( sql ) {
91- const hasTargets = await sql . begin ( 'isolation level repeatable read' , async ( tx ) => {
92- await tx `SET LOCAL statement_timeout = '120s'`
93- await tx `SET LOCAL lock_timeout = '1s'`
94- await tx `CREATE TABLE IF NOT EXISTS search_embedding_cleanup_progress (
103+ up : ( sql ) => retireSearchEmbeddings ( sql ) ,
104+ }
105+
106+ /** Retires every captured target, resuming the saved cursor, paced by `pacing`. */
107+ export async function retireSearchEmbeddings (
108+ sql : Sql ,
109+ pacing : RetirementPacing = DEFAULT_RETIREMENT_PACING
110+ ) : Promise < void > {
111+ if (
112+ ! ( pacing . pauseRatio >= 0 ) ||
113+ ! Number . isInteger ( pacing . maxRows ) ||
114+ pacing . maxRows < ROW_LIMIT . min ||
115+ pacing . maxRows > ROW_LIMIT . max
116+ ) {
117+ throw new Error (
118+ `Search retirement pacing needs a pause ratio of at least 0 and ${ ROW_LIMIT . min } -${ ROW_LIMIT . max } max rows`
119+ )
120+ }
121+ const pause = ( pageMs : number ) => sleep ( Math . min ( pageMs * pacing . pauseRatio , MAX_PAGE_PAUSE_MS ) )
122+ const hasTargets = await sql . begin ( 'isolation level repeatable read' , async ( tx ) => {
123+ await tx `SET LOCAL statement_timeout = '120s'`
124+ await tx `SET LOCAL lock_timeout = '1s'`
125+ await tx `CREATE TABLE IF NOT EXISTS search_embedding_cleanup_progress (
95126 id integer PRIMARY KEY CHECK (id = 1), knowledge_base_id text NOT NULL,
96127 phase text NOT NULL CHECK (phase IN ('documents', 'embeddings', 'done')),
97128 after_id text NOT NULL
98129 )`
99- const [ existing ] = await tx < Progress [ ] > `
130+ const [ existing ] = await tx < Progress [ ] > `
100131 SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress WHERE id = 1 FOR UPDATE`
101- const [ snapshot ] =
102- await tx `SELECT to_regclass('search_embedding_cleanup_targets') AS relation`
103- if ( snapshot . relation ) return Boolean ( existing )
104- if ( existing && existing . phase !== 'done' ) {
105- const [ target ] = await tx `SELECT id FROM knowledge_base
132+ const [ snapshot ] = await tx `SELECT to_regclass('search_embedding_cleanup_targets') AS relation`
133+ if ( snapshot . relation ) return Boolean ( existing )
134+ if ( existing && existing . phase !== 'done' ) {
135+ const [ target ] = await tx `SELECT id FROM knowledge_base
106136 WHERE id = ${ existing . knowledge_base_id } AND is_search_index FOR SHARE`
107- if ( ! target ) throw new Error ( 'Cleanup target is no longer a Search knowledge base' )
108- }
137+ if ( ! target ) throw new Error ( 'Cleanup target is no longer a Search knowledge base' )
138+ }
109139
110- /** Creating the snapshot and resetting a legacy cursor commit atomically, once. */
111- await tx `CREATE TABLE search_embedding_cleanup_targets (knowledge_base_id text PRIMARY KEY)`
112- let afterId = ''
113- for ( ; ; ) {
114- const [ page ] = await tx < { after_id : string | null } [ ] > `
140+ /** Creating the snapshot and resetting a legacy cursor commit atomically, once. */
141+ await tx `CREATE TABLE search_embedding_cleanup_targets (knowledge_base_id text PRIMARY KEY)`
142+ let afterId = ''
143+ for ( ; ; ) {
144+ const [ page ] = await tx < { after_id : string | null } [ ] > `
115145 WITH targets AS (
116146 INSERT INTO search_embedding_cleanup_targets (knowledge_base_id)
117147 SELECT id FROM knowledge_base WHERE is_search_index AND id > ${ afterId }
118148 ORDER BY id LIMIT ${ SCAN_PAGE_SIZE }
119149 RETURNING knowledge_base_id
120150 ) SELECT max(knowledge_base_id) AS after_id FROM targets`
121- if ( page . after_id === null ) break
122- afterId = page . after_id
123- }
124- const [ first ] = await tx < { knowledge_base_id : string } [ ] > `
151+ if ( page . after_id === null ) break
152+ afterId = page . after_id
153+ }
154+ const [ first ] = await tx < { knowledge_base_id : string } [ ] > `
125155 SELECT knowledge_base_id FROM search_embedding_cleanup_targets ORDER BY knowledge_base_id LIMIT 1`
126- if ( ! first ) return false
127- await tx `ANALYZE search_embedding_cleanup_targets`
128- await tx `ALTER TABLE search_embedding_cleanup_progress
156+ if ( ! first ) return false
157+ await tx `ANALYZE search_embedding_cleanup_targets`
158+ await tx `ALTER TABLE search_embedding_cleanup_progress
129159 ADD COLUMN IF NOT EXISTS reindexed_through text NOT NULL DEFAULT '',
130160 ADD COLUMN IF NOT EXISTS vacuumed_tables integer NOT NULL DEFAULT 0`
131- await tx `INSERT INTO search_embedding_cleanup_progress (id, knowledge_base_id, phase, after_id)
161+ await tx `INSERT INTO search_embedding_cleanup_progress (id, knowledge_base_id, phase, after_id)
132162 VALUES (1, ${ first . knowledge_base_id } , 'documents', '')
133163 ON CONFLICT (id) DO UPDATE SET knowledge_base_id = EXCLUDED.knowledge_base_id,
134164 phase = 'documents', after_id = '',
135165 reindexed_through = '', vacuumed_tables = 0`
136- return true
137- } )
138- if ( ! hasTargets ) return
166+ return true
167+ } )
168+ if ( ! hasTargets ) return
139169
140- const startedAt = Date . now ( )
141- let batches = 0
142- let mutated = 0
143- let rowLimit : number = ROW_LIMIT . initial
144- /** The largest limit the run may still try: half of the smallest limit that timed out. */
145- let ceiling : number = ROW_LIMIT . max
146- for ( ; ; ) {
147- /** Timed around the whole call, so the synchronous-replication wait at commit counts. */
148- const pageStartedAt = performance . now ( )
149- let page : PageResult
150- try {
151- page = await retirePage ( sql , rowLimit )
152- } catch ( error ) {
153- if ( ! ( error instanceof PageMutationTimeout ) ) throw error
154- /** The timed-out page rolled back with its cursor, so it is retried with fewer rows. */
155- if ( rowLimit <= ROW_LIMIT . min ) throw error . timeout
156- ceiling = halve ( rowLimit )
157- rowLimit = ceiling
158- logger . warn ( 'Search retirement page timed out; retrying with fewer rows' , { rowLimit } )
159- await sleep ( Math . min ( performance . now ( ) - pageStartedAt , MAX_PAGE_PAUSE_MS ) )
160- continue
161- }
162- if ( page . done ) break
163- const pageMs = performance . now ( ) - pageStartedAt
164- batches ++
165- mutated += page . mutated
166- if ( page . transition ) {
167- /** A phase change may include a full recheck, which says nothing about page cost. */
168- logger . info ( 'Search retirement phase changed' , {
169- transition : page . transition ,
170- batches,
171- mutated,
172- } )
173- } else if ( pageMs > SLOW_PAGE_MS ) {
174- rowLimit = halve ( rowLimit )
175- logger . warn ( 'Search retirement page was slow; halving the row limit' , {
176- pageMs : Math . round ( pageMs ) ,
177- rowLimit,
178- } )
179- } else if ( pageMs < FAST_PAGE_MS ) {
180- rowLimit = Math . min ( ceiling , rowLimit * 2 )
181- }
182- if ( batches % 10 === 0 ) {
183- logger . info ( 'Search embedding retirement progress' , {
184- batches,
185- phase : page . phase ,
186- afterId : page . afterId ,
187- mutated,
188- rowLimit,
189- elapsedMs : Date . now ( ) - startedAt ,
190- } )
191- }
192- await sleep ( Math . min ( pageMs , MAX_PAGE_PAUSE_MS ) )
170+ const startedAt = Date . now ( )
171+ let batches = 0
172+ let mutated = 0
173+ let rowLimit = Math . min ( ROW_LIMIT . initial , pacing . maxRows )
174+ /** The largest limit the run may still try: half of the smallest limit that timed out. */
175+ let ceiling = pacing . maxRows
176+ for ( ; ; ) {
177+ /** Timed around the whole call, so the synchronous-replication wait at commit counts. */
178+ const pageStartedAt = performance . now ( )
179+ let page : PageResult
180+ try {
181+ page = await retirePage ( sql , rowLimit )
182+ } catch ( error ) {
183+ if ( ! ( error instanceof PageMutationTimeout ) ) throw error
184+ /** The timed-out page rolled back with its cursor, so it is retried with fewer rows. */
185+ if ( rowLimit <= ROW_LIMIT . min ) throw error . timeout
186+ ceiling = halve ( rowLimit )
187+ rowLimit = ceiling
188+ logger . warn ( 'Search retirement page timed out; retrying with fewer rows' , { rowLimit } )
189+ await pause ( performance . now ( ) - pageStartedAt )
190+ continue
191+ }
192+ if ( page . done ) break
193+ const pageMs = performance . now ( ) - pageStartedAt
194+ batches ++
195+ mutated += page . mutated
196+ if ( page . transition ) {
197+ /** A phase change may include a full recheck, which says nothing about page cost. */
198+ logger . info ( 'Search retirement phase changed' , {
199+ transition : page . transition ,
200+ batches,
201+ mutated,
202+ } )
203+ } else if ( pageMs > SLOW_PAGE_MS ) {
204+ rowLimit = halve ( rowLimit )
205+ logger . warn ( 'Search retirement page was slow; halving the row limit' , {
206+ pageMs : Math . round ( pageMs ) ,
207+ rowLimit,
208+ } )
209+ } else if ( pageMs < FAST_PAGE_MS ) {
210+ rowLimit = Math . min ( ceiling , rowLimit * 2 )
211+ }
212+ if ( batches % 10 === 0 ) {
213+ logger . info ( 'Search embedding retirement progress' , {
214+ batches,
215+ phase : page . phase ,
216+ afterId : page . afterId ,
217+ mutated,
218+ rowLimit,
219+ elapsedMs : Date . now ( ) - startedAt ,
220+ } )
193221 }
194- logger . info ( 'Selected Search knowledge bases retired' , {
195- batches,
196- mutated,
197- elapsedMs : Date . now ( ) - startedAt ,
198- } )
199- } ,
222+ await pause ( pageMs )
223+ }
224+ logger . info ( 'Selected Search knowledge bases retired' , {
225+ batches,
226+ mutated,
227+ elapsedMs : Date . now ( ) - startedAt ,
228+ } )
200229}
201230
202231/**
@@ -380,17 +409,36 @@ async function validateTargetMarkers(tx: TransactionSql): Promise<void> {
380409 }
381410}
382411
383- /** The standalone entry resumes the deployment cursor and journals only a completed retirement. */
412+ /**
413+ * The operator entry: resumes the saved cursor and, with `--maintenance`, also rebuilds the indexes,
414+ * vacuums, and journals the completed cleanup. `--pause-ratio` and `--max-rows` set the pacing.
415+ */
384416if ( import . meta. main ) {
417+ const { values } = parseArgs ( {
418+ options : {
419+ maintenance : { type : 'boolean' , default : false } ,
420+ 'pause-ratio' : { type : 'string' } ,
421+ 'max-rows' : { type : 'string' } ,
422+ } ,
423+ } )
424+ const pacing : RetirementPacing = {
425+ pauseRatio : Number ( values [ 'pause-ratio' ] ?? DEFAULT_RETIREMENT_PACING . pauseRatio ) ,
426+ maxRows : Number ( values [ 'max-rows' ] ?? DEFAULT_RETIREMENT_PACING . maxRows ) ,
427+ }
385428 const url = resolveMigrationDatabaseUrl ( )
386429 if ( ! url ) throw new Error ( 'DATABASE_URL is required for Search retirement' )
387430 const sql = postgres ( url , { max : 1 , max_lifetime : null , onnotice : ( ) => undefined } )
388431 try {
389- const { runScriptMigrations } = await import ( '@sim/db/script-migrations/index' )
390- const { retireAllSearchEmbeddingsMigration } = await import (
391- '@sim/db/script-migrations/0029_retire_all_search_embeddings'
392- )
393- await runScriptMigrations ( sql , [ retireAllSearchEmbeddingsMigration ] )
432+ if ( values . maintenance ) {
433+ const { runScriptMigrations } = await import ( '@sim/db/script-migrations/index' )
434+ const { retireAllSearchEmbeddings } = await import (
435+ '@sim/db/script-migrations/0029_retire_all_search_embeddings'
436+ )
437+ await runScriptMigrations ( sql , [ retireAllSearchEmbeddings ( pacing ) ] )
438+ } else {
439+ await retireSearchEmbeddings ( sql , pacing )
440+ logger . info ( 'Search retirement pass finished; run with --maintenance off-peak to complete it' )
441+ }
394442 } finally {
395443 await sql . end ( )
396444 }
0 commit comments