diff --git a/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts index 0f4085e794..c51179896b 100644 --- a/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts +++ b/plugins/events-backend/src/service/hub/DatabaseEventBusStore.ts @@ -109,14 +109,31 @@ class DatabaseEventBusListener { await this.#ensureConnection(); const updatePromise = new Promise<{ topic: string }>((resolve, reject) => { - const listener = { topics, resolve, reject }; + const listener = { + topics, + resolve(result: { topic: string }) { + resolve(result); + cleanup(); + }, + reject(err: Error) { + reject(err); + cleanup(); + }, + }; this.#listeners.add(listener); - signal.addEventListener('abort', () => { + const onAbort = () => { this.#listeners.delete(listener); this.#maybeTimeoutConnection(); reject(signal.reason); - }); + cleanup(); + }; + + function cleanup() { + signal.removeEventListener('abort', onAbort); + } + + signal.addEventListener('abort', onAbort); }); // Ignore unhandled rejections diff --git a/plugins/events-backend/src/service/hub/MemoryEventBusStore.ts b/plugins/events-backend/src/service/hub/MemoryEventBusStore.ts index 2e4dba34bd..b7b1d78052 100644 --- a/plugins/events-backend/src/service/hub/MemoryEventBusStore.ts +++ b/plugins/events-backend/src/service/hub/MemoryEventBusStore.ts @@ -115,13 +115,26 @@ export class MemoryEventBusStore implements EventBusStore { } return new Promise<{ topic: string }>((resolve, reject) => { - const listener = { topics: sub.topics, resolve }; + const listener = { + topics: sub.topics, + resolve(result: { topic: string }) { + resolve(result); + cleanup(); + }, + }; this.#listeners.add(listener); - options.signal.addEventListener('abort', () => { + const onAbort = () => { this.#listeners.delete(listener); reject(options.signal.reason); - }); + cleanup(); + }; + + function cleanup() { + options.signal.removeEventListener('abort', onAbort); + } + + options.signal.addEventListener('abort', onAbort); }); }, };