Reuse whereInArray helper in deleteMarkEntities and fix test formatting

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
Signed-off-by: Fredrik Adelöw <freben@gmail.com>
This commit is contained in:
Fredrik Adelöw
2026-05-19 13:30:21 +02:00
parent 32f0dfe7f8
commit 0f8c977018
2 changed files with 119 additions and 117 deletions
@@ -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');
});
},
},
);
@@ -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;