feature: refactor processing engine to remove listeners publicly and provide a subscription via catalog builder

Signed-off-by: Hasan Oezdemir <21654050+nodify-at@users.noreply.github.com>
This commit is contained in:
Hasan Oezdemir
2022-07-05 16:29:41 +02:00
parent bc3497f66a
commit 751323a50b
4 changed files with 33 additions and 14 deletions
+5 -3
View File
@@ -145,6 +145,10 @@ export class CatalogBuilder {
processingInterval: ProcessingIntervalFunction,
): CatalogBuilder;
setProcessingIntervalSeconds(seconds: number): CatalogBuilder;
// @alpha (undocumented)
subscribe(
catalogProcessingErrorListeners: CatalogProcessingErrorListener[],
): void;
}
// @alpha
@@ -202,15 +206,13 @@ export type CatalogPermissionRule<TParams extends unknown[] = unknown[]> =
// @public (undocumented)
export interface CatalogProcessingEngine {
// (undocumented)
addErrorListener?(errorListener: CatalogProcessingErrorListener): void;
// (undocumented)
start(): Promise<void>;
// (undocumented)
stop(): Promise<void>;
}
// @public
// @alpha
export interface CatalogProcessingErrorListener {
// (undocumented)
onError(
@@ -34,7 +34,6 @@ const CACHE_TTL = 5;
export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine {
private readonly tracker = progressTracker();
private readonly errorListeners: CatalogProcessingErrorListener[] = [];
private stopFunc?: () => void;
constructor(
@@ -44,6 +43,7 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine {
private readonly stitcher: Stitcher,
private readonly createHash: () => Hash,
private readonly pollingIntervalMs: number = 1000,
private readonly catalogProcessingErrorListener?: CatalogProcessingErrorListener,
) {}
async start() {
@@ -157,9 +157,11 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine {
// just store the errors and trigger a stich so that they become visible to
// the outside.
if (!result.ok) {
// notify the error listeners if the entity can not be processed.
this.errorListeners.forEach(listener =>
listener.onError(unprocessedEntity, result, resultHash),
// notify the error listener if the entity can not be processed.
this.catalogProcessingErrorListener?.onError(
unprocessedEntity,
result,
resultHash,
);
await this.processingDatabase.transaction(async tx => {
@@ -231,10 +233,6 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine {
this.stopFunc = undefined;
}
}
addErrorListener(errorListener: CatalogProcessingErrorListener) {
this.errorListeners.push(errorListener);
}
}
// Helps wrap the timing and logging behaviors
@@ -66,13 +66,12 @@ export type DeferredEntity = {
export interface CatalogProcessingEngine {
start(): Promise<void>;
stop(): Promise<void>;
addErrorListener?(errorListener: CatalogProcessingErrorListener): void;
}
/**
* An error listener for catalog processing engine. It can be used to listen and track entity errors.
*
* @public
* @alpha
*/
export interface CatalogProcessingErrorListener {
onError(
@@ -56,7 +56,10 @@ import {
} from '../modules/core/PlaceholderProcessor';
import { defaultEntityDataParser } from '../modules/util/parse';
import { LocationAnalyzer } from '../ingestion/types';
import { CatalogProcessingEngine } from '../processing/types';
import {
CatalogProcessingEngine,
CatalogProcessingErrorListener,
} from '../processing';
import { DefaultProcessingDatabase } from '../database/DefaultProcessingDatabase';
import { applyDatabaseMigrations } from '../database/migrations';
import { DefaultCatalogProcessingEngine } from '../processing/DefaultCatalogProcessingEngine';
@@ -133,6 +136,7 @@ export class CatalogBuilder {
private processors: CatalogProcessor[];
private processorsReplace: boolean;
private parser: CatalogProcessorParser | undefined;
private catalogProcessingErrorListeners?: CatalogProcessingErrorListener[];
private processingInterval: ProcessingIntervalFunction =
createRandomProcessingInterval({
minSeconds: 100,
@@ -447,6 +451,14 @@ export class CatalogBuilder {
orchestrator,
stitcher,
() => createHash('sha1'),
1000,
{
onError: async (unprocessedEntity, result, resultHash) => {
this.catalogProcessingErrorListeners?.forEach(listener =>
listener.onError(unprocessedEntity, result, resultHash),
);
},
},
);
const locationAnalyzer =
@@ -478,6 +490,14 @@ export class CatalogBuilder {
};
}
/**
* @alpha
* @param catalogProcessingErrorListeners - a list of listeners to get notified if an error occurs while processing an entity
*/
subscribe(catalogProcessingErrorListeners: CatalogProcessingErrorListener[]) {
this.catalogProcessingErrorListeners = catalogProcessingErrorListeners;
}
private buildEntityPolicy(): EntityPolicy {
const entityPolicies: EntityPolicy[] = this.entityPoliciesReplace
? [new SchemaValidEntityPolicy(), ...this.entityPolicies]