compute deltas based on hashes just like full mutations

Signed-off-by: Fredrik Adelöw <freben@gmail.com>
This commit is contained in:
Fredrik Adelöw
2024-11-19 16:47:10 +01:00
parent d589a90b10
commit f159b258c9
3 changed files with 194 additions and 21 deletions
+5
View File
@@ -0,0 +1,5 @@
---
'@backstage/plugin-catalog-backend': patch
---
Compute deltas more efficiently, which generally leads to less wasted processing cycles
@@ -19,13 +19,14 @@ import {
TestDatabaseId,
TestDatabases,
} from '@backstage/backend-test-utils';
import { Entity } from '@backstage/catalog-model';
import { Entity, stringifyEntityRef } from '@backstage/catalog-model';
import { Knex } from 'knex';
import * as uuid from 'uuid';
import { DefaultProviderDatabase } from './DefaultProviderDatabase';
import { applyDatabaseMigrations } from './migrations';
import { DbRefreshStateReferencesRow, DbRefreshStateRow } from './tables';
import { LoggerService } from '@backstage/backend-plugin-api';
import { generateStableHash } from './util';
jest.setTimeout(60_000);
@@ -690,13 +691,8 @@ describe('DefaultProviderDatabase', () => {
it.each(databases.eachSupportedId())(
'should successfully fall back from batch to individual mode on conflicts, %p',
async databaseId => {
const fakeLogger = {
debug: jest.fn(),
};
const { knex, db } = await createDatabase(
databaseId,
fakeLogger as any,
);
const fakeLogger = mockServices.logger.mock();
const { knex, db } = await createDatabase(databaseId, fakeLogger);
await createLocations(knex, ['component:default/a']);
@@ -744,11 +740,8 @@ describe('DefaultProviderDatabase', () => {
it.each(databases.eachSupportedId())(
'should gracefully handle accidental duplicate refresh state references when deletion happens during a full sync, %p',
async databaseId => {
const fakeLogger = { debug: jest.fn() };
const { knex, db } = await createDatabase(
databaseId,
fakeLogger as any,
);
const fakeLogger = mockServices.logger.mock();
const { knex, db } = await createDatabase(databaseId, fakeLogger);
await createLocations(knex, ['component:default/a']);
@@ -773,5 +766,162 @@ describe('DefaultProviderDatabase', () => {
expect(state).toEqual([]);
},
);
it.each(databases.eachSupportedId())(
'should properly translate deltas into add/update/remove, %p',
async databaseId => {
const fakeLogger = mockServices.logger.mock();
const { knex, db } = await createDatabase(databaseId, fakeLogger);
const entity1Before: Entity = {
apiVersion: '1',
kind: 'k',
metadata: { namespace: 'ns', name: 'n1' },
};
const entity1After: Entity = {
...entity1Before,
apiVersion: '2',
};
const entity2Before: Entity = {
apiVersion: '1',
kind: 'k',
metadata: { namespace: 'ns', name: 'n2' },
};
const entity2After: Entity = {
...entity2Before,
apiVersion: '2',
};
const entity3: Entity = {
apiVersion: '1',
kind: 'k',
metadata: { namespace: 'ns', name: 'n3' },
};
const entity4: Entity = {
apiVersion: '1',
kind: 'k',
metadata: { namespace: 'ns', name: 'n4' },
};
await insertRefreshStateRow(knex, {
entity_id: 'id1',
entity_ref: stringifyEntityRef(entity1Before),
last_discovery_at: new Date(),
next_update_at: new Date(),
errors: '[]',
unprocessed_entity: JSON.stringify(entity1Before),
unprocessed_hash: generateStableHash(entity1Before),
});
await insertRefRow(knex, {
source_key: 'my-provider',
target_entity_ref: stringifyEntityRef(entity1Before),
});
await insertRefreshStateRow(knex, {
entity_id: 'id2',
entity_ref: stringifyEntityRef(entity2Before),
last_discovery_at: new Date(),
next_update_at: new Date(),
errors: '[]',
unprocessed_entity: JSON.stringify(entity2Before),
unprocessed_hash: generateStableHash(entity2After), // lie about the hash!
});
await insertRefRow(knex, {
source_key: 'my-provider',
target_entity_ref: stringifyEntityRef(entity2Before),
});
await insertRefreshStateRow(knex, {
entity_id: 'id4',
entity_ref: stringifyEntityRef(entity4),
last_discovery_at: new Date(),
next_update_at: new Date(),
errors: '[]',
unprocessed_entity: JSON.stringify(entity4),
unprocessed_hash: generateStableHash(entity4),
});
await insertRefRow(knex, {
source_key: 'my-provider',
target_entity_ref: stringifyEntityRef(entity4),
});
await db.transaction(async tx => {
await db.replaceUnprocessedEntities(tx, {
type: 'delta',
sourceKey: 'my-provider',
added: [
// we used the right hashes for entity1, so this will turn into an update
{ entity: entity1After },
// we lied about the hash for entity2, so this will become a no-op
{ entity: entity2After },
// this didn't exist, so will become an add
{ entity: entity3 },
],
removed: [{ entityRef: stringifyEntityRef(entity4) }],
});
});
const state = await knex<DbRefreshStateRow>('refresh_state')
.select(['entity_ref', 'unprocessed_entity', 'unprocessed_hash'])
.orderBy('entity_ref');
expect(state).toEqual([
{
entity_ref: stringifyEntityRef(entity1After),
unprocessed_entity: JSON.stringify(entity1After),
unprocessed_hash: generateStableHash(entity1After),
},
{
entity_ref: stringifyEntityRef(entity2After),
unprocessed_entity: JSON.stringify(entity2Before), // didn't change
unprocessed_hash: generateStableHash(entity2After),
},
{
entity_ref: stringifyEntityRef(entity3),
unprocessed_entity: JSON.stringify(entity3),
unprocessed_hash: generateStableHash(entity3),
},
]);
},
);
it.each(databases.eachSupportedId())(
'can handle large deltas without exploding, %p',
async databaseId => {
const fakeLogger = mockServices.logger.mock();
const { knex, db } = await createDatabase(databaseId, fakeLogger);
const count = 10000;
const padded = (n: number) => String(n).padStart(8, '0');
const entities = Array.from({ length: count }, (_, i) => ({
entity: {
apiVersion: '1',
kind: 'k',
metadata: { namespace: 'ns', name: padded(i) },
},
}));
await db.transaction(async tx => {
await db.replaceUnprocessedEntities(tx, {
type: 'delta',
sourceKey: 'my-provider',
added: entities,
removed: [],
});
});
const state = await knex<DbRefreshStateRow>('refresh_state')
.select(['entity_ref', 'unprocessed_entity', 'unprocessed_hash'])
.orderBy('entity_ref');
expect(state).toHaveLength(count);
expect(state[0]).toEqual({
entity_ref: stringifyEntityRef(entities[0].entity),
unprocessed_entity: JSON.stringify(entities[0].entity),
unprocessed_hash: generateStableHash(entities[0].entity),
});
},
);
});
});
@@ -213,14 +213,32 @@ export class DefaultProviderDatabase implements ProviderDatabase {
toRemove: string[];
}> {
if (options.type === 'delta') {
return {
toAdd: [],
toUpsert: options.added.map(e => ({
deferred: e,
hash: generateStableHash(e.entity),
})),
toRemove: options.removed.map(e => e.entityRef),
};
const toAdd = new Array<{ deferred: DeferredEntity; hash: string }>();
const toUpsert = new Array<{ deferred: DeferredEntity; hash: string }>();
const toRemove = options.removed.map(e => e.entityRef);
for (const chunk of lodash.chunk(options.added, 1000)) {
const entityRefs = chunk.map(e => stringifyEntityRef(e.entity));
const rows = await tx<DbRefreshStateRow>('refresh_state')
.select(['entity_ref', 'unprocessed_hash'])
.whereIn('entity_ref', entityRefs);
const oldHashes = new Map(
rows.map(row => [row.entity_ref, row.unprocessed_hash]),
);
chunk.forEach((deferred, i) => {
const entityRef = entityRefs[i];
const newHash = generateStableHash(deferred.entity);
const oldHash = oldHashes.get(entityRef);
if (oldHash === undefined) {
toAdd.push({ deferred, hash: newHash });
} else if (newHash !== oldHash) {
toUpsert.push({ deferred, hash: newHash });
}
});
}
return { toAdd, toUpsert, toRemove };
}
// Grab all of the existing references from the same source, and their locationKeys as well