11import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings'
2- import { ScriptMigrationDeferred } from '@sim/db/script-migrations/types'
2+ import { maintainSearchRetirementMigration } from '@sim/db/script-migrations/0028_maintain_search_retirement'
3+ import { runScriptMigrations } from '@sim/db/script-migrations/index'
34import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
5+ import { sleep } from '@sim/utils/helpers'
46import { generateId } from '@sim/utils/id'
57import postgres , { type Sql } from 'postgres'
68import { afterAll , beforeAll , beforeEach , describe , expect , it } from 'vitest'
@@ -27,7 +29,9 @@ describe('retiring dormant Search embeddings', () => {
2729 await sql `CREATE TABLE embedding (
2830 id text PRIMARY KEY, knowledge_base_id text REFERENCES knowledge_base(id),
2931 document_id text REFERENCES document(id))`
30- await sql `CREATE TABLE embedding_search (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)`
32+ await sql `CREATE TABLE embedding_search (
33+ id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE,
34+ vector public.vector(3) NOT NULL DEFAULT '[1,2,3]')`
3135 await sql `CREATE TABLE embedding_keyword_search (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)`
3236 await sql `CREATE TABLE embedding_keyword_tin (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)`
3337 await sql `CREATE TABLE embedding_secret_provenance (embedding_id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)`
@@ -43,34 +47,34 @@ describe('retiring dormant Search embeddings', () => {
4347 await sql `TRUNCATE knowledge_base, document, embedding, embedding_search,
4448 embedding_keyword_search, embedding_keyword_tin, embedding_secret_provenance`
4549 await sql `DROP TABLE IF EXISTS search_embedding_cleanup_progress`
50+ await sql `DROP TABLE IF EXISTS script_migrations`
4651 await sql `INSERT INTO knowledge_base VALUES ('search', true), ('ordinary', false)`
4752 await sql `INSERT INTO document (id, knowledge_base_id, processing_queue_token)
4853 VALUES ('search-doc', 'search', 'old-dispatch'), ('ordinary-doc', 'ordinary', 'keep-dispatch')`
4954 await sql `INSERT INTO embedding
5055 SELECT lpad(i::text, 5, '0'), CASE WHEN i % 2 = 0 THEN 'search' ELSE 'ordinary' END,
5156 CASE WHEN i % 2 = 0 THEN 'search-doc' ELSE 'ordinary-doc' END
5257 FROM generate_series(1, 1002) i`
53- await sql `INSERT INTO embedding_search SELECT id FROM embedding`
58+ await sql `INSERT INTO embedding_search (id) SELECT id FROM embedding`
5459 await sql `INSERT INTO embedding_keyword_search SELECT id FROM embedding`
5560 await sql `INSERT INTO embedding_keyword_tin SELECT id FROM embedding`
5661 await sql `INSERT INTO embedding_secret_provenance SELECT id FROM embedding`
5762 } )
5863
5964 async function pass ( ) {
60- try {
61- await retireSearchEmbeddingsMigration . up ( sql )
62- return true
63- } catch ( error ) {
64- if ( error instanceof ScriptMigrationDeferred ) return false
65- throw error
66- }
65+ await runScriptMigrations ( sql , [ retireSearchEmbeddingsMigration ] )
66+ const receipts = await sql `SELECT name FROM script_migrations
67+ WHERE name = '0027_retire_search_embeddings'`
68+ return receipts . length === 1
6769 }
6870
69- it ( 'does nothing without Search data and defers an ambiguous target without changing data' , async ( ) => {
71+ it ( 'does nothing without Search data and fails an ambiguous target without changing data' , async ( ) => {
7072 await sql `UPDATE knowledge_base SET is_search_index = false`
7173 expect ( await pass ( ) ) . toBe ( true )
74+ await sql `DELETE FROM script_migrations`
7275 await sql `UPDATE knowledge_base SET is_search_index = true`
73- expect ( await pass ( ) ) . toBe ( false )
76+ await expect ( pass ( ) ) . rejects . toThrow ( 'ambiguous' )
77+ expect ( await sql `SELECT name FROM script_migrations` ) . toHaveLength ( 0 )
7478 expect ( ( await sql `SELECT count(*)::int AS n FROM embedding` ) [ 0 ] . n ) . toBe ( 1002 )
7579 expect (
7680 ( await sql `SELECT user_excluded FROM document WHERE id = 'search-doc'` ) [ 0 ] . user_excluded
@@ -109,15 +113,22 @@ describe('retiring dormant Search embeddings', () => {
109113 } )
110114
111115 it ( 'rolls back failed pages and resumes the frozen target, retiring documents inserted behind the cursor' , async ( ) => {
116+ await sql `INSERT INTO embedding
117+ SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 26002) i`
112118 await sql `CREATE TABLE deletion_blocker (id text REFERENCES embedding(id))`
113- await sql `INSERT INTO deletion_blocker VALUES ('00002 ')`
119+ await sql `INSERT INTO deletion_blocker VALUES ('26002 ')`
114120 await expect ( pass ( ) ) . rejects . toThrow ( )
115121 const before = await sql `SELECT * FROM search_embedding_cleanup_progress`
122+ const [ remaining ] = await sql `SELECT count(*)::int AS n FROM embedding`
123+ expect ( remaining . n ) . toBeGreaterThan ( 501 )
124+ expect ( remaining . n ) . toBeLessThan ( 26002 )
125+ expect ( await sql `SELECT name FROM script_migrations` ) . toHaveLength ( 0 )
116126 await expect ( pass ( ) ) . rejects . toThrow ( )
117127 expect ( await sql `SELECT * FROM search_embedding_cleanup_progress` ) . toEqual ( before )
118- expect ( ( await sql `SELECT count(*)::int AS n FROM embedding` ) [ 0 ] . n ) . toBe ( 1002 )
128+ expect ( ( await sql `SELECT count(*)::int AS n FROM embedding` ) [ 0 ] . n ) . toBe ( remaining . n )
119129 await sql `DROP TABLE deletion_blocker`
120130 await sql `INSERT INTO document (id, knowledge_base_id) VALUES ('aaa-late-document', 'search')`
131+ await sql `INSERT INTO embedding VALUES ('00000', 'search', 'aaa-late-document')`
121132 await sql `INSERT INTO knowledge_base VALUES ('other-search', true)`
122133 await sql `INSERT INTO document (id, knowledge_base_id, processing_queue_token)
123134 VALUES ('other-search-doc', 'other-search', 'keep-dispatch')`
@@ -135,15 +146,193 @@ describe('retiring dormant Search embeddings', () => {
135146 ) . toEqual ( { user_excluded : false , processing_queue_token : 'keep-dispatch' } )
136147 } )
137148
138- it ( 'defers at the page budget and resumes without skipping remaining chunks' , async ( ) => {
149+ it ( 'finishes beyond the former page budget and journals completion in one invocation' , async ( ) => {
150+ await sql `INSERT INTO document (id, knowledge_base_id)
151+ SELECT 'bulk-doc-' || i::text, 'search' FROM generate_series(1, 51002) i`
139152 await sql `INSERT INTO embedding
140153 SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 51002) i`
141- expect ( await pass ( ) ) . toBe ( false )
142- const [ remaining ] =
143- await sql `SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`
144- expect ( remaining . n ) . toBeGreaterThan ( 0 )
145- expect ( remaining . n ) . toBeLessThan ( 50501 )
146- expect ( await pass ( ) ) . toBe ( true )
154+ await runScriptMigrations ( sql , [
155+ retireSearchEmbeddingsMigration ,
156+ maintainSearchRetirementMigration ,
157+ ] )
158+ expect ( await sql `SELECT name FROM script_migrations` ) . toHaveLength ( 2 )
147159 expect ( ( await sql `SELECT count(*)::int AS n FROM embedding` ) [ 0 ] . n ) . toBe ( 501 )
160+ expect (
161+ (
162+ await sql `SELECT count(*)::int AS n FROM document
163+ WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
164+ ) [ 0 ] . n
165+ ) . toBe ( 0 )
148166 } , 60_000 )
167+
168+ it ( 'rejects inconsistent document ownership before deleting any chunk in the page' , async ( ) => {
169+ await sql `UPDATE embedding SET document_id = 'ordinary-doc' WHERE id = '00002'`
170+ await expect ( pass ( ) ) . rejects . toThrow ( 'Search content changed after retirement' )
171+ expect ( ( await sql `SELECT count(*)::int AS n FROM embedding` ) [ 0 ] . n ) . toBe ( 1002 )
172+ expect ( await sql `SELECT name FROM script_migrations` ) . toHaveLength ( 0 )
173+ await sql `UPDATE embedding SET document_id = 'search-doc' WHERE id = '00002'`
174+ expect ( await pass ( ) ) . toBe ( true )
175+ expect ( ( await sql `SELECT count(*)::int AS n FROM embedding` ) [ 0 ] . n ) . toBe ( 501 )
176+ } )
177+
178+ it ( 'retries a briefly locked document page and completes in the same invocation' , async ( ) => {
179+ const blocker = postgres ( readTestDatabaseUrl ( ) , {
180+ max : 1 ,
181+ connection : { search_path : schema } ,
182+ onnotice : ( ) => undefined ,
183+ } )
184+ let signalLocked ! : ( ) => void
185+ const locked = new Promise < void > ( ( resolve ) => {
186+ signalLocked = resolve
187+ } )
188+ let release ! : ( ) => void
189+ const released = new Promise < void > ( ( resolve ) => {
190+ release = resolve
191+ } )
192+ const holding = blocker . begin ( async ( tx ) => {
193+ await tx `SELECT id FROM document WHERE id = 'search-doc' FOR UPDATE`
194+ signalLocked ( )
195+ await released
196+ } )
197+ try {
198+ await locked
199+ const result = pass ( ) . then (
200+ ( complete ) => ( { complete, error : undefined } ) ,
201+ ( error : unknown ) => ( { complete : false , error } )
202+ )
203+ await sleep ( 1_500 )
204+ release ( )
205+ await holding
206+ expect ( await result ) . toEqual ( { complete : true , error : undefined } )
207+ expect (
208+ ( await sql `SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'` ) [ 0 ]
209+ . n
210+ ) . toBe ( 0 )
211+ } finally {
212+ release ( )
213+ await holding
214+ await blocker . end ( )
215+ }
216+ } )
217+
218+ it ( 'maintains an already-retired index and resumes failed vacuum bookkeeping without rebuilding it again' , async ( ) => {
219+ await pass ( )
220+ await sql `CREATE INDEX retirement_hnsw_idx ON embedding_search USING hnsw (vector public.vector_l2_ops)`
221+ const [ original ] = await sql `SELECT to_regclass('retirement_hnsw_idx')::oid AS oid`
222+ await sql `CREATE FUNCTION interrupt_maintenance_checkpoint() RETURNS trigger LANGUAGE plpgsql AS $$
223+ BEGIN
224+ IF NEW.vacuumed_tables = 2 THEN RAISE EXCEPTION 'interrupt maintenance checkpoint'; END IF;
225+ RETURN NEW;
226+ END $$`
227+ await sql `CREATE TRIGGER interrupt_maintenance_checkpoint BEFORE UPDATE ON search_embedding_cleanup_progress
228+ FOR EACH ROW EXECUTE FUNCTION interrupt_maintenance_checkpoint()`
229+ try {
230+ await expect ( runScriptMigrations ( sql , [ maintainSearchRetirementMigration ] ) ) . rejects . toThrow (
231+ 'interrupt maintenance checkpoint'
232+ )
233+ const [ rebuilt ] =
234+ await sql `SELECT c.oid, i.indisvalid FROM pg_class c JOIN pg_index i ON i.indexrelid = c.oid
235+ WHERE c.oid = to_regclass('retirement_hnsw_idx')`
236+ expect ( rebuilt . oid ) . not . toBe ( original . oid )
237+ expect ( rebuilt . indisvalid ) . toBe ( true )
238+ expect (
239+ await sql `SELECT reindexed_through, vacuumed_tables FROM search_embedding_cleanup_progress`
240+ ) . toEqual ( [ { reindexed_through : 'retirement_hnsw_idx' , vacuumed_tables : 1 } ] )
241+ expect (
242+ await sql `SELECT name FROM script_migrations WHERE name = '0028_maintain_search_retirement'`
243+ ) . toHaveLength ( 0 )
244+ await sql `DROP TRIGGER interrupt_maintenance_checkpoint ON search_embedding_cleanup_progress`
245+ await runScriptMigrations ( sql , [ maintainSearchRetirementMigration ] )
246+ expect ( ( await sql `SELECT to_regclass('retirement_hnsw_idx')::oid AS oid` ) [ 0 ] . oid ) . toBe (
247+ rebuilt . oid
248+ )
249+ expect (
250+ await sql `SELECT name FROM script_migrations WHERE name = '0028_maintain_search_retirement'`
251+ ) . toHaveLength ( 1 )
252+ expect (
253+ (
254+ await sql `SELECT reltuples::int AS n FROM pg_class WHERE oid = to_regclass('embedding_search')`
255+ ) [ 0 ] . n
256+ ) . toBe ( 501 )
257+ await runScriptMigrations ( sql , [ maintainSearchRetirementMigration ] )
258+ expect ( ( await sql `SELECT to_regclass('retirement_hnsw_idx')::oid AS oid` ) [ 0 ] . oid ) . toBe (
259+ rebuilt . oid
260+ )
261+ } finally {
262+ await sql `DROP TRIGGER IF EXISTS interrupt_maintenance_checkpoint ON search_embedding_cleanup_progress`
263+ await sql `DROP FUNCTION interrupt_maintenance_checkpoint()`
264+ await sql `DROP INDEX retirement_hnsw_idx`
265+ }
266+ } )
267+
268+ it ( 'recovers a canceled concurrent rebuild and completes cleanup plus maintenance in one rerun' , async ( ) => {
269+ await sql `CREATE INDEX retirement_hnsw_idx ON embedding_search USING hnsw (vector public.vector_l2_ops)`
270+ const blocker = postgres ( readTestDatabaseUrl ( ) , {
271+ max : 1 ,
272+ connection : { search_path : schema } ,
273+ onnotice : ( ) => undefined ,
274+ } )
275+ let signalLocked ! : ( ) => void
276+ const locked = new Promise < void > ( ( resolve ) => {
277+ signalLocked = resolve
278+ } )
279+ let release ! : ( ) => void
280+ const released = new Promise < void > ( ( resolve ) => {
281+ release = resolve
282+ } )
283+ const holding = blocker . begin ( async ( tx ) => {
284+ await tx `UPDATE embedding_search SET vector = '[3,2,1]' WHERE id = '00001'`
285+ signalLocked ( )
286+ await released
287+ } )
288+ const [ { pid } ] = await sql `SELECT pg_backend_pid() AS pid`
289+ const migrations = [ retireSearchEmbeddingsMigration , maintainSearchRetirementMigration ]
290+ let outcome : Promise < unknown > | undefined
291+ try {
292+ await locked
293+ outcome = runScriptMigrations ( sql , migrations ) . then (
294+ ( ) => undefined ,
295+ ( error : unknown ) => error
296+ )
297+ let waiting = false
298+ for ( let attempt = 0 ; attempt < 200 ; attempt ++ ) {
299+ const [ progress ] =
300+ await admin `SELECT phase FROM pg_stat_progress_create_index WHERE pid = ${ pid } `
301+ if ( progress ?. phase === 'waiting for writers before build' ) {
302+ waiting = true
303+ break
304+ }
305+ await sleep ( 25 )
306+ }
307+ expect ( waiting ) . toBe ( true )
308+ await admin `SELECT pg_cancel_backend(${ pid } )`
309+ expect ( await outcome ) . toMatchObject ( { code : '57014' } )
310+ release ( )
311+ await holding
312+ const leftovers =
313+ await sql `SELECT c.relname FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid
314+ WHERE i.indrelid = to_regclass('embedding_search') AND NOT i.indisvalid`
315+ expect ( leftovers . length ) . toBeGreaterThan ( 0 )
316+ expect (
317+ ( await sql `SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'` ) [ 0 ]
318+ . n
319+ ) . toBe ( 0 )
320+ await runScriptMigrations ( sql , migrations )
321+ expect (
322+ await sql `SELECT c.relname FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid
323+ WHERE i.indrelid = to_regclass('embedding_search') AND NOT i.indisvalid`
324+ ) . toHaveLength ( 0 )
325+ expect ( ( await sql `SELECT count(*)::int AS n FROM embedding_search` ) [ 0 ] . n ) . toBe ( 501 )
326+ expect (
327+ await sql `SELECT name FROM script_migrations WHERE name = '0028_maintain_search_retirement'`
328+ ) . toHaveLength ( 1 )
329+ } finally {
330+ release ( )
331+ await holding
332+ await admin `SELECT pg_cancel_backend(${ pid } )`
333+ await outcome
334+ await blocker . end ( )
335+ await sql `DROP INDEX retirement_hnsw_idx`
336+ }
337+ } )
149338} )
0 commit comments