feat(queue): add typed application queue lifecycle
Quality / quality (ubuntu-latest) (push) Failing after 10m55s
Quality / quality (windows-latest) (push) Canceled after 0s

This commit is contained in:
2026-08-23 09:22:05 +05:30
parent 9332522614
commit 9377a69d88
17 changed files with 667 additions and 13 deletions
+39
View File
@@ -118,6 +118,45 @@ await queue.add("email", { to: "a@b.com" }, { delayMs: 5000, maxAttempts: 3 });
queue.start(); // begin polling; queue.stop() to halt
```
### Application queues with `defineQueue`
WrNexus applications can place typed queue modules in `app/queues`. The dev and
production servers discover them, start them after runtime initialization, and
shut them down gracefully:
```ts
import { defineQueue } from "@wrnexus/queue";
export default defineQueue({
name: "email",
jobs: {
send: {
maxAttempts: 3,
idempotency: ({ messageId }: { messageId: number }) => `message:${messageId}`,
run: async ({ messageId }) => sendMessage(messageId),
},
},
});
```
Callers get a typed producer and inspection helpers:
```ts
const job = await email.send.add({ messageId: 42 });
await email.send.status(job.id);
await email.send.cancel(job.id);
await email.send.retry(job.id);
```
`createDurableQueue` also exposes `start()`, `stop()`, `isRunning()`, `health()`
and `get()` so applications do not need to build their own polling middleware.
### SQLite
Use `sqliteQueueStore(db)` from `@wrnexus/queue/sqlite` with a WrNexus-compatible
SQLite client. `installSqliteQueueSchema(db)` installs the queue table and its
due-job and idempotency indexes.
Use `context.signal` in network/database calls so forced shutdown and active
cancellation finish promptly. For process termination, prefer
`await queue.shutdown()`; use `{ force: true }` only after your grace period.
+3 -2
View File
@@ -1,11 +1,12 @@
{
"name": "@wrnexus/queue",
"version": "0.8.10",
"version": "0.8.11",
"private": true,
"type": "module",
"main": "src/index.ts",
"exports": {
".": "./src/index.ts"
".": "./src/index.ts",
"./sqlite": "./src/sqlite.ts"
},
"dependencies": {
"@wrnexus/core": "workspace:*",
+139
View File
@@ -0,0 +1,139 @@
import type { AddOptions, Job, JobContext, JobHandler } from "./index.ts";
import {
createDurableQueue,
type DurableQueue,
type DurableQueueOptions,
type QueueHealth,
} from "./durable.ts";
export type JobStatus = "queued" | "completed" | "failed" | "cancelled" | "missing";
export interface DefinedJobOptions<T> extends Omit<AddOptions, "idempotencyKey"> {
run: (data: T, context: JobContext & { job: Job<T> }) => void | Promise<void>;
failed?: (data: T, error: unknown, job: Job<T>) => void | Promise<void>;
idempotency?: (data: T) => string | undefined;
validate?: (data: unknown) => data is T;
}
export type QueueJobDefinitions = Record<string, DefinedJobOptions<any>>;
export interface DefineQueueOptions<TJobs extends QueueJobDefinitions> {
name: string;
jobs: TJobs;
queue?: DurableQueue;
options?: DurableQueueOptions;
}
type DataOf<T> = T extends DefinedJobOptions<infer I> ? I : never;
export interface DefinedJob<T> {
readonly name: string;
add(data: T, options?: AddOptions): Promise<Job<T>>;
addMany(entries: readonly { data: T; options?: AddOptions }[]): Promise<Job<T>[]>;
get(id: string): Promise<Job<T> | null>;
status(id: string): Promise<JobStatus>;
list(): Promise<Job<T>[]>;
cancel(id: string): Promise<boolean>;
retry(id: string): Promise<boolean>;
}
export type DefinedJobs<TJobs extends QueueJobDefinitions> = {
[K in keyof TJobs]: DefinedJob<DataOf<TJobs[K]>>;
};
export interface DefinedQueue<TJobs extends QueueJobDefinitions> {
readonly name: string;
readonly jobs: DefinedJobs<TJobs>;
readonly raw: DurableQueue;
start(): void;
stop(): void;
shutdown(options?: { force?: boolean }): Promise<void>;
health(): Promise<QueueHealth>;
drain(): Promise<number>;
}
function fullName(queue: string, job: string): string {
return `${queue}:${job}`;
}
export function defineQueue<TJobs extends QueueJobDefinitions>(
definition: DefineQueueOptions<TJobs>,
): DefinedQueue<TJobs> & DefinedJobs<TJobs> {
const name = definition.name.trim();
if (!name) throw new TypeError("queue name cannot be empty");
const queue =
definition.queue ??
createDurableQueue({
...definition.options,
async onDeadLetter(job, error) {
await definition.options?.onDeadLetter?.(job, error);
const prefix = `${name}:`;
if (!job.name.startsWith(prefix)) return;
const failed = definition.jobs[job.name.slice(prefix.length)]?.failed;
await failed?.(job.data, error, job);
},
});
const jobs: Record<string, DefinedJob<unknown>> = {};
for (const [shortName, jobDefinition] of Object.entries(definition.jobs)) {
const jobName = fullName(name, shortName);
const handler: JobHandler<unknown> = async (job, context) => {
if (jobDefinition.validate && !jobDefinition.validate(job.data)) {
throw new TypeError(`WRN-QUEUE-PAYLOAD: invalid payload for '${jobName}'`);
}
await jobDefinition.run(job.data, { ...context, job });
};
queue.process(jobName, handler);
jobs[shortName] = {
name: jobName,
add(data, options = {}) {
if (jobDefinition.validate && !jobDefinition.validate(data)) {
return Promise.reject(
new TypeError(`WRN-QUEUE-PAYLOAD: invalid payload for '${jobName}'`),
);
}
return queue.add(jobName, data, {
maxAttempts: jobDefinition.maxAttempts,
delayMs: jobDefinition.delayMs,
repeat: jobDefinition.repeat,
priority: jobDefinition.priority,
idempotencyKey: jobDefinition.idempotency?.(data),
...options,
});
},
addMany(entries) {
return Promise.all(entries.map((entry) => this.add(entry.data, entry.options)));
},
async get(id) {
const job = await queue.get(id);
return job?.name === jobName ? (job as Job<unknown>) : null;
},
async status(id) {
return queue.status(id);
},
async list() {
return (await queue.list(jobName)) as Job<unknown>[];
},
cancel(id) {
return queue.cancel(id);
},
retry(id) {
return queue.retry(id);
},
};
}
const result = {
name,
jobs: jobs as DefinedJobs<TJobs>,
raw: queue,
start: () => queue.start(),
stop: () => queue.stop(),
shutdown: (options?: { force?: boolean }) => queue.shutdown(options),
health: () => queue.health(),
drain: () => queue.drain(),
...(jobs as DefinedJobs<TJobs>),
};
return result;
}
+160 -2
View File
@@ -17,10 +17,54 @@ export interface QueueStore {
findByIdempotencyKey?(name: string, key: string): Promise<Job | null>;
size?(): Promise<number>;
claim?(id: string, worker: string, leaseUntil: number, now?: number): Promise<boolean>;
release?(id: string, worker: string): Promise<void>;
archive?(record: QueueJobRecord): Promise<void>;
history?(id?: string): Promise<QueueJobRecord[]>;
}
export type QueueJobState = "completed" | "failed" | "cancelled";
export interface QueueJobRecord {
job: Job;
state: QueueJobState;
finishedAt: number;
error?: string;
}
export type QueueEventType =
| "job.added"
| "job.started"
| "job.completed"
| "job.retrying"
| "job.failed"
| "job.cancelled"
| "worker.started"
| "worker.stopped"
| "worker.error";
export interface QueueEvent {
type: QueueEventType;
at: number;
workerId: string;
job?: Job;
error?: unknown;
durationMs?: number;
}
export interface QueueHealth {
running: boolean;
accepting: boolean;
draining: boolean;
active: number;
pending: number;
failed: number;
workerId: string;
lastPollAt?: number;
lastError?: unknown;
}
export function memoryQueueStore(): QueueStore {
const jobs = new Map<string, Job>();
const records = new Map<string, QueueJobRecord>();
return {
async put(job) {
@@ -59,6 +103,14 @@ export function memoryQueueStore(): QueueStore {
async size() {
return jobs.size;
},
async archive(record) {
records.set(record.job.id, structuredClone(record));
},
async history(id) {
return [...records.values()]
.filter((record) => !id || record.job.id === id)
.map((record) => structuredClone(record));
},
};
}
@@ -69,16 +121,27 @@ export interface DurableQueueOptions {
concurrency?: number;
capacity?: number;
leaseMs?: number;
pollMs?: number;
backoff?: (attempt: number) => number;
now?: () => number;
onDeadLetter?: (job: Job, error: unknown) => void | Promise<void>;
context?: (job: Job, signal: AbortSignal) => ExecutionContext;
onEvent?: (event: QueueEvent) => void | Promise<void>;
/** Maintenance hook run before each poll, for application-specific recovery. */
beforeDrain?: (queue: DurableQueue) => void | Promise<void>;
}
export interface DurableQueue {
add<T>(name: string, data: T, options?: AddOptions): Promise<Job<T>>;
process<T>(name: string, handler: JobHandler<T>): void;
drain(): Promise<number>;
start(): void;
stop(): void;
isRunning(): boolean;
health(): Promise<QueueHealth>;
get(id: string): Promise<Job | null>;
status(id: string): Promise<"queued" | QueueJobState | "missing">;
history(id?: string): Promise<QueueJobRecord[]>;
cancel(id: string): Promise<boolean>;
shutdown(options?: { force?: boolean }): Promise<void>;
list(name?: string): Promise<Job[]>;
@@ -114,11 +177,19 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
const concurrency = positiveInteger(options.concurrency ?? 10, "queue concurrency");
const capacity = positiveInteger(options.capacity ?? 10_000, "queue capacity");
const leaseMs = positiveInteger(options.leaseMs ?? 30_000, "queue leaseMs");
const pollMs = positiveInteger(options.pollMs ?? 250, "queue pollMs");
let sequence = 0;
let draining = false;
let accepting = true;
let timer: ReturnType<typeof setInterval> | null = null;
let lastPollAt: number | undefined;
let lastError: unknown;
const active = new Map<string, { controller: AbortController; promise: Promise<boolean> }>();
const emit = async (event: Omit<QueueEvent, "at" | "workerId">) => {
await options.onEvent?.({ ...event, at: now(), workerId });
};
function beginJob(job: Job): Promise<boolean> {
const controller = new AbortController();
const promise = runJob(job, controller).finally(() => active.delete(job.id));
@@ -135,6 +206,8 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
}
job.attempts += 1;
const startedAt = now();
void emit({ type: "job.started", job: structuredClone(job) });
try {
await handler(job, {
@@ -153,8 +226,14 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
job.runAt = now() + job.repeat;
await store.put(job);
} else {
await store.archive?.({ job: structuredClone(job), state: "completed", finishedAt: now() });
await store.remove(job.id);
}
await emit({
type: "job.completed",
job: structuredClone(job),
durationMs: now() - startedAt,
});
} catch (error) {
if (controller.signal.aborted) {
await store.remove(job.id);
@@ -167,17 +246,25 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
);
job.runAt = now() + delay;
await store.put(job);
await emit({ type: "job.retrying", job: structuredClone(job), error });
} else {
await store.remove(job.id);
deadLetters.set(job.id, structuredClone(job));
await store.archive?.({
job: structuredClone(job),
state: "failed",
finishedAt: now(),
error: error instanceof Error ? error.message : String(error),
});
await options.onDeadLetter?.(structuredClone(job), error);
await emit({ type: "job.failed", job: structuredClone(job), error });
}
}
return true;
}
return {
const durableQueue: DurableQueue = {
async add<T>(name: string, data: T, add: AddOptions = {}) {
if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down");
if (!name.trim()) throw new TypeError("queue job name cannot be empty");
@@ -216,6 +303,7 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
createdAt,
};
await store.put(job);
await emit({ type: "job.added", job: structuredClone(job) });
return structuredClone(job);
},
@@ -233,7 +321,9 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
async drain() {
if (draining) return 0;
draining = true;
lastPollAt = now();
try {
if (options.beforeDrain) await options.beforeDrain(durableQueue);
const due = await store.due(now(), concurrency);
const results = await Promise.all(due.map(beginJob));
return results.filter(Boolean).length;
@@ -250,19 +340,80 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
return store.list(name);
},
get(id) {
return store.get(id);
},
async status(id) {
if (await store.get(id)) return "queued";
const record = (await store.history?.(id))?.at(-1);
if (record) return record.state;
if (deadLetters.has(id)) return "failed";
return "missing";
},
async history(id) {
return (await store.history?.(id)) ?? [];
},
start() {
if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down");
if (timer) return;
timer = setInterval(() => {
void this.drain().catch(async (error: unknown) => {
lastError = error;
await emit({ type: "worker.error", error });
});
}, pollMs);
timer.unref?.();
void emit({ type: "worker.started" });
},
stop() {
if (timer) clearInterval(timer);
timer = null;
void emit({ type: "worker.stopped" });
},
isRunning() {
return timer !== null;
},
async health() {
return {
running: timer !== null,
accepting,
draining,
active: active.size,
pending: store.size ? await store.size() : (await store.list()).length,
failed: deadLetters.size,
workerId,
lastPollAt,
lastError,
};
},
async cancel(id) {
const running = active.get(id);
if (running) {
running.controller.abort();
const job = await store.get(id);
if (job) await store.archive?.({ job, state: "cancelled", finishedAt: now() });
await emit({ type: "job.cancelled" });
return true;
}
if (!(await store.get(id))) return false;
const job = await store.get(id);
if (job) await store.archive?.({ job, state: "cancelled", finishedAt: now() });
await store.remove(id);
await emit({ type: "job.cancelled" });
return true;
},
async shutdown(shutdownOptions = {}) {
accepting = false;
if (timer) clearInterval(timer);
timer = null;
if (shutdownOptions.force) {
for (const { controller } of active.values()) controller.abort();
}
@@ -270,7 +421,13 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
},
async retry(id) {
const job = deadLetters.get(id);
let job = deadLetters.get(id);
if (!job) {
const record = (await store.history?.(id))?.find(
(candidate) => candidate.state === "failed",
);
job = record?.job;
}
if (!job) return false;
deadLetters.delete(id);
job.attempts = 0;
@@ -279,4 +436,5 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
return true;
},
};
return durableQueue;
}
+19
View File
@@ -298,6 +298,25 @@ export function createQueue(options: QueueOptions = {}): Queue {
}
export { memoryQueueStore, createDurableQueue } from "./durable.ts";
export type { QueueStore, DurableQueue, DurableQueueOptions } from "./durable.ts";
export type {
QueueEvent,
QueueEventType,
QueueHealth,
QueueJobRecord,
QueueJobState,
} from "./durable.ts";
export { defineQueue } from "./defined.ts";
export type {
DefinedJob,
DefinedJobOptions,
DefinedJobs,
DefinedQueue,
DefineQueueOptions,
JobStatus,
QueueJobDefinitions,
} from "./defined.ts";
export { installSqliteQueueSchema, sqliteQueueSchema, sqliteQueueStore } from "./sqlite.ts";
export type { SqliteQueueClient } from "./sqlite.ts";
export { redisQueueStore, postgresQueueStore, POSTGRES_QUEUE_SCHEMA } from "./stores.ts";
export type { RedisQueueClient, SqlQueueClient } from "./stores.ts";
export {
+2 -1
View File
@@ -85,10 +85,11 @@ export interface QueueDashboardSnapshot {
export async function queueDashboardSnapshot(queue: DurableQueue): Promise<QueueDashboardSnapshot> {
const jobs = await queue.list();
const history = await queue.history();
return {
generatedAt: Date.now(),
pending: jobs.length,
failed: queue.failed().length,
failed: history.filter((record) => record.state === "failed").length || queue.failed().length,
byName: jobs.reduce<Record<string, number>>((counts, job) => {
counts[job.name] = (counts[job.name] ?? 0) + 1;
return counts;
+167
View File
@@ -0,0 +1,167 @@
import type { Job } from "./index.ts";
import type { QueueJobRecord, QueueStore } from "./durable.ts";
export interface SqliteQueueClient {
exec(sql: string, parameters?: unknown[]): Promise<{ changes: number }>;
one<T>(sql: string, parameters?: unknown[]): Promise<T | null>;
all<T>(sql: string, parameters?: unknown[]): Promise<T[]>;
}
interface QueueJobRow {
id: string;
name: string;
payload: string;
run_at: number;
priority: number;
idempotency_key: string | null;
}
const SAFE_TABLE = /^[a-z_][a-z0-9_]*$/i;
export function sqliteQueueSchema(table = "wrnexus_jobs"): string {
if (!SAFE_TABLE.test(table)) throw new Error("Invalid queue table name");
return `CREATE TABLE IF NOT EXISTS ${table} (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
payload TEXT NOT NULL,
run_at INTEGER NOT NULL,
priority INTEGER NOT NULL DEFAULT 0,
idempotency_key TEXT,
lease_owner TEXT,
lease_until INTEGER
);
CREATE INDEX IF NOT EXISTS ${table}_due ON ${table} (run_at, priority DESC);
CREATE UNIQUE INDEX IF NOT EXISTS ${table}_idempotency ON ${table} (name, idempotency_key)
WHERE idempotency_key IS NOT NULL;
CREATE TABLE IF NOT EXISTS ${table}_history (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
state TEXT NOT NULL,
payload TEXT NOT NULL,
error TEXT,
finished_at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS ${table}_history_finished ON ${table}_history (finished_at DESC);`;
}
export async function installSqliteQueueSchema(
db: SqliteQueueClient,
table = "wrnexus_jobs",
): Promise<void> {
for (const statement of sqliteQueueSchema(table)
.split(";")
.map((value) => value.trim())) {
if (statement) await db.exec(statement);
}
}
export function sqliteQueueStore(db: SqliteQueueClient, table = "wrnexus_jobs"): QueueStore {
if (!SAFE_TABLE.test(table)) throw new Error("Invalid queue table name");
const decode = (row: QueueJobRow): Job => JSON.parse(row.payload) as Job;
return {
async put(job) {
await db.exec(
`INSERT INTO ${table} (id,name,payload,run_at,priority,idempotency_key,lease_owner,lease_until)
VALUES (?,?,?,?,?,?,NULL,NULL)
ON CONFLICT(id) DO UPDATE SET name=excluded.name,payload=excluded.payload,
run_at=excluded.run_at,priority=excluded.priority,idempotency_key=excluded.idempotency_key,
lease_owner=NULL,lease_until=NULL`,
[
job.id,
job.name,
JSON.stringify(job),
job.runAt,
job.priority,
job.idempotencyKey ?? null,
],
);
},
async get(id) {
const row = await db.one<QueueJobRow>(`SELECT * FROM ${table} WHERE id = ?`, [id]);
return row ? decode(row) : null;
},
async remove(id) {
await db.exec(`DELETE FROM ${table} WHERE id = ?`, [id]);
},
async due(now, limit) {
const rows = await db.all<QueueJobRow>(
`SELECT * FROM ${table} WHERE run_at <= ? AND (lease_until IS NULL OR lease_until < ?)
ORDER BY priority DESC,run_at ASC LIMIT ?`,
[now, now, limit],
);
return rows.map(decode);
},
async list(name) {
const rows = name
? await db.all<QueueJobRow>(`SELECT * FROM ${table} WHERE name = ? ORDER BY run_at`, [name])
: await db.all<QueueJobRow>(`SELECT * FROM ${table} ORDER BY run_at`);
return rows.map(decode);
},
async findByIdempotencyKey(name, key) {
const row = await db.one<QueueJobRow>(
`SELECT * FROM ${table} WHERE name = ? AND idempotency_key = ? LIMIT 1`,
[name, key],
);
return row ? decode(row) : null;
},
async size() {
const row = await db.one<{ total: number }>(`SELECT COUNT(*) AS total FROM ${table}`);
return Number(row?.total ?? 0);
},
async claim(id, worker, leaseUntil, now = Date.now()) {
const result = await db.exec(
`UPDATE ${table} SET lease_owner = ?,lease_until = ?
WHERE id = ? AND (lease_until IS NULL OR lease_until < ?)`,
[worker, leaseUntil, id, now],
);
return result.changes > 0;
},
async release(id, worker) {
await db.exec(
`UPDATE ${table} SET lease_owner = NULL,lease_until = NULL WHERE id = ? AND lease_owner = ?`,
[id, worker],
);
},
async archive(record) {
await db.exec(
`INSERT INTO ${table}_history (id,name,state,payload,error,finished_at) VALUES (?,?,?,?,?,?)
ON CONFLICT(id) DO UPDATE SET state=excluded.state,payload=excluded.payload,
error=excluded.error,finished_at=excluded.finished_at`,
[
record.job.id,
record.job.name,
record.state,
JSON.stringify(record.job),
record.error ?? null,
record.finishedAt,
],
);
},
async history(id) {
const rows = id
? await db.all<{
state: QueueJobRecord["state"];
payload: string;
error: string | null;
finished_at: number;
}>(
`SELECT state,payload,error,finished_at FROM ${table}_history WHERE id = ? ORDER BY finished_at DESC`,
[id],
)
: await db.all<{
state: QueueJobRecord["state"];
payload: string;
error: string | null;
finished_at: number;
}>(
`SELECT state,payload,error,finished_at FROM ${table}_history ORDER BY finished_at DESC`,
);
return rows.map((row) => ({
job: JSON.parse(row.payload) as Job,
state: row.state,
finishedAt: row.finished_at,
error: row.error ?? undefined,
}));
},
};
}
+62
View File
@@ -0,0 +1,62 @@
import { expect, test } from "bun:test";
import { createDurableQueue, defineQueue, memoryQueueStore } from "../src/index.ts";
test("defineQueue creates typed producers and registers workers", async () => {
const seen: number[] = [];
const email = defineQueue({
name: "email",
options: { store: memoryQueueStore() },
jobs: {
send: {
maxAttempts: 4,
idempotency: (data: { messageId: number }) => `message:${data.messageId}`,
validate: (data: unknown): data is { messageId: number } =>
typeof data === "object" &&
data !== null &&
Number.isInteger((data as { messageId?: unknown }).messageId),
run: async ({ messageId }: { messageId: number }) => void seen.push(messageId),
},
},
});
const first = await email.send.add({ messageId: 7 });
const duplicate = await email.jobs.send.add({ messageId: 7 });
expect(duplicate.id).toBe(first.id);
expect(await email.send.status(first.id)).toBe("queued");
expect(await email.drain()).toBe(1);
expect(seen).toEqual([7]);
expect(await email.send.status(first.id)).toBe("completed");
});
test("durable queues expose lifecycle, health and events", async () => {
const events: string[] = [];
const queue = createDurableQueue({
pollMs: 10,
onEvent: (event) => void events.push(event.type),
});
queue.process("work", async () => {});
queue.start();
expect(queue.isRunning()).toBe(true);
await queue.add("work", {});
await queue.drain();
const health = await queue.health();
expect(health.pending).toBe(0);
expect(events).toContain("job.added");
expect(events).toContain("job.completed");
queue.stop();
expect(queue.isRunning()).toBe(false);
await queue.shutdown();
});
test("defineQueue validates payloads before persistence", async () => {
const queue = defineQueue({
name: "safe",
jobs: {
number: {
validate: (data: unknown): data is number => typeof data === "number",
run: async (_data: number) => {},
},
},
});
await expect(queue.number.add("no" as never)).rejects.toThrow("WRN-QUEUE-PAYLOAD");
});