Do not create new event route, rely on events backend instead

Signed-off-by: Damon Kaswell <damon.kaswell1@hp.com>
This commit is contained in:
Damon Kaswell
2023-01-19 08:44:30 -08:00
parent 00ad09b5b4
commit 3c1eae9324
5 changed files with 84 additions and 119 deletions
@@ -51,23 +51,22 @@ export class IncrementalCatalogBuilder {
// @public
export interface IncrementalEntityProvider<TCursor, TContext> {
around(burst: (context: TContext) => Promise<void>): Promise<void>;
eventHandler?: {
onEvent: (params: EventParams) =>
| undefined
| {
added: DeferredEntity[];
removed: {
entityRef: string;
}[];
};
supportsEventTopics: () => string[];
};
getProviderName(): string;
next(
context: TContext,
cursor?: TCursor,
): Promise<EntityIteratorResult<TCursor>>;
onEvent?: (params: EventParams) =>
| {
delta:
| {
added: DeferredEntity[];
removed: {
entityRef: string;
}[];
}
| undefined;
}
| undefined;
}
// @public (undocumented)
@@ -28,7 +28,6 @@ export class IncrementalIngestionEngine
{
private readonly restLength: Duration;
private readonly backoff: DurationObjectUnits[];
private readonly providerEventTopic: string;
private manager: IncrementalIngestionDatabaseManager;
@@ -41,7 +40,6 @@ export class IncrementalIngestionEngine
{ minutes: 30 },
{ hours: 3 },
];
this.providerEventTopic = `${options.provider.getProviderName()}-push`;
}
async taskFn(signal: AbortSignal) {
@@ -337,68 +335,68 @@ export class IncrementalIngestionEngine
async onEvent(params: EventParams): Promise<void> {
const { topic } = params;
if (topic !== this.providerEventTopic) {
if (!this.supportsEventTopics().includes(topic)) {
return;
}
const { logger, provider, connection } = this.options;
const providerName = provider.getProviderName();
logger.debug(
`incremental-engine: Received ${this.providerEventTopic} event`,
);
logger.debug(`incremental-engine: ${providerName} received ${topic} event`);
if (!provider.onEvent) {
if (!provider.eventHandler) {
return;
}
const update = provider.onEvent(params);
const delta = provider.eventHandler.onEvent(params);
if (update) {
if (update.delta) {
if (update.delta.added.length > 0) {
const ingestionRecord = await this.manager.getCurrentIngestionRecord(
providerName,
if (delta) {
if (delta.added.length > 0) {
const ingestionRecord = await this.manager.getCurrentIngestionRecord(
providerName,
);
if (!ingestionRecord) {
logger.debug(
`incremental-engine: ${providerName} skipping delta addition because incremental ingestion is restarting.`,
);
} else {
const mark =
ingestionRecord.status === 'resting'
? await this.manager.getLastMark(ingestionRecord.id)
: await this.manager.getFirstMark(ingestionRecord.id);
if (!ingestionRecord) {
logger.debug(
`incremental-engine: Skipping delta addition because incremental ingestion is restarting.`,
if (!mark) {
throw new Error(
`Cannot apply delta, page records are missing! Please re-run incremental ingestion for ${providerName}.`,
);
} else {
const mark =
ingestionRecord.status === 'resting'
? await this.manager.getLastMark(ingestionRecord.id)
: await this.manager.getFirstMark(ingestionRecord.id);
if (!mark) {
throw new Error(
`Cannot apply delta, page records are missing! Please re-run incremental ingestion for ${providerName}.`,
);
}
await this.manager.createMarkEntities(mark.id, update.delta.added);
}
await this.manager.createMarkEntities(mark.id, delta.added);
}
if (update.delta.removed.length > 0) {
await this.manager.deleteEntityRecordsByRef(update.delta.removed);
}
await connection.applyMutation({
type: 'delta',
...update.delta,
});
logger.debug(
`incremental-engine: Processed ${this.providerEventTopic} event`,
);
} else {
logger.warn(
`incremental-engine: Rejected ${this.providerEventTopic} event - empty or invalid`,
);
}
if (delta.removed.length > 0) {
await this.manager.deleteEntityRecordsByRef(delta.removed);
}
await connection.applyMutation({
type: 'delta',
...delta,
});
logger.debug(
`incremental-engine: ${providerName} processed delta from '${topic}' event`,
);
} else {
logger.warn(
`incremental-engine: Rejected delta from '${topic}' event - empty or invalid`,
);
}
}
supportsEventTopics(): string[] {
return [this.providerEventTopic];
const { provider } = this.options;
const topics = provider.eventHandler
? provider.eventHandler.supportsEventTopics()
: [];
return topics;
}
}
@@ -81,8 +81,8 @@ export class WrapperProviders {
).createRouter();
}
private async startProvider<TCursor, TContext>(
provider: IncrementalEntityProvider<TCursor, TContext>,
private async startProvider(
provider: IncrementalEntityProvider<unknown, unknown>,
providerOptions: IncrementalEntityProviderOptions,
connection: EntityProviderConnection,
) {
@@ -15,28 +15,21 @@
*/
import { errorHandler } from '@backstage/backend-common';
import { stringifyError } from '@backstage/errors';
import { EventBroker, EventPublisher } from '@backstage/plugin-events-node';
import express from 'express';
import Router from 'express-promise-router';
import { Logger } from 'winston';
import { IncrementalIngestionDatabaseManager } from '../database/IncrementalIngestionDatabaseManager';
import { PROVIDER_BASE_PATH, PROVIDER_CLEANUP, PROVIDER_HEALTH } from './paths';
export class IncrementalProviderRouter implements EventPublisher {
export class IncrementalProviderRouter {
private manager: IncrementalIngestionDatabaseManager;
private logger: Logger;
private eventBroker: EventBroker | undefined;
constructor(manager: IncrementalIngestionDatabaseManager, logger: Logger) {
this.manager = manager;
this.logger = logger;
}
async setEventBroker(eventBroker: EventBroker): Promise<void> {
this.eventBroker = eventBroker;
}
async createRouter() {
const router = Router();
router.use(express.json());
@@ -239,43 +232,6 @@ export class IncrementalProviderRouter implements EventPublisher {
});
});
router.post(`${PROVIDER_BASE_PATH}/event`, async (req, res) => {
const { provider } = req.params;
const topic = `${provider}-push`;
const eventPayload = req.body;
if (!this.eventBroker) {
res.status(500).json({
success: false,
provider,
message: `The payload could not be processed!`,
});
throw new Error('Event broker not initialized!');
}
try {
await this.eventBroker.publish({
topic,
eventPayload,
});
res.json({
success: true,
provider,
message: 'Payload submitted.',
});
} catch (e) {
res.status(500).json({
success: false,
provider,
message: `There was an error submitting the payload: ${stringifyError(
e,
)}`,
});
}
});
router.use(errorHandler());
return router;
@@ -78,23 +78,35 @@ export interface IncrementalEntityProvider<TCursor, TContext> {
around(burst: (context: TContext) => Promise<void>): Promise<void>;
/**
* This method accepts an incoming event for the provider, and
* (optionally) maps the payload to an object containing a delta
* mutation.
* If set, the IncrementalEntityProvider will receive and respond to
* events.
*
* If a valid delta is returned by this method, it will be ingested
* automatically by the provider.
* This system acts as a wrapper for the Backstage events bus, and
* requires the events backend to function. It does not provide its
* own events backend. See {@link https://github.com/backstage/backstage/tree/master/plugins/events-backend}.
*/
onEvent?: (params: EventParams) =>
| {
delta:
| {
added: DeferredEntity[];
removed: { entityRef: string }[];
}
| undefined;
}
| undefined;
eventHandler?: {
/**
* This method accepts an incoming event for the provider, and
* optionally maps the payload to an object containing a delta
* mutation.
*
* If a valid delta is returned by this method, it will be ingested
* automatically by the provider.
*/
onEvent: (params: EventParams) =>
| undefined
| {
added: DeferredEntity[];
removed: { entityRef: string }[];
};
/**
* This method returns an array of topics for the IncrementalEntityProvider
* to respond to.
*/
supportsEventTopics: () => string[];
};
}
/**