Small cleanup
This commit is contained in:
@@ -44,14 +44,6 @@ function exponentialDelay(initialDelay: number, tries: number, cap: number = Inf
|
|||||||
return Math.min(initialDelay * (2 ** tries), cap);
|
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<G extends string, T extends Record<string, EventTypeOptions<any, any>>> {
|
export class EventQueue<G extends string, T extends Record<string, EventTypeOptions<any, any>>> {
|
||||||
private static readonly DEQUEUE_RETRY_DELAY_MS = secToMs(10);
|
private static readonly DEQUEUE_RETRY_DELAY_MS = secToMs(10);
|
||||||
private static readonly DEQUEUE_RETRY_LIMIT = 3;
|
private static readonly DEQUEUE_RETRY_LIMIT = 3;
|
||||||
@@ -140,7 +132,7 @@ export class EventQueue<G extends string, T extends Record<string, EventTypeOpti
|
|||||||
};
|
};
|
||||||
|
|
||||||
log.trace({ name: front.id }, "dequeuing event");
|
log.trace({ name: front.id }, "dequeuing event");
|
||||||
await this.processEvent(front);
|
await this.process(front);
|
||||||
|
|
||||||
return front;
|
return front;
|
||||||
}
|
}
|
||||||
@@ -179,18 +171,18 @@ export class EventQueue<G extends string, T extends Record<string, EventTypeOpti
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private async processEvent(event: Pick<Selectable<DB["events"]>, "id" | "type" | "errors" | "payload" | "rate_limits">) {
|
private async process(event: Pick<Selectable<DB["events"]>, "id" | "type" | "errors" | "payload" | "rate_limits">) {
|
||||||
log.info({ id: event.id, type: event.type, payload: event.payload }, "processing event");
|
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}`);
|
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]!;
|
const cfg = this.types[event.type as keyof T]!;
|
||||||
|
|
||||||
try {
|
try {
|
||||||
Value.Assert(cfg.schemas.payload, event.payload);
|
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);
|
await this.Events.fulfill(event.id);
|
||||||
}
|
}
|
||||||
catch (err) {
|
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) {
|
if (err instanceof RateLimitError && event.rate_limits < EventQueue.RATE_LIMIT_RETRY_LIMIT) {
|
||||||
log.warn(
|
log.warn(
|
||||||
@@ -222,20 +214,6 @@ export class EventQueue<G extends string, T extends Record<string, EventTypeOpti
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private async setEventState(id: string, state: JsonValue) {
|
|
||||||
await db.updateTable("events")
|
|
||||||
.set({ state })
|
|
||||||
.where("id", "=", id)
|
|
||||||
.execute();
|
|
||||||
}
|
|
||||||
|
|
||||||
private async setLastError(id: string, err: unknown) {
|
|
||||||
await db.updateTable("events")
|
|
||||||
.set({ last_error: serializeError(err) })
|
|
||||||
.where("id", "=", id)
|
|
||||||
.execute();
|
|
||||||
}
|
|
||||||
|
|
||||||
private async retry(id: string, delayMs: number, type: RetryType) {
|
private async retry(id: string, delayMs: number, type: RetryType) {
|
||||||
await db.updateTable("events")
|
await db.updateTable("events")
|
||||||
.set((eb) => ({
|
.set((eb) => ({
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
import type { Static } from "@sinclair/typebox";
|
import type { Static } from "@sinclair/typebox";
|
||||||
import { t } from "elysia";
|
import { t } from "elysia";
|
||||||
import { sql, type Selectable, type Transaction } from "kysely";
|
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 { db } from "../database/db";
|
||||||
import { log } from "../util";
|
import { log } from "../util";
|
||||||
import { minToMs } from "@blade-and-brawn/domain";
|
import { minToMs } from "@blade-and-brawn/domain";
|
||||||
@@ -21,9 +21,25 @@ export class EventsService {
|
|||||||
static readonly CONCURRENCY_TIMEOUT_MS = minToMs(10);
|
static readonly CONCURRENCY_TIMEOUT_MS = minToMs(10);
|
||||||
|
|
||||||
private _lastCleanDate: Date = new Date(0);
|
private _lastCleanDate: Date = new Date(0);
|
||||||
|
|
||||||
get lastCleanDate() { return this._lastCleanDate; };
|
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<Selectable<DB["events"]>, "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<DB>) {
|
async get(id: string, trx?: Transaction<DB>) {
|
||||||
return await (trx ?? db).selectFrom("events")
|
return await (trx ?? db).selectFrom("events")
|
||||||
.selectAll()
|
.selectAll()
|
||||||
@@ -73,13 +89,18 @@ export class EventsService {
|
|||||||
.execute();
|
.execute();
|
||||||
}
|
}
|
||||||
|
|
||||||
static status(event: Pick<Selectable<DB["events"]>, "available_at" | "started_at" | "ended_at" | "has_failed">): EventStatus {
|
async setState(id: string, state: JsonValue) {
|
||||||
if (!event.started_at) {
|
await db.updateTable("events")
|
||||||
if (Date.now() < event.available_at.getTime()) return "waiting";
|
.set({ state })
|
||||||
return "available";
|
.where("id", "=", id)
|
||||||
};
|
.execute();
|
||||||
if (!event.ended_at) return "processing";
|
}
|
||||||
return event.has_failed ? "failed" : "fulfilled";
|
|
||||||
|
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<DB>): Promise<void> {
|
async fail(id: string, trx?: Transaction<DB>): Promise<void> {
|
||||||
|
|||||||
Reference in New Issue
Block a user