feat: make queues durable by default and add seed helpers
This commit is contained in:
@@ -6,11 +6,11 @@ Part of the **WrNexus** framework — an SSR-first, Bun-native full-stack web fr
|
||||
|
||||
## Overview
|
||||
|
||||
`@wrnexus/queue` is a server-side in-process job queue. You register named
|
||||
`@wrnexus/queue` is a server-side durable job queue. You register named
|
||||
workers, enqueue jobs (optionally delayed or recurring), and let the queue poll
|
||||
and run them on a timer — with per-job retry limits and doubling backoff between
|
||||
attempts. The default store lives in memory; the design allows a pluggable driver
|
||||
to back it with Redis/SQL for durability across restarts. Reach for it when you
|
||||
attempts. `defineQueue` persists to `.wrnexus/queue.sqlite` by default, while
|
||||
low-level `createQueue` remains an intentionally in-memory primitive. Reach for it when you
|
||||
need to defer work (emails, webhooks, cleanup) off the request path without a
|
||||
heavyweight external broker. Tests can drive it deterministically via `drain()`.
|
||||
|
||||
@@ -100,6 +100,20 @@ interface Job<T = unknown> {
|
||||
|
||||
## Usage
|
||||
|
||||
Application queues belong in `app/queues/*.ts` and default-export `defineQueue(...)`.
|
||||
Configure persistence once in `wrnexus.config.ts`:
|
||||
|
||||
```ts
|
||||
export default {
|
||||
queue: { storage: "database", databaseName: "default" },
|
||||
// Or omit queue entirely for .wrnexus/queue.sqlite.
|
||||
};
|
||||
```
|
||||
|
||||
Use `storage: "memory"` only for disposable work. A selected database must exist;
|
||||
WrNexus fails startup/first use with a clear error instead of silently falling back.
|
||||
Each job may declare `success(data, job)` and `failed(data, error, job)` lifecycle hooks.
|
||||
|
||||
Register workers, enqueue jobs, then start the poller:
|
||||
|
||||
```ts
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@wrnexus/queue",
|
||||
"version": "0.8.11",
|
||||
"version": "0.8.12",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "src/index.ts",
|
||||
@@ -10,6 +10,7 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@wrnexus/core": "workspace:*",
|
||||
"@wrnexus/db": "workspace:*",
|
||||
"@wrnexus/rpc": "workspace:*"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
import { getDb, hasDb, registerDb, type Db } from "@wrnexus/db";
|
||||
import { connectFromConfig } from "@wrnexus/db/connect";
|
||||
import { memoryQueueStore, type QueueStore } from "./durable.ts";
|
||||
import { installSqliteQueueSchema, sqliteQueueStore } from "./sqlite.ts";
|
||||
|
||||
export type QueueStorage = "sqlite" | "database" | "memory";
|
||||
|
||||
export interface QueueStorageConfig {
|
||||
/** Durable SQLite is the default. Use memory only for disposable/test queues. */
|
||||
storage?: QueueStorage;
|
||||
/** Named `databases` connection, or `default` for the app's `db`. */
|
||||
databaseName?: string;
|
||||
table?: string;
|
||||
/** Used internally to resolve the default `.wrnexus/queue.sqlite` path. */
|
||||
appRoot?: string;
|
||||
}
|
||||
|
||||
let configured: QueueStorageConfig = { storage: "sqlite" };
|
||||
|
||||
export function configureQueueStorage(config: QueueStorageConfig = {}): void {
|
||||
configured = { storage: "sqlite", ...config };
|
||||
}
|
||||
|
||||
export function queueStorageConfig(): Readonly<QueueStorageConfig> {
|
||||
return { ...configured };
|
||||
}
|
||||
|
||||
function lazyStore(resolve: () => Promise<QueueStore>): QueueStore {
|
||||
let pending: Promise<QueueStore> | undefined;
|
||||
const ready = () => (pending ??= resolve());
|
||||
return {
|
||||
async put(job) { return (await ready()).put(job); },
|
||||
async get(id) { return (await ready()).get(id); },
|
||||
async remove(id) { return (await ready()).remove(id); },
|
||||
async due(now, limit) { return (await ready()).due(now, limit); },
|
||||
async list(name) { return (await ready()).list(name); },
|
||||
async findByIdempotencyKey(name, key) {
|
||||
return (await ready()).findByIdempotencyKey?.(name, key) ?? null;
|
||||
},
|
||||
async size() { return (await ready()).size?.() ?? (await (await ready()).list()).length; },
|
||||
async claim(id, worker, leaseUntil, now) {
|
||||
return (await ready()).claim?.(id, worker, leaseUntil, now) ?? true;
|
||||
},
|
||||
async release(id, worker) { await (await ready()).release?.(id, worker); },
|
||||
async archive(record) { await (await ready()).archive?.(record); },
|
||||
async history(id) { return (await ready()).history?.(id) ?? []; },
|
||||
};
|
||||
}
|
||||
|
||||
async function resolveConfiguredStore(): Promise<QueueStore> {
|
||||
const config = configured;
|
||||
if (config.storage === "memory") return memoryQueueStore();
|
||||
|
||||
const databaseName = config.databaseName ?? "default";
|
||||
let db: Db;
|
||||
if (config.storage === "database") {
|
||||
if (!hasDb(databaseName)) {
|
||||
throw new Error(
|
||||
`WRN-QUEUE-DATABASE: database '${databaseName}' is not configured. ` +
|
||||
"Add it to db/databases or choose storage: 'sqlite'.",
|
||||
);
|
||||
}
|
||||
db = getDb(databaseName);
|
||||
} else {
|
||||
if (hasDb("__wrnexus_queue")) db = getDb("__wrnexus_queue");
|
||||
else {
|
||||
db = connectFromConfig(
|
||||
{ driver: "sqlite", url: "file:./.wrnexus/queue.sqlite" },
|
||||
config.appRoot ?? process.cwd(),
|
||||
);
|
||||
registerDb("__wrnexus_queue", db);
|
||||
}
|
||||
}
|
||||
|
||||
if (db.driver.dialect !== "sqlite") {
|
||||
throw new Error(
|
||||
`WRN-QUEUE-DATABASE: configured queue storage currently requires SQLite; ` +
|
||||
`database '${databaseName}' uses ${db.driver.dialect}. Pass a custom queue store for that driver.`,
|
||||
);
|
||||
}
|
||||
const table = config.table ?? "wrnexus_jobs";
|
||||
await installSqliteQueueSchema(db, table);
|
||||
return sqliteQueueStore(db, table);
|
||||
}
|
||||
|
||||
/** Lazy global store used by defineQueue; runtime config is read on first operation. */
|
||||
export function configuredQueueStore(): QueueStore {
|
||||
return lazyStore(resolveConfiguredStore);
|
||||
}
|
||||
@@ -5,11 +5,13 @@ import {
|
||||
type DurableQueueOptions,
|
||||
type QueueHealth,
|
||||
} from "./durable.ts";
|
||||
import { configuredQueueStore } from "./configured.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>;
|
||||
success?: (data: T, 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;
|
||||
@@ -65,6 +67,7 @@ export function defineQueue<TJobs extends QueueJobDefinitions>(
|
||||
definition.queue ??
|
||||
createDurableQueue({
|
||||
...definition.options,
|
||||
store: definition.options?.store ?? configuredQueueStore(),
|
||||
async onDeadLetter(job, error) {
|
||||
await definition.options?.onDeadLetter?.(job, error);
|
||||
const prefix = `${name}:`;
|
||||
@@ -82,6 +85,7 @@ export function defineQueue<TJobs extends QueueJobDefinitions>(
|
||||
throw new TypeError(`WRN-QUEUE-PAYLOAD: invalid payload for '${jobName}'`);
|
||||
}
|
||||
await jobDefinition.run(job.data, { ...context, job });
|
||||
await jobDefinition.success?.(job.data, job);
|
||||
};
|
||||
queue.process(jobName, handler);
|
||||
|
||||
|
||||
@@ -317,6 +317,8 @@ export type {
|
||||
} from "./defined.ts";
|
||||
export { installSqliteQueueSchema, sqliteQueueSchema, sqliteQueueStore } from "./sqlite.ts";
|
||||
export type { SqliteQueueClient } from "./sqlite.ts";
|
||||
export { configureQueueStorage, configuredQueueStore, queueStorageConfig } from "./configured.ts";
|
||||
export type { QueueStorage, QueueStorageConfig } from "./configured.ts";
|
||||
export { redisQueueStore, postgresQueueStore, POSTGRES_QUEUE_SCHEMA } from "./stores.ts";
|
||||
export type { RedisQueueClient, SqlQueueClient } from "./stores.ts";
|
||||
export {
|
||||
|
||||
@@ -3,6 +3,7 @@ import { createDurableQueue, defineQueue, memoryQueueStore } from "../src/index.
|
||||
|
||||
test("defineQueue creates typed producers and registers workers", async () => {
|
||||
const seen: number[] = [];
|
||||
const completed: number[] = [];
|
||||
const email = defineQueue({
|
||||
name: "email",
|
||||
options: { store: memoryQueueStore() },
|
||||
@@ -15,6 +16,7 @@ test("defineQueue creates typed producers and registers workers", async () => {
|
||||
data !== null &&
|
||||
Number.isInteger((data as { messageId?: unknown }).messageId),
|
||||
run: async ({ messageId }: { messageId: number }) => void seen.push(messageId),
|
||||
success: async ({ messageId }: { messageId: number }) => void completed.push(messageId),
|
||||
},
|
||||
},
|
||||
});
|
||||
@@ -25,6 +27,7 @@ test("defineQueue creates typed producers and registers workers", async () => {
|
||||
expect(await email.send.status(first.id)).toBe("queued");
|
||||
expect(await email.drain()).toBe(1);
|
||||
expect(seen).toEqual([7]);
|
||||
expect(completed).toEqual([7]);
|
||||
expect(await email.send.status(first.id)).toBe("completed");
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user