diff --git a/apps/api/src/services/event-queue.ts b/apps/api/src/services/event-queue.ts index 6cc440f..88fe5e1 100644 --- a/apps/api/src/services/event-queue.ts +++ b/apps/api/src/services/event-queue.ts @@ -44,14 +44,6 @@ function exponentialDelay(initialDelay: number, tries: number, cap: number = Inf return Math.min(initialDelay * (2 ** tries), cap); } -function serializeError(err: unknown): JsonValue { - if (err instanceof Error) { - const payload = "payload" in err ? { payload: err.payload as JsonValue } : {}; - return { name: err.name, message: err.message, ...payload }; - } - return String(err); -} - export class EventQueue>> { private static readonly DEQUEUE_RETRY_DELAY_MS = secToMs(10); private static readonly DEQUEUE_RETRY_LIMIT = 3; @@ -140,7 +132,7 @@ export class EventQueue, "id" | "type" | "errors" | "payload" | "rate_limits">) { + private async process(event: Pick, "id" | "type" | "errors" | "payload" | "rate_limits">) { log.info({ id: event.id, type: event.type, payload: event.payload }, "processing event"); if (!(event.type in this.types)) throw new Error(`event type not registered in this queue (${this.group}): ${event.type}`); const cfg = this.types[event.type as keyof T]!; try { Value.Assert(cfg.schemas.payload, event.payload); - await cfg.processor(event.payload, (state) => this.setEventState(event.id, state), event.id); + await cfg.processor(event.payload, (state) => this.Events.setState(event.id, state), event.id); await this.Events.fulfill(event.id); } catch (err) { - await this.setLastError(event.id, err); + await this.Events.setLastError(event.id, err); if (err instanceof RateLimitError && event.rate_limits < EventQueue.RATE_LIMIT_RETRY_LIMIT) { log.warn( @@ -222,20 +214,6 @@ export class EventQueue ({ diff --git a/apps/api/src/services/events.ts b/apps/api/src/services/events.ts index ba0e4e3..7788ed7 100644 --- a/apps/api/src/services/events.ts +++ b/apps/api/src/services/events.ts @@ -1,7 +1,7 @@ import type { Static } from "@sinclair/typebox"; import { t } from "elysia"; import { sql, type Selectable, type Transaction } from "kysely"; -import type { DB } from "../database/out/db"; +import type { DB, JsonValue } from "../database/out/db"; import { db } from "../database/db"; import { log } from "../util"; import { minToMs } from "@blade-and-brawn/domain"; @@ -21,9 +21,25 @@ export class EventsService { static readonly CONCURRENCY_TIMEOUT_MS = minToMs(10); private _lastCleanDate: Date = new Date(0); - get lastCleanDate() { return this._lastCleanDate; }; + private static serializeError(err: unknown): JsonValue { + if (err instanceof Error) { + const payload = "payload" in err ? { payload: err.payload as JsonValue } : {}; + return { name: err.name, message: err.message, ...payload }; + } + return String(err); + } + + static status(event: Pick, "available_at" | "started_at" | "ended_at" | "has_failed">): EventStatus { + if (!event.started_at) { + if (Date.now() < event.available_at.getTime()) return "waiting"; + return "available"; + }; + if (!event.ended_at) return "processing"; + return event.has_failed ? "failed" : "fulfilled"; + } + async get(id: string, trx?: Transaction) { return await (trx ?? db).selectFrom("events") .selectAll() @@ -73,13 +89,18 @@ export class EventsService { .execute(); } - static status(event: Pick, "available_at" | "started_at" | "ended_at" | "has_failed">): EventStatus { - if (!event.started_at) { - if (Date.now() < event.available_at.getTime()) return "waiting"; - return "available"; - }; - if (!event.ended_at) return "processing"; - return event.has_failed ? "failed" : "fulfilled"; + async setState(id: string, state: JsonValue) { + await db.updateTable("events") + .set({ state }) + .where("id", "=", id) + .execute(); + } + + async setLastError(id: string, err: unknown) { + await db.updateTable("events") + .set({ last_error: EventsService.serializeError(err) }) + .where("id", "=", id) + .execute(); } async fail(id: string, trx?: Transaction): Promise {