diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.test.ts b/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.test.ts index 155b45cacb..91975edfb5 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.test.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.test.ts @@ -122,123 +122,123 @@ describe.each(databases.eachSupportedId())( expect(typeof result.total).toBe('number'); }); - it('createMarkEntities handles existing and new refs correctly', async () => { - const knex = await databases.init(databaseId); - await knex.migrate.latest({ directory: migrationsDir }); + it('createMarkEntities handles existing and new refs correctly', async () => { + const knex = await databases.init(databaseId); + await knex.migrate.latest({ directory: migrationsDir }); - const manager = new IncrementalIngestionDatabaseManager({ client: knex }); - const { ingestionId } = (await manager.createProviderIngestionRecord( - 'testProvider', - ))!; + const manager = new IncrementalIngestionDatabaseManager({ client: knex }); + const { ingestionId } = (await manager.createProviderIngestionRecord( + 'testProvider', + ))!; - const markId1 = uuid(); - await manager.createMark({ - record: { - id: markId1, - ingestion_id: ingestionId, - sequence: 1, - cursor: { data: 1 }, - }, + const markId1 = uuid(); + await manager.createMark({ + record: { + id: markId1, + ingestion_id: ingestionId, + sequence: 1, + cursor: { data: 1 }, + }, + }); + + const makeEntity = (name: string): DeferredEntity => ({ + entity: { + apiVersion: 'backstage.io/v1alpha1', + kind: 'Component', + metadata: { namespace: 'default', name }, + }, + }); + + // First batch: create 3 entities + await manager.createMarkEntities(markId1, [ + makeEntity('a'), + makeEntity('b'), + makeEntity('c'), + ]); + + const rows1 = await knex('ingestion_mark_entities').select('ref'); + expect(rows1).toHaveLength(3); + + // Second batch with overlap: b and c already exist, d is new. + // Existing refs should be updated to the new mark, new refs inserted. + const markId2 = uuid(); + await manager.createMark({ + record: { + id: markId2, + ingestion_id: ingestionId, + sequence: 2, + cursor: { data: 2 }, + }, + }); + + await manager.createMarkEntities(markId2, [ + makeEntity('b'), + makeEntity('c'), + makeEntity('d'), + ]); + + const rows2 = await knex('ingestion_mark_entities') + .select('ref', 'ingestion_mark_id') + .orderBy('ref'); + expect(rows2).toHaveLength(4); + + // a stays on markId1, b and c moved to markId2, d is new on markId2 + expect( + rows2.find(r => r.ref === 'component:default/a')?.ingestion_mark_id, + ).toBe(markId1); + expect( + rows2.find(r => r.ref === 'component:default/b')?.ingestion_mark_id, + ).toBe(markId2); + expect( + rows2.find(r => r.ref === 'component:default/c')?.ingestion_mark_id, + ).toBe(markId2); + expect( + rows2.find(r => r.ref === 'component:default/d')?.ingestion_mark_id, + ).toBe(markId2); }); - const makeEntity = (name: string): DeferredEntity => ({ - entity: { - apiVersion: 'backstage.io/v1alpha1', - kind: 'Component', - metadata: { namespace: 'default', name }, - }, + it('deleteEntityRecordsByRef removes matching refs', async () => { + const knex = await databases.init(databaseId); + await knex.migrate.latest({ directory: migrationsDir }); + + const manager = new IncrementalIngestionDatabaseManager({ client: knex }); + const { ingestionId } = (await manager.createProviderIngestionRecord( + 'testProvider', + ))!; + + const markId = uuid(); + await manager.createMark({ + record: { + id: markId, + ingestion_id: ingestionId, + sequence: 1, + cursor: { data: 1 }, + }, + }); + + const makeEntity = (name: string): DeferredEntity => ({ + entity: { + apiVersion: 'backstage.io/v1alpha1', + kind: 'Component', + metadata: { namespace: 'default', name }, + }, + }); + + await manager.createMarkEntities(markId, [ + makeEntity('x'), + makeEntity('y'), + makeEntity('z'), + ]); + + // Delete two of the three + await manager.deleteEntityRecordsByRef([ + { entityRef: 'component:default/x' }, + { entityRef: 'component:default/z' }, + ]); + + const remaining = await knex('ingestion_mark_entities').select('ref'); + expect(remaining).toHaveLength(1); + expect(remaining[0].ref).toBe('component:default/y'); }); - - // First batch: create 3 entities - await manager.createMarkEntities(markId1, [ - makeEntity('a'), - makeEntity('b'), - makeEntity('c'), - ]); - - const rows1 = await knex('ingestion_mark_entities').select('ref'); - expect(rows1).toHaveLength(3); - - // Second batch with overlap: b and c already exist, d is new. - // Existing refs should be updated to the new mark, new refs inserted. - const markId2 = uuid(); - await manager.createMark({ - record: { - id: markId2, - ingestion_id: ingestionId, - sequence: 2, - cursor: { data: 2 }, - }, - }); - - await manager.createMarkEntities(markId2, [ - makeEntity('b'), - makeEntity('c'), - makeEntity('d'), - ]); - - const rows2 = await knex('ingestion_mark_entities') - .select('ref', 'ingestion_mark_id') - .orderBy('ref'); - expect(rows2).toHaveLength(4); - - // a stays on markId1, b and c moved to markId2, d is new on markId2 - expect( - rows2.find(r => r.ref === 'component:default/a')?.ingestion_mark_id, - ).toBe(markId1); - expect( - rows2.find(r => r.ref === 'component:default/b')?.ingestion_mark_id, - ).toBe(markId2); - expect( - rows2.find(r => r.ref === 'component:default/c')?.ingestion_mark_id, - ).toBe(markId2); - expect( - rows2.find(r => r.ref === 'component:default/d')?.ingestion_mark_id, - ).toBe(markId2); - }); - - it('deleteEntityRecordsByRef removes matching refs', async () => { - const knex = await databases.init(databaseId); - await knex.migrate.latest({ directory: migrationsDir }); - - const manager = new IncrementalIngestionDatabaseManager({ client: knex }); - const { ingestionId } = (await manager.createProviderIngestionRecord( - 'testProvider', - ))!; - - const markId = uuid(); - await manager.createMark({ - record: { - id: markId, - ingestion_id: ingestionId, - sequence: 1, - cursor: { data: 1 }, - }, - }); - - const makeEntity = (name: string): DeferredEntity => ({ - entity: { - apiVersion: 'backstage.io/v1alpha1', - kind: 'Component', - metadata: { namespace: 'default', name }, - }, - }); - - await manager.createMarkEntities(markId, [ - makeEntity('x'), - makeEntity('y'), - makeEntity('z'), - ]); - - // Delete two of the three - await manager.deleteEntityRecordsByRef([ - { entityRef: 'component:default/x' }, - { entityRef: 'component:default/z' }, - ]); - - const remaining = await knex('ingestion_mark_entities').select('ref'); - expect(remaining).toHaveLength(1); - expect(remaining[0].ref).toBe('component:default/y'); - }); -}, + }, ); diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts b/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts index 78dc646d72..ecf821364f 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts @@ -94,9 +94,11 @@ export class IncrementalIngestionDatabaseManager { const allIds = ids.map(entry => entry.id); if (this.client.client.config.client === 'pg') { - return await tx('ingestion_mark_entities') - .delete() - .whereRaw('?? = ANY(?)', ['id', allIds]); + return await this.whereInArray( + tx('ingestion_mark_entities').delete(), + 'id', + allIds, + ); } let deleted = 0;