diff --git a/apps/api/src/server.ts b/apps/api/src/server.ts index 6b319f7..e8d7065 100644 --- a/apps/api/src/server.ts +++ b/apps/api/src/server.ts @@ -321,11 +321,16 @@ app.listen(3000, async () => { for (const queue of queues) { (async () => { while (true) { - // clean, if ready - if ((Date.now() - queue.lastCleanDate.getTime()) >= queue.cfg.concurrency.timeoutMs) - await queue.clean().catch((err) => log.error(err)); - // drain - await queue.drain().catch((err) => log.error(err)); + try { + // clean, if ready + if ((Date.now() - queue.lastCleanDate.getTime()) >= queue.cfg.concurrency.timeoutMs) + await queue.clean().catch((err) => log.error({ name: queue.type, err })); + // drain + await queue.drain(); + } + catch (err) { + log.error({ name: queue.type, err }, "error occurred during queue management"); + } await sleep(EVENT_QUEUE_MANAGE_DELAY_MS); } })(); diff --git a/apps/api/src/services/event-queue.ts b/apps/api/src/services/event-queue.ts index f45759f..bffcf6c 100644 --- a/apps/api/src/services/event-queue.ts +++ b/apps/api/src/services/event-queue.ts @@ -38,9 +38,11 @@ type EventQueueServiceOptions { private _lastCleanDate: Date = new Date(0); readonly cfg: EventQueueServiceOptions["cfg"]; + readonly type: EventType constructor(private readonly opt: EventQueueServiceOptions) { this.cfg = opt.cfg; + this.type = opt.type; } get lastCleanDate() { return this._lastCleanDate; };