From 79f73d9f92510077a09f402f849b7c758f455f9e Mon Sep 17 00:00:00 2001 From: Patrik Oldsberg Date: Sat, 25 May 2024 09:38:23 +0200 Subject: [PATCH] events-backend: added cleanup of old events Signed-off-by: Patrik Oldsberg --- .../src/service/EventsPlugin.ts | 15 ++- .../src/service/hub/DatabaseEventBusStore.ts | 92 ++++++++++++++++++- .../src/service/hub/createEventBusRouter.ts | 4 + 3 files changed, 107 insertions(+), 4 deletions(-) diff --git a/plugins/events-backend/src/service/EventsPlugin.ts b/plugins/events-backend/src/service/EventsPlugin.ts index 176d7e2f34..e8769b4a09 100644 --- a/plugins/events-backend/src/service/EventsPlugin.ts +++ b/plugins/events-backend/src/service/EventsPlugin.ts @@ -77,10 +77,19 @@ export const eventsPlugin = createBackendPlugin({ events: eventsServiceRef, database: coreServices.database, logger: coreServices.logger, + scheduler: coreServices.scheduler, httpAuth: coreServices.httpAuth, router: coreServices.httpRouter, }, - async init({ config, events, database, logger, httpAuth, router }) { + async init({ + config, + events, + database, + logger, + scheduler, + httpAuth, + router, + }) { const ingresses = Object.fromEntries( extensionPoint.httpPostIngresses.map(ingress => [ ingress.topic, @@ -97,7 +106,9 @@ export const eventsPlugin = createBackendPlugin({ const eventsRouter = Router(); http.bind(eventsRouter); - router.use(await createEventBusRouter({ database, logger, httpAuth })); + router.use( + await createEventBusRouter({ database, logger, httpAuth, scheduler }), + ); router.use(eventsRouter); router.addAuthPolicy({ diff --git a/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts index e5e69d697d..a6260cd0be 100644 --- a/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts +++ b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts @@ -19,9 +19,15 @@ import { Knex } from 'knex'; import { DatabaseService, LoggerService, + SchedulerService, resolvePackagePath, } from '@backstage/backend-plugin-api'; import { ForwardedError, NotFoundError } from '@backstage/errors'; +import { HumanDuration, durationToMilliseconds } from '@backstage/types'; + +const WINDOW_SIZE_DEFAULT = 10000; +const WINDOW_MIN_AGE_DEFAULT = { minutes: 10 }; +const WINDOW_MAX_AGE_DEFAULT = { days: 1 }; const MAX_BATCH_SIZE = 10; const LISTENER_CONNECTION_TIMEOUT_MS = 60_000; @@ -208,6 +214,15 @@ export class DatabaseEventBusStore implements EventBusStore { static async create(options: { database: DatabaseService; logger: LoggerService; + scheduler: SchedulerService; + window?: { + /** Events within this range will never be deleted */ + minAge?: HumanDuration; + /** Events outside of this age will always be deleted */ + maxAge?: HumanDuration; + /** Events outside of this count will be deleted if they are outside the minAge window */ + size?: number; + }; }): Promise { const db = await options.database.getClient(); @@ -224,16 +239,43 @@ export class DatabaseEventBusStore implements EventBusStore { options.logger.info('DatabaseEventBusStore migrations ran successfully'); } - return new DatabaseEventBusStore(db, options.logger); + const store = new DatabaseEventBusStore( + db, + options.logger, + options.window?.size ?? WINDOW_SIZE_DEFAULT, + durationToMilliseconds(options.window?.minAge ?? WINDOW_MIN_AGE_DEFAULT), + durationToMilliseconds(options.window?.maxAge ?? WINDOW_MAX_AGE_DEFAULT), + ); + + options.scheduler.scheduleTask({ + id: 'event-bus-cleanup', + frequency: { seconds: 10 }, + timeout: { minutes: 1 }, + fn: store.#cleanup, + }); + + return store; } readonly #db: Knex; readonly #logger: LoggerService; + readonly #windowSize: number; + readonly #windowMinAge: number; + readonly #windowMaxAge: number; readonly #listener: DatabaseEventBusListener; - private constructor(db: Knex, logger: LoggerService) { + private constructor( + db: Knex, + logger: LoggerService, + windowSize: number, + windowMinAge: number, + windowMaxAge: number, + ) { this.#db = db; this.#logger = logger; + this.#windowSize = windowSize; + this.#windowMinAge = windowMinAge; + this.#windowMaxAge = windowMaxAge; this.#listener = new DatabaseEventBusListener(db.client, logger); } @@ -430,4 +472,50 @@ export class DatabaseEventBusStore implements EventBusStore { options.signal.addEventListener('abort', cancel); } } + + #cleanup = async () => { + try { + const eventCount = await this.#db + .delete() + .from(TABLE_EVENTS) + // Delete any events that are outside both the min age and size window + .where('created_at', '<', new Date(Date.now() - this.#windowMinAge)) + .andWhere( + 'id', + '=', + this.#db + .raw( + this.#db + .select('id') + .from(TABLE_EVENTS) + .orderBy('id', 'desc') + .offset(this.#windowSize), + ) + .wrap('ANY(ARRAY(', '))'), + ) + // If events are outside the max age they will always be deleted + .orWhere('created_at', '<', new Date(Date.now() - this.#windowMaxAge)); + this.#logger.info( + `Event cleanup resulted in ${eventCount} old events being deleted`, + ); + } catch (error) { + this.#logger.error('Event cleanup failed', error); + } + + try { + // Delete any subscribers that aren't keeping up with current events + const subscriberCount = await this.#db + .delete() + .from(TABLE_SUBSCRIPTIONS) + .where('read_until', '<', (q: Knex.QueryBuilder) => + q.select(this.#db.raw('MIN(id)')).from(TABLE_EVENTS), + ); + + this.#logger.info( + `Subscription cleanup resulted in ${subscriberCount} stale subscribers being deleted`, + ); + } catch (error) { + this.#logger.error('Subscription cleanup failed', error); + } + }; } diff --git a/plugins/events-backend/src/service/hub/createEventBusRouter.ts b/plugins/events-backend/src/service/hub/createEventBusRouter.ts index 14237650fc..30cd137104 100644 --- a/plugins/events-backend/src/service/hub/createEventBusRouter.ts +++ b/plugins/events-backend/src/service/hub/createEventBusRouter.ts @@ -18,6 +18,7 @@ import { DatabaseService, HttpAuthService, LoggerService, + SchedulerService, } from '@backstage/backend-plugin-api'; import { Handler } from 'express'; import Router from 'express-promise-router'; @@ -32,12 +33,14 @@ const DEFAULT_NOTIFY_TIMEOUT_MS = 55_000; // Just below 60s, which is a common H export async function createEventBusRouter(options: { logger: LoggerService; database: DatabaseService; + scheduler: SchedulerService; httpAuth: HttpAuthService; notifyTimeoutMs?: number; // for testing }): Promise { const { database, httpAuth, + scheduler, notifyTimeoutMs = DEFAULT_NOTIFY_TIMEOUT_MS, } = options; const logger = options.logger.child({ type: 'EventBus' }); @@ -50,6 +53,7 @@ export async function createEventBusRouter(options: { store = await DatabaseEventBusStore.create({ database, logger, + scheduler, }); } else { logger.info('Database is not PostgreSQL, using memory store');