catalog-backend: refactor defferal of entites + fix dissapearing ref bug
Signed-off-by: Patrik Oldsberg <poldsberg@gmail.com>
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
'@backstage/plugin-catalog-backend': patch
|
||||
---
|
||||
|
||||
Fixed a bug where internal references within the catalog were broken when new entities where added through entity providers, such as registering a new location or adding one in configuration. These broken references then caused some entities to be incorrectly marked as orphaned and prevented refresh from working properly.
|
||||
@@ -66,123 +66,6 @@ describe('Default Processing Database', () => {
|
||||
await db<DbRefreshStateRow>('refresh_state').insert(ref);
|
||||
};
|
||||
|
||||
describe('addUprocessedEntities', () => {
|
||||
function mockEntity(name: string, type: string): Entity {
|
||||
return {
|
||||
apiVersion: '1',
|
||||
kind: 'Component',
|
||||
metadata: {
|
||||
name,
|
||||
},
|
||||
spec: {
|
||||
type,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
it.each(databases.eachSupportedId())(
|
||||
'updates refresh state with varying location keys, %p',
|
||||
async databaseId => {
|
||||
const mockWarn = jest.fn();
|
||||
const { db } = await createDatabase(databaseId, {
|
||||
debug: jest.fn(),
|
||||
warn: mockWarn,
|
||||
} as unknown as Logger);
|
||||
await db.transaction(async tx => {
|
||||
const knexTx = tx as Knex.Transaction;
|
||||
|
||||
const steps = [
|
||||
{
|
||||
locationKey: undefined,
|
||||
expectedLocationKey: null,
|
||||
type: 'a',
|
||||
expectedType: 'a',
|
||||
},
|
||||
{
|
||||
locationKey: undefined,
|
||||
expectedLocationKey: null,
|
||||
type: 'b',
|
||||
expectedType: 'b',
|
||||
},
|
||||
{
|
||||
locationKey: 'x',
|
||||
expectedLocationKey: 'x',
|
||||
type: 'c',
|
||||
expectedType: 'c',
|
||||
},
|
||||
{
|
||||
locationKey: 'y',
|
||||
expectedLocationKey: 'x',
|
||||
type: 'd',
|
||||
expectedType: 'c',
|
||||
expectConflict: true,
|
||||
},
|
||||
{
|
||||
locationKey: undefined,
|
||||
expectedLocationKey: 'x',
|
||||
type: 'e',
|
||||
expectedType: 'c',
|
||||
expectConflict: true,
|
||||
},
|
||||
{
|
||||
locationKey: 'x',
|
||||
expectedLocationKey: 'x',
|
||||
type: 'f',
|
||||
expectedType: 'f',
|
||||
},
|
||||
];
|
||||
for (const step of steps) {
|
||||
mockWarn.mockClear();
|
||||
|
||||
await db.addUnprocessedEntities(tx, {
|
||||
sourceKey: 'testing',
|
||||
entities: [
|
||||
{
|
||||
entity: mockEntity('1', step.type),
|
||||
locationKey: step.locationKey,
|
||||
},
|
||||
],
|
||||
});
|
||||
|
||||
if (step.expectConflict) {
|
||||
// eslint-disable-next-line jest/no-conditional-expect
|
||||
expect(mockWarn).toHaveBeenCalledWith(
|
||||
expect.stringMatching(/^Detected conflicting entityRef/),
|
||||
);
|
||||
} else {
|
||||
// eslint-disable-next-line jest/no-conditional-expect
|
||||
expect(mockWarn).not.toHaveBeenCalled();
|
||||
}
|
||||
|
||||
const entities = await knexTx<DbRefreshStateRow>(
|
||||
'refresh_state',
|
||||
).select();
|
||||
expect(entities).toEqual([
|
||||
expect.objectContaining({
|
||||
entity_ref: 'component:default/1',
|
||||
location_key: step.expectedLocationKey,
|
||||
}),
|
||||
]);
|
||||
const entity = JSON.parse(entities[0].unprocessed_entity) as Entity;
|
||||
expect(entity.spec?.type).toEqual(step.expectedType);
|
||||
|
||||
await expect(
|
||||
knexTx<DbRefreshStateReferencesRow>(
|
||||
'refresh_state_references',
|
||||
).select(),
|
||||
).resolves.toEqual([
|
||||
expect.objectContaining({
|
||||
source_key: 'testing',
|
||||
target_entity_ref: 'component:default/1',
|
||||
}),
|
||||
]);
|
||||
}
|
||||
});
|
||||
},
|
||||
60_000,
|
||||
);
|
||||
});
|
||||
|
||||
describe('updateProcessedEntity', () => {
|
||||
let id: string;
|
||||
let processedEntity: Entity;
|
||||
@@ -408,6 +291,139 @@ describe('Default Processing Database', () => {
|
||||
},
|
||||
60_000,
|
||||
);
|
||||
|
||||
it.each(databases.eachSupportedId())(
|
||||
'updates unprocessed entities with varying location keys, %p',
|
||||
async databaseId => {
|
||||
const mockLogger = {
|
||||
debug: jest.fn(),
|
||||
error: jest.fn(),
|
||||
warn: jest.fn(),
|
||||
};
|
||||
const { knex, db } = await createDatabase(
|
||||
databaseId,
|
||||
mockLogger as unknown as Logger,
|
||||
);
|
||||
|
||||
await insertRefreshStateRow(knex, {
|
||||
entity_id: id,
|
||||
entity_ref: 'location:default/fakelocation',
|
||||
unprocessed_entity: '{}',
|
||||
processed_entity: '{}',
|
||||
errors: '[]',
|
||||
next_update_at: '2021-04-01 13:37:00',
|
||||
last_discovery_at: '2021-04-01 13:37:00',
|
||||
});
|
||||
|
||||
await db.transaction(async tx => {
|
||||
const knexTx = tx as Knex.Transaction;
|
||||
|
||||
const steps = [
|
||||
{
|
||||
locationKey: undefined,
|
||||
expectedLocationKey: null,
|
||||
testKey: 'a',
|
||||
expectedTestKey: 'a',
|
||||
},
|
||||
{
|
||||
locationKey: undefined,
|
||||
expectedLocationKey: null,
|
||||
testKey: 'b',
|
||||
expectedTestKey: 'b',
|
||||
},
|
||||
{
|
||||
locationKey: 'x',
|
||||
expectedLocationKey: 'x',
|
||||
testKey: 'c',
|
||||
expectedTestKey: 'c',
|
||||
},
|
||||
{
|
||||
locationKey: 'y',
|
||||
expectedLocationKey: 'x',
|
||||
testKey: 'd',
|
||||
expectedTestKey: 'c',
|
||||
expectConflict: true,
|
||||
},
|
||||
{
|
||||
locationKey: undefined,
|
||||
expectedLocationKey: 'x',
|
||||
testKey: 'e',
|
||||
expectedTestKey: 'c',
|
||||
expectConflict: true,
|
||||
},
|
||||
{
|
||||
locationKey: 'x',
|
||||
expectedLocationKey: 'x',
|
||||
testKey: 'f',
|
||||
expectedTestKey: 'f',
|
||||
},
|
||||
];
|
||||
for (const step of steps) {
|
||||
mockLogger.debug.mockClear();
|
||||
mockLogger.warn.mockClear();
|
||||
mockLogger.error.mockClear();
|
||||
|
||||
await db.updateProcessedEntity(tx, {
|
||||
id,
|
||||
processedEntity,
|
||||
resultHash: '',
|
||||
relations: [],
|
||||
deferredEntities: [
|
||||
{
|
||||
entity: {
|
||||
apiVersion: '1',
|
||||
kind: 'Component',
|
||||
metadata: {
|
||||
name: '1',
|
||||
},
|
||||
spec: {
|
||||
type: step.testKey,
|
||||
},
|
||||
},
|
||||
locationKey: step.locationKey,
|
||||
},
|
||||
],
|
||||
});
|
||||
|
||||
if (step.expectConflict) {
|
||||
// eslint-disable-next-line jest/no-conditional-expect
|
||||
expect(mockLogger.warn).toHaveBeenCalledWith(
|
||||
expect.stringMatching(/^Detected conflicting entityRef/),
|
||||
);
|
||||
} else {
|
||||
// eslint-disable-next-line jest/no-conditional-expect
|
||||
expect(mockLogger.warn).not.toHaveBeenCalled();
|
||||
}
|
||||
|
||||
const [, unprocessed] = await knexTx<DbRefreshStateRow>(
|
||||
'refresh_state',
|
||||
).select();
|
||||
expect(unprocessed).toEqual(
|
||||
expect.objectContaining({
|
||||
entity_ref: 'component:default/1',
|
||||
location_key: step.expectedLocationKey,
|
||||
}),
|
||||
);
|
||||
const entity = JSON.parse(unprocessed.unprocessed_entity) as Entity;
|
||||
expect(entity.spec?.type).toEqual(step.expectedTestKey);
|
||||
|
||||
await expect(
|
||||
knexTx<DbRefreshStateReferencesRow>(
|
||||
'refresh_state_references',
|
||||
).select(),
|
||||
).resolves.toEqual([
|
||||
expect.objectContaining({
|
||||
source_entity_ref: 'location:default/fakelocation',
|
||||
target_entity_ref: 'component:default/1',
|
||||
}),
|
||||
]);
|
||||
|
||||
expect(mockLogger.error).not.toHaveBeenCalled();
|
||||
}
|
||||
});
|
||||
},
|
||||
60_000,
|
||||
);
|
||||
});
|
||||
|
||||
describe('updateEntityCache', () => {
|
||||
|
||||
@@ -31,7 +31,6 @@ import {
|
||||
DbRelationsRow,
|
||||
} from './tables';
|
||||
import {
|
||||
AddUnprocessedEntitiesOptions,
|
||||
GetProcessableEntitiesResult,
|
||||
ProcessingDatabase,
|
||||
RefreshStateItem,
|
||||
@@ -151,63 +150,6 @@ export class DefaultProcessingDatabase implements ProcessingDatabase {
|
||||
.where('entity_id', id);
|
||||
}
|
||||
|
||||
private deduplicateRelations(rows: DbRelationsRow[]): DbRelationsRow[] {
|
||||
return lodash.uniqBy(
|
||||
rows,
|
||||
r => `${r.source_entity_ref}:${r.target_entity_ref}:${r.type}`,
|
||||
);
|
||||
}
|
||||
|
||||
private async createDelta(
|
||||
tx: Knex.Transaction,
|
||||
options: ReplaceUnprocessedEntitiesOptions,
|
||||
): Promise<{ toAdd: DeferredEntity[]; toRemove: string[] }> {
|
||||
if (options.type === 'delta') {
|
||||
return {
|
||||
toAdd: options.added,
|
||||
toRemove: options.removed.map(e => stringifyEntityRef(e.entity)),
|
||||
};
|
||||
}
|
||||
|
||||
// Grab all of the existing references from the same source, and their locationKeys as well
|
||||
const oldRefs = await tx<DbRefreshStateReferencesRow>(
|
||||
'refresh_state_references',
|
||||
)
|
||||
.where({ source_key: options.sourceKey })
|
||||
.leftJoin<DbRefreshStateRow>('refresh_state', {
|
||||
target_entity_ref: 'entity_ref',
|
||||
})
|
||||
.select(['target_entity_ref', 'location_key']);
|
||||
|
||||
const items = options.items.map(deferred => ({
|
||||
deferred,
|
||||
ref: stringifyEntityRef(deferred.entity),
|
||||
}));
|
||||
|
||||
const oldRefsSet = new Map(
|
||||
oldRefs.map(r => [r.target_entity_ref, r.location_key]),
|
||||
);
|
||||
const newRefsSet = new Set(items.map(item => item.ref));
|
||||
|
||||
const toAdd = new Array<DeferredEntity>();
|
||||
const toRemove = oldRefs
|
||||
.map(row => row.target_entity_ref)
|
||||
.filter(ref => !newRefsSet.has(ref));
|
||||
|
||||
for (const item of items) {
|
||||
if (!oldRefsSet.has(item.ref)) {
|
||||
// Add any entity that does not exist in the database
|
||||
toAdd.push(item.deferred);
|
||||
} else if (oldRefsSet.get(item.ref) !== item.deferred.locationKey) {
|
||||
// Remove and then re-add any entity that exists, but with a different location key
|
||||
toRemove.push(item.ref);
|
||||
toAdd.push(item.deferred);
|
||||
}
|
||||
}
|
||||
|
||||
return { toAdd, toRemove };
|
||||
}
|
||||
|
||||
async replaceUnprocessedEntities(
|
||||
txOpaque: Transaction,
|
||||
options: ReplaceUnprocessedEntitiesOptions,
|
||||
@@ -332,143 +274,38 @@ export class DefaultProcessingDatabase implements ProcessingDatabase {
|
||||
}
|
||||
|
||||
if (toAdd.length) {
|
||||
await this.addUnprocessedEntities(tx, {
|
||||
sourceKey: options.sourceKey,
|
||||
entities: toAdd,
|
||||
});
|
||||
}
|
||||
}
|
||||
for (const { entity, locationKey } of toAdd) {
|
||||
const entityRef = stringifyEntityRef(entity);
|
||||
|
||||
/**
|
||||
* Add a set of deferred entities for processing.
|
||||
* The entities will be added at the front of the processing queue.
|
||||
*/
|
||||
async addUnprocessedEntities(
|
||||
txOpaque: Transaction,
|
||||
options: AddUnprocessedEntitiesOptions,
|
||||
): Promise<void> {
|
||||
const tx = txOpaque as Knex.Transaction;
|
||||
|
||||
// Keeps track of the entities that we end up inserting to update refresh_state_references afterwards
|
||||
const stateReferences = new Array<string>();
|
||||
const conflictingStateReferences = new Array<string>();
|
||||
|
||||
// Upsert all of the unprocessed entities into the refresh_state table, by
|
||||
// their entity ref.
|
||||
for (const { entity, locationKey } of options.entities) {
|
||||
const entityRef = stringifyEntityRef(entity);
|
||||
const serializedEntity = JSON.stringify(entity);
|
||||
|
||||
// We optimistically try to update any existing refresh state first, as this is by far
|
||||
// the most common case.
|
||||
const refreshResult = await tx<DbRefreshStateRow>('refresh_state')
|
||||
.update({
|
||||
unprocessed_entity: serializedEntity,
|
||||
location_key: locationKey,
|
||||
last_discovery_at: tx.fn.now(),
|
||||
// We only get to this point if a processed entity actually had any changes, or
|
||||
// if an entity provider requested this mutation, meaning that we can safely
|
||||
// bump the deferred entities to the front of the queue for immediate processing.
|
||||
next_update_at: tx.fn.now(),
|
||||
})
|
||||
.where('entity_ref', entityRef)
|
||||
.andWhere(inner => {
|
||||
if (!locationKey) {
|
||||
return inner.whereNull('location_key');
|
||||
}
|
||||
return inner
|
||||
.where('location_key', locationKey)
|
||||
.orWhereNull('location_key');
|
||||
});
|
||||
|
||||
if (refreshResult === 0) {
|
||||
// In the event that we can't update an existing refresh state, we first try to insert a new row
|
||||
try {
|
||||
let query = tx('refresh_state').insert<any>({
|
||||
entity_id: uuid(),
|
||||
entity_ref: entityRef,
|
||||
unprocessed_entity: serializedEntity,
|
||||
errors: '',
|
||||
location_key: locationKey,
|
||||
next_update_at: tx.fn.now(),
|
||||
last_discovery_at: tx.fn.now(),
|
||||
});
|
||||
|
||||
// TODO(Rugvip): only tested towards Postgres and SQLite
|
||||
// We have to do this because the only way to detect if there was a conflict with
|
||||
// SQLite is to catch the error, while Postgres needs to ignore the conflict to not
|
||||
// break the ongoing transaction.
|
||||
if (tx.client.config.client !== 'sqlite3') {
|
||||
query = query.onConflict('entity_ref').ignore();
|
||||
let ok = await this.insertUnprocessedEntity(tx, entity, locationKey);
|
||||
if (!ok) {
|
||||
ok = await this.updateUnprocessedEntity(tx, entity, locationKey);
|
||||
}
|
||||
|
||||
const result: { /* postgres */ rowCount?: number } = await query;
|
||||
if (result.rowCount === 0) {
|
||||
throw new ConflictError(
|
||||
'Insert failed due to conflicting entity_ref',
|
||||
if (ok) {
|
||||
await tx('refresh_state_references').insert<any>({
|
||||
source_key: options.sourceKey,
|
||||
target_entity_ref: entityRef,
|
||||
});
|
||||
} else {
|
||||
const conflictingKey = await this.checkLocationKeyConflict(
|
||||
tx,
|
||||
entityRef,
|
||||
locationKey,
|
||||
);
|
||||
if (conflictingKey) {
|
||||
this.options.logger.warn(
|
||||
`Source ${options.sourceKey} detected conflicting entityRef ${entityRef} already referenced by ${conflictingKey} and now also ${locationKey}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
if (
|
||||
!error.message.includes('UNIQUE constraint failed') &&
|
||||
error.name !== 'ConflictError'
|
||||
) {
|
||||
throw error;
|
||||
}
|
||||
// If the row can't be inserted, we have a conflict, but it could be either
|
||||
// because of a conflicting locationKey or a race with another instance, so check
|
||||
// whether the conflicting entity has the same entityRef but a different locationKey
|
||||
const [conflictingEntity] = await tx<DbRefreshStateRow>(
|
||||
'refresh_state',
|
||||
)
|
||||
.where({ entity_ref: entityRef })
|
||||
.select();
|
||||
|
||||
// If the location key matches it means we just had a race trigger, which we can safely ignore
|
||||
if (
|
||||
!conflictingEntity ||
|
||||
conflictingEntity.location_key !== locationKey
|
||||
) {
|
||||
this.options.logger.warn(
|
||||
`Detected conflicting entityRef ${entityRef} already referenced by ${conflictingEntity.location_key} and now also ${locationKey}`,
|
||||
);
|
||||
conflictingStateReferences.push(entityRef);
|
||||
continue;
|
||||
}
|
||||
this.options.logger.error(
|
||||
`Failed to add '${entityRef}' from source '${options.sourceKey}', ${error}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Skipped on locationKey conflict
|
||||
stateReferences.push(entityRef);
|
||||
}
|
||||
|
||||
// Replace all references for the originating entity or source and then create new ones
|
||||
if ('sourceKey' in options) {
|
||||
await tx<DbRefreshStateReferencesRow>('refresh_state_references')
|
||||
.whereNotIn('target_entity_ref', conflictingStateReferences)
|
||||
.andWhere({ source_key: options.sourceKey })
|
||||
.delete();
|
||||
await tx.batchInsert(
|
||||
'refresh_state_references',
|
||||
stateReferences.map(entityRef => ({
|
||||
source_key: options.sourceKey,
|
||||
target_entity_ref: entityRef,
|
||||
})),
|
||||
BATCH_SIZE,
|
||||
);
|
||||
} else {
|
||||
await tx<DbRefreshStateReferencesRow>('refresh_state_references')
|
||||
.whereNotIn('target_entity_ref', conflictingStateReferences)
|
||||
.andWhere({ source_entity_ref: options.sourceEntityRef })
|
||||
.delete();
|
||||
await tx.batchInsert(
|
||||
'refresh_state_references',
|
||||
stateReferences.map(entityRef => ({
|
||||
source_entity_ref: options.sourceEntityRef,
|
||||
target_entity_ref: entityRef,
|
||||
})),
|
||||
BATCH_SIZE,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -599,4 +436,246 @@ export class DefaultProcessingDatabase implements ProcessingDatabase {
|
||||
throw rethrowError(e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Attempts to update an existing refresh state row, returning true if it was
|
||||
* updated and false if there was no with a matching entity ref and location key.
|
||||
*
|
||||
* Updating the entity will also caused it to be scheduled for immediate processing.
|
||||
*/
|
||||
private async updateUnprocessedEntity(
|
||||
tx: Knex.Transaction,
|
||||
entity: Entity,
|
||||
locationKey?: string,
|
||||
): Promise<boolean> {
|
||||
const entityRef = stringifyEntityRef(entity);
|
||||
const serializedEntity = JSON.stringify(entity);
|
||||
|
||||
// We optimistically try to update any existing refresh state first, as this is by far
|
||||
// the most common case.
|
||||
const refreshResult = await tx<DbRefreshStateRow>('refresh_state')
|
||||
.update({
|
||||
unprocessed_entity: serializedEntity,
|
||||
location_key: locationKey,
|
||||
last_discovery_at: tx.fn.now(),
|
||||
// We only get to this point if a processed entity actually had any changes, or
|
||||
// if an entity provider requested this mutation, meaning that we can safely
|
||||
// bump the deferred entities to the front of the queue for immediate processing.
|
||||
next_update_at: tx.fn.now(),
|
||||
})
|
||||
.where('entity_ref', entityRef)
|
||||
.andWhere(inner => {
|
||||
if (!locationKey) {
|
||||
return inner.whereNull('location_key');
|
||||
}
|
||||
return inner
|
||||
.where('location_key', locationKey)
|
||||
.orWhereNull('location_key');
|
||||
});
|
||||
|
||||
return refreshResult === 1;
|
||||
}
|
||||
|
||||
/**
|
||||
* Attempts to insert a new refresh state row for the given entity, returning
|
||||
* true if successful and false if there was a conflict.
|
||||
*/
|
||||
private async insertUnprocessedEntity(
|
||||
tx: Knex.Transaction,
|
||||
entity: Entity,
|
||||
locationKey?: string,
|
||||
): Promise<boolean> {
|
||||
const entityRef = stringifyEntityRef(entity);
|
||||
const serializedEntity = JSON.stringify(entity);
|
||||
|
||||
// In the event that we can't update an existing refresh state, we first try to insert a new row
|
||||
try {
|
||||
let query = tx('refresh_state').insert<any>({
|
||||
entity_id: uuid(),
|
||||
entity_ref: entityRef,
|
||||
unprocessed_entity: serializedEntity,
|
||||
errors: '',
|
||||
location_key: locationKey,
|
||||
next_update_at: tx.fn.now(),
|
||||
last_discovery_at: tx.fn.now(),
|
||||
});
|
||||
|
||||
// TODO(Rugvip): only tested towards Postgres and SQLite
|
||||
// We have to do this because the only way to detect if there was a conflict with
|
||||
// SQLite is to catch the error, while Postgres needs to ignore the conflict to not
|
||||
// break the ongoing transaction.
|
||||
if (tx.client.config.client !== 'sqlite3') {
|
||||
query = query.onConflict('entity_ref').ignore();
|
||||
}
|
||||
|
||||
const result: { /* postgres */ rowCount?: number; length?: number } =
|
||||
await query;
|
||||
return result.rowCount === 1 || result.length === 1;
|
||||
} catch (error) {
|
||||
// SQLite reached this rather than the rowCount check above
|
||||
if (error.message.includes('UNIQUE constraint failed')) {
|
||||
return false;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks whether a refresh state exists for the given entity that has a
|
||||
* location key that does not match the provided location key.
|
||||
*
|
||||
* @returns The conflicting key if there is one.
|
||||
*/
|
||||
private async checkLocationKeyConflict(
|
||||
tx: Knex.Transaction,
|
||||
entityRef: string,
|
||||
locationKey?: string,
|
||||
): Promise<string | undefined> {
|
||||
const row = await tx<DbRefreshStateRow>('refresh_state')
|
||||
.select('location_key')
|
||||
.where('entity_ref', entityRef)
|
||||
.first();
|
||||
|
||||
const conflictingKey = row?.location_key;
|
||||
|
||||
// If there's no existing key we can't have a conflict
|
||||
if (!conflictingKey) {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
if (conflictingKey !== locationKey) {
|
||||
return conflictingKey;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
private deduplicateRelations(rows: DbRelationsRow[]): DbRelationsRow[] {
|
||||
return lodash.uniqBy(
|
||||
rows,
|
||||
r => `${r.source_entity_ref}:${r.target_entity_ref}:${r.type}`,
|
||||
);
|
||||
}
|
||||
|
||||
private async createDelta(
|
||||
tx: Knex.Transaction,
|
||||
options: ReplaceUnprocessedEntitiesOptions,
|
||||
): Promise<{ toAdd: DeferredEntity[]; toRemove: string[] }> {
|
||||
if (options.type === 'delta') {
|
||||
return {
|
||||
toAdd: options.added,
|
||||
toRemove: options.removed.map(e => stringifyEntityRef(e.entity)),
|
||||
};
|
||||
}
|
||||
|
||||
// Grab all of the existing references from the same source, and their locationKeys as well
|
||||
const oldRefs = await tx<DbRefreshStateReferencesRow>(
|
||||
'refresh_state_references',
|
||||
)
|
||||
.where({ source_key: options.sourceKey })
|
||||
.leftJoin<DbRefreshStateRow>('refresh_state', {
|
||||
target_entity_ref: 'entity_ref',
|
||||
})
|
||||
.select(['target_entity_ref', 'location_key']);
|
||||
|
||||
const items = options.items.map(deferred => ({
|
||||
deferred,
|
||||
ref: stringifyEntityRef(deferred.entity),
|
||||
}));
|
||||
|
||||
const oldRefsSet = new Map(
|
||||
oldRefs.map(r => [r.target_entity_ref, r.location_key]),
|
||||
);
|
||||
const newRefsSet = new Set(items.map(item => item.ref));
|
||||
|
||||
const toAdd = new Array<DeferredEntity>();
|
||||
const toRemove = oldRefs
|
||||
.map(row => row.target_entity_ref)
|
||||
.filter(ref => !newRefsSet.has(ref));
|
||||
|
||||
for (const item of items) {
|
||||
if (!oldRefsSet.has(item.ref)) {
|
||||
// Add any entity that does not exist in the database
|
||||
toAdd.push(item.deferred);
|
||||
} else if (oldRefsSet.get(item.ref) !== item.deferred.locationKey) {
|
||||
// Remove and then re-add any entity that exists, but with a different location key
|
||||
toRemove.push(item.ref);
|
||||
toAdd.push(item.deferred);
|
||||
}
|
||||
}
|
||||
|
||||
return { toAdd, toRemove };
|
||||
}
|
||||
|
||||
/**
|
||||
* Add a set of deferred entities for processing.
|
||||
* The entities will be added at the front of the processing queue.
|
||||
*/
|
||||
private async addUnprocessedEntities(
|
||||
txOpaque: Transaction,
|
||||
options: {
|
||||
sourceEntityRef: string;
|
||||
entities: DeferredEntity[];
|
||||
},
|
||||
): Promise<void> {
|
||||
const tx = txOpaque as Knex.Transaction;
|
||||
|
||||
// Keeps track of the entities that we end up inserting to update refresh_state_references afterwards
|
||||
const stateReferences = new Array<string>();
|
||||
const conflictingStateReferences = new Array<string>();
|
||||
|
||||
// Upsert all of the unprocessed entities into the refresh_state table, by
|
||||
// their entity ref.
|
||||
for (const { entity, locationKey } of options.entities) {
|
||||
const entityRef = stringifyEntityRef(entity);
|
||||
|
||||
const updated = await this.updateUnprocessedEntity(
|
||||
tx,
|
||||
entity,
|
||||
locationKey,
|
||||
);
|
||||
if (updated) {
|
||||
stateReferences.push(entityRef);
|
||||
continue;
|
||||
}
|
||||
|
||||
const inserted = await this.insertUnprocessedEntity(
|
||||
tx,
|
||||
entity,
|
||||
locationKey,
|
||||
);
|
||||
if (inserted) {
|
||||
stateReferences.push(entityRef);
|
||||
continue;
|
||||
}
|
||||
|
||||
// If the row can't be inserted, we have a conflict, but it could be either
|
||||
// because of a conflicting locationKey or a race with another instance, so check
|
||||
// whether the conflicting entity has the same entityRef but a different locationKey
|
||||
const conflictingKey = await this.checkLocationKeyConflict(
|
||||
tx,
|
||||
entityRef,
|
||||
locationKey,
|
||||
);
|
||||
if (conflictingKey) {
|
||||
this.options.logger.warn(
|
||||
`Detected conflicting entityRef ${entityRef} already referenced by ${conflictingKey} and now also ${locationKey}`,
|
||||
);
|
||||
conflictingStateReferences.push(entityRef);
|
||||
}
|
||||
}
|
||||
|
||||
// Replace all references for the originating entity or source and then create new ones
|
||||
await tx<DbRefreshStateReferencesRow>('refresh_state_references')
|
||||
.whereNotIn('target_entity_ref', conflictingStateReferences)
|
||||
.andWhere({ source_entity_ref: options.sourceEntityRef })
|
||||
.delete();
|
||||
await tx.batchInsert(
|
||||
'refresh_state_references',
|
||||
stateReferences.map(entityRef => ({
|
||||
source_entity_ref: options.sourceEntityRef,
|
||||
target_entity_ref: entityRef,
|
||||
})),
|
||||
BATCH_SIZE,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -20,16 +20,6 @@ import { DateTime } from 'luxon';
|
||||
import { Transaction } from '../../database/types';
|
||||
import { DeferredEntity } from '../processing/types';
|
||||
|
||||
export type AddUnprocessedEntitiesOptions =
|
||||
| {
|
||||
sourceEntityRef: string;
|
||||
entities: DeferredEntity[];
|
||||
}
|
||||
| {
|
||||
sourceKey: string;
|
||||
entities: DeferredEntity[];
|
||||
};
|
||||
|
||||
export type AddUnprocessedEntitiesResult = {};
|
||||
|
||||
export type UpdateProcessedEntityOptions = {
|
||||
|
||||
Reference in New Issue
Block a user