Improve 'tries' naming scheme and implement exponential retry delay with jitter
This commit is contained in:
@@ -5,11 +5,13 @@ import { sleep } from "bun";
|
||||
import { db } from "../database/db";
|
||||
import { AssertError, Value } from "@sinclair/typebox/value";
|
||||
import type { TSchema, Static } from "@sinclair/typebox";
|
||||
import { RateLimitError } from "@blade-and-brawn/domain";
|
||||
import { minToMs, RateLimitError, secToMs } from "@blade-and-brawn/domain";
|
||||
|
||||
export type EventSource = "portal" | "printful" | "webflow";
|
||||
export type EventStatus = "waiting" | "available" | "processing" | "failed" | "fulfilled";
|
||||
|
||||
type RetryType = "error" | "ratelimit";
|
||||
|
||||
type EventProcessor<P extends TSchema, S extends TSchema> =
|
||||
(payload: Static<P>, setState: (state: Static<S>) => Promise<void>, id: string) => Promise<void>;
|
||||
|
||||
@@ -19,9 +21,9 @@ export type EventTypeOptions<P extends TSchema, S extends TSchema> = {
|
||||
};
|
||||
|
||||
type EventQueueCfg = {
|
||||
retry: {
|
||||
error: {
|
||||
limit: number,
|
||||
delayMs: number,
|
||||
initialRetryDelayMs: number,
|
||||
},
|
||||
concurrency: {
|
||||
limit: number,
|
||||
@@ -36,10 +38,20 @@ export function defineEventType<P extends TSchema, S extends TSchema>(opt: Event
|
||||
return opt;
|
||||
}
|
||||
|
||||
function addJitter(delay: number, jitterFactor: number): number {
|
||||
const jitterMultiplier = 1 + (Math.random() - 0.5) * jitterFactor;
|
||||
return delay * jitterMultiplier;
|
||||
}
|
||||
|
||||
function exponentialDelay(initialDelay: number, tries: number, cap: number = Infinity) {
|
||||
return Math.min(initialDelay * (2 ** tries), cap);
|
||||
}
|
||||
|
||||
export class EventQueue<G extends string, T extends Record<string, EventTypeOptions<any, any>>> {
|
||||
private static readonly DEQUEUE_RETRY_DELAY_MS = 10000; // 10 seconds
|
||||
private static readonly DEQUEUE_RETRY_DELAY_MS = secToMs(10);
|
||||
private static readonly DEQUEUE_RETRY_LIMIT = 3;
|
||||
private static readonly RATE_LIMIT_RETRY_LIMIT = 15;
|
||||
private static readonly JITTER_FACTOR = 0.2;
|
||||
|
||||
private _lastCleanDate: Date = new Date(0);
|
||||
readonly cfg: EventQueueCfg;
|
||||
@@ -181,7 +193,7 @@ export class EventQueue<G extends string, T extends Record<string, EventTypeOpti
|
||||
|
||||
async front(trx?: Transaction<DB>) {
|
||||
const result = await (trx ?? db).selectFrom("events")
|
||||
.select(["id", "type", "payload", "available_at", "tries", "rate_limits"])
|
||||
.select(["id", "type", "payload", "available_at", "errors", "rate_limits"])
|
||||
.where("group", "=", this.group)
|
||||
.where("started_at", "is", null)
|
||||
.where("available_at", "<=", new Date())
|
||||
@@ -222,7 +234,7 @@ export class EventQueue<G extends string, T extends Record<string, EventTypeOpti
|
||||
.execute();
|
||||
}
|
||||
|
||||
private async processEvent(event: Pick<Selectable<DB["events"]>, "id" | "type" | "tries" | "payload" | "rate_limits">) {
|
||||
private async processEvent(event: Pick<Selectable<DB["events"]>, "id" | "type" | "errors" | "payload" | "rate_limits">) {
|
||||
log.info({ id: event.id, type: event.type, payload: event.payload }, "processing event");
|
||||
if (!(event.type in this.events)) throw new Error(`event type not registered in this queue (${this.group}): ${event.type}`);
|
||||
const cfg = this.events[event.type as keyof T]!;
|
||||
@@ -235,26 +247,30 @@ export class EventQueue<G extends string, T extends Record<string, EventTypeOpti
|
||||
catch (err) {
|
||||
if (err instanceof RateLimitError && event.rate_limits < EventQueue.RATE_LIMIT_RETRY_LIMIT) {
|
||||
log.warn(
|
||||
{ tries: event.tries, delay: err.retryDelayMs, source: err.source, payload: err.payload },
|
||||
{ rateLimits: event.rate_limits, delay: err.retryDelayMs, source: err.source, payload: err.payload },
|
||||
"encountered a rate limit error"
|
||||
);
|
||||
await this.retry(event.id, err.retryDelayMs, { isRateLimit: true });
|
||||
await this.retry(event.id, err.retryDelayMs, "ratelimit");
|
||||
return;
|
||||
}
|
||||
|
||||
if (err instanceof AssertError) {
|
||||
log.error({ tries: event.tries }, "could not process event: malformed payload");
|
||||
log.error({ errors: event.errors }, "could not process event: malformed payload");
|
||||
await this.fail(event.id);
|
||||
return;
|
||||
}
|
||||
|
||||
if ((event.tries - event.rate_limits) < this.cfg.retry.limit) {
|
||||
log.warn({ err, tries: event.tries, delay: this.cfg.retry.delayMs }, "encountered an error during event processing, retrying")
|
||||
await this.retry(event.id, this.cfg.retry.delayMs);
|
||||
if (event.errors < this.cfg.error.limit) {
|
||||
const retryDelayMs = addJitter(
|
||||
exponentialDelay(this.cfg.error.initialRetryDelayMs, event.errors),
|
||||
EventQueue.JITTER_FACTOR
|
||||
);
|
||||
log.warn({ err, errors: event.errors, retryDelayMs: retryDelayMs }, "encountered an error during event processing, retrying")
|
||||
await this.retry(event.id, retryDelayMs, "error");
|
||||
return;
|
||||
}
|
||||
|
||||
log.error({ err, tries: event.tries }, "event failed");
|
||||
log.error({ err, errors: event.errors, rateLimits: event.rate_limits }, "event failed");
|
||||
await this.fail(event.id);
|
||||
}
|
||||
}
|
||||
@@ -266,11 +282,11 @@ export class EventQueue<G extends string, T extends Record<string, EventTypeOpti
|
||||
.execute();
|
||||
}
|
||||
|
||||
private async retry(id: string, delayMs: number, opt: { isRateLimit?: boolean } = {}) {
|
||||
private async retry(id: string, delayMs: number, type: RetryType) {
|
||||
await db.updateTable("events")
|
||||
.set((eb) => ({
|
||||
tries: eb("tries", "+", 1),
|
||||
rate_limits: opt.isRateLimit ? eb("rate_limits", "+", 1) : eb.ref("rate_limits"),
|
||||
errors: type === "error" ? eb("errors", "+", 1) : eb.ref("errors"),
|
||||
rate_limits: type === "ratelimit" ? eb("rate_limits", "+", 1) : eb.ref("rate_limits"),
|
||||
started_at: null,
|
||||
ended_at: null,
|
||||
has_failed: false,
|
||||
|
||||
Reference in New Issue
Block a user