diff --git a/plugins/catalog-backend/src/database/operations/stitcher/performStitching.test.ts b/plugins/catalog-backend/src/database/operations/stitcher/performStitching.test.ts index eca65c4985..c6349379e7 100644 --- a/plugins/catalog-backend/src/database/operations/stitcher/performStitching.test.ts +++ b/plugins/catalog-backend/src/database/operations/stitcher/performStitching.test.ts @@ -90,12 +90,7 @@ it.each(databases.eachSupportedId())( knex, logger, entityRef: 'k:ns/n', - stitchTicket: ( - await knex('stitch_queue') - .where('entity_ref', 'k:ns/n') - .select('stitch_ticket') - .first() - )?.stitch_ticket, + stitchTicket: await getStitchTicket(knex, 'k:ns/n'), }); entities = await knex('final_entities'); @@ -185,12 +180,7 @@ it.each(databases.eachSupportedId())( knex, logger, entityRef: 'k:ns/n', - stitchTicket: ( - await knex('stitch_queue') - .where('entity_ref', 'k:ns/n') - .select('stitch_ticket') - .first() - )?.stitch_ticket, + stitchTicket: await getStitchTicket(knex, 'k:ns/n'), }); entities = await knex('final_entities'); @@ -218,12 +208,7 @@ it.each(databases.eachSupportedId())( knex, logger, entityRef: 'k:ns/n', - stitchTicket: ( - await knex('stitch_queue') - .where('entity_ref', 'k:ns/n') - .select('stitch_ticket') - .first() - )?.stitch_ticket, + stitchTicket: await getStitchTicket(knex, 'k:ns/n'), }); entities = await knex('final_entities'); @@ -357,12 +342,7 @@ describe.each(databases.eachSupportedId())( knex, logger: stitchLogger, entityRef: 'k:ns/n', - stitchTicket: ( - await knex('stitch_queue') - .where('entity_ref', 'k:ns/n') - .select('stitch_ticket') - .first() - )?.stitch_ticket, + stitchTicket: await getStitchTicket(knex, 'k:ns/n'), }), ).resolves.toBe('changed'); @@ -395,12 +375,7 @@ describe.each(databases.eachSupportedId())( // First stitch: create the final_entities row with a valid ticket await markForStitching({ knex, entityRefs: ['k:ns/n'] }); - const validTicket = ( - await knex('stitch_queue') - .where('entity_ref', 'k:ns/n') - .select('stitch_ticket') - .first() - )?.stitch_ticket; + const validTicket = await getStitchTicket(knex, 'k:ns/n'); const result1 = await performStitching({ knex, @@ -427,12 +402,7 @@ describe.each(databases.eachSupportedId())( }); await markForStitching({ knex, entityRefs: ['k:ns/n'] }); - const freshTicket = ( - await knex('stitch_queue') - .where('entity_ref', 'k:ns/n') - .select('stitch_ticket') - .first() - )?.stitch_ticket; + const freshTicket = await getStitchTicket(knex, 'k:ns/n'); // Attempt to stitch with a WRONG ticket (simulating a stale worker) const result2 = await performStitching({ @@ -466,3 +436,17 @@ describe.each(databases.eachSupportedId())( }); }, ); + +async function getStitchTicket( + knex: import('knex').Knex, + entityRef: string, +): Promise { + const row = await knex('stitch_queue') + .where('entity_ref', entityRef) + .select('stitch_ticket') + .first(); + if (!row) { + throw new Error(`No stitch_queue entry for ${entityRef}`); + } + return row.stitch_ticket; +} diff --git a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts index b473891a7f..bd6eb0bd47 100644 --- a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts +++ b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts @@ -72,7 +72,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - createHash: () => hash, scheduler: mockServices.scheduler(), events: mockServices.events.mock(), @@ -142,7 +141,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, events: mockServices.events.mock(), @@ -228,7 +226,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, events: mockServices.events.mock(), @@ -308,7 +305,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, events: mockServices.events.mock(), @@ -370,7 +366,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, @@ -488,7 +483,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, @@ -596,7 +590,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, @@ -682,7 +675,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, @@ -773,7 +765,6 @@ describe('DefaultCatalogProcessingEngine', () => { processingDatabase: db, knex: {} as any, orchestrator: orchestrator, - scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100,