From c412d48980ed3ed6780fe3c056611f6ef6825647 Mon Sep 17 00:00:00 2001 From: Patrik Oldsberg Date: Thu, 12 Sep 2024 23:17:17 +0200 Subject: [PATCH] events-backend: add cleanup test for DatabaseEventBusStore Signed-off-by: Patrik Oldsberg --- .../service/hub/DatabaseEventBusStore.test.ts | 62 +++++++++++++++++++ .../src/service/hub/DatabaseEventBusStore.ts | 16 +++++ 2 files changed, 78 insertions(+) create mode 100644 plugins/events-backend/src/service/hub/DatabaseEventBusStore.test.ts diff --git a/plugins/events-backend/src/service/hub/DatabaseEventBusStore.test.ts b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.test.ts new file mode 100644 index 0000000000..ed9fa3a7af --- /dev/null +++ b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.test.ts @@ -0,0 +1,62 @@ +/* + * Copyright 2024 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { TestDatabases, mockServices } from '@backstage/backend-test-utils'; +import { DatabaseEventBusStore } from './DatabaseEventBusStore'; + +const logger = mockServices.logger.mock(); + +const databases = TestDatabases.create({ + ids: ['POSTGRES_9', 'POSTGRES_13', 'POSTGRES_16'], +}); + +describe('DatabaseEventBusStore', () => { + it.each(databases.eachSupportedId())( + 'should clean up old events, %p', + async databaseId => { + const db = await databases.init(databaseId); + const store = await DatabaseEventBusStore.forTest({ logger, db }); + + await store.upsertSubscription('tester-1', ['test']); + await store.upsertSubscription('tester-2', ['test']); + + for (let i = 0; i < 10; ++i) { + await store.publish({ + params: { topic: 'test', eventPayload: { n: i } }, + }); + } + + const { events: events1 } = await store.readSubscription('tester-1'); + expect(events1.length).toBe(10); + + await store.clean(); + + await expect(store.readSubscription('tester-2')).rejects.toThrow( + "Subscription with ID 'tester-2' not found", + ); + + await store.upsertSubscription('tester-3', ['test']); + + // Reset read pointer to read form the beginning + await db('event_bus_subscriptions').select({ id: 'tester-3' }).update({ + read_until: 0, + }); + + const { events: events3 } = await store.readSubscription('tester-3'); + expect(events3.length).toBe(5); + }, + ); +}); diff --git a/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts index e8f0da444d..744e9d3569 100644 --- a/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts +++ b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts @@ -298,6 +298,22 @@ export class DatabaseEventBusStore implements EventBusStore { return store; } + /** @internal */ + static async forTest({ db, logger }: { db: Knex; logger: LoggerService }) { + await db.migrate.latest({ directory: migrationsDir }); + + const store = new DatabaseEventBusStore( + db, + logger, + new DatabaseEventBusListener(db.client, logger), + 5, + 0, + 10, + ); + + return Object.assign(store, { clean: () => store.#cleanup() }); + } + readonly #db: Knex; readonly #logger: LoggerService; readonly #listener: DatabaseEventBusListener;