release: WRNexusJS 0.8.0
This commit is contained in:
@@ -1,3 +1,4 @@
|
||||
import { createExecutionContext, type ExecutionContext } from "@wrnexus/core";
|
||||
import type { AddOptions, Job, JobHandler } from "./index.ts";
|
||||
|
||||
/**
|
||||
@@ -55,16 +56,21 @@ export interface DurableQueueOptions {
|
||||
workerId?: string;
|
||||
maxAttempts?: number;
|
||||
concurrency?: number;
|
||||
capacity?: number;
|
||||
leaseMs?: number;
|
||||
backoff?: (attempt: number) => number;
|
||||
now?: () => number;
|
||||
onDeadLetter?: (job: Job, error: unknown) => void | Promise<void>;
|
||||
context?: (job: Job, signal: AbortSignal) => ExecutionContext;
|
||||
}
|
||||
|
||||
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>;
|
||||
cancel(id: string): Promise<boolean>;
|
||||
shutdown(options?: { force?: boolean }): Promise<void>;
|
||||
list(name?: string): Promise<Job[]>;
|
||||
failed(): Job[];
|
||||
retry(id: string): Promise<boolean>;
|
||||
}
|
||||
@@ -91,11 +97,21 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
const workerId = options.workerId?.trim() || `worker-${crypto.randomUUID()}`;
|
||||
const defaultMaxAttempts = positiveInteger(options.maxAttempts ?? 3, "queue maxAttempts");
|
||||
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");
|
||||
let sequence = 0;
|
||||
let draining = false;
|
||||
let accepting = true;
|
||||
const active = new Map<string, { controller: AbortController; promise: Promise<boolean> }>();
|
||||
|
||||
async function runJob(job: Job): Promise<boolean> {
|
||||
function beginJob(job: Job): Promise<boolean> {
|
||||
const controller = new AbortController();
|
||||
const promise = runJob(job, controller).finally(() => active.delete(job.id));
|
||||
active.set(job.id, { controller, promise });
|
||||
return promise;
|
||||
}
|
||||
|
||||
async function runJob(job: Job, controller: AbortController): Promise<boolean> {
|
||||
const handler = handlers.get(job.name);
|
||||
if (!handler) return false;
|
||||
|
||||
@@ -106,7 +122,16 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
job.attempts += 1;
|
||||
|
||||
try {
|
||||
await handler(job);
|
||||
await handler(job, {
|
||||
signal: controller.signal,
|
||||
execution:
|
||||
options.context?.(job, controller.signal) ??
|
||||
createExecutionContext({
|
||||
kind: job.repeat ? "cron" : "queue",
|
||||
signal: controller.signal,
|
||||
metadata: { jobId: job.id, jobName: job.name, attempt: job.attempts },
|
||||
}),
|
||||
});
|
||||
|
||||
if (job.repeat && job.repeat > 0) {
|
||||
job.attempts = 0;
|
||||
@@ -116,6 +141,10 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
await store.remove(job.id);
|
||||
}
|
||||
} catch (error) {
|
||||
if (controller.signal.aborted) {
|
||||
await store.remove(job.id);
|
||||
return true;
|
||||
}
|
||||
if (job.attempts < job.maxAttempts) {
|
||||
const delay = nonNegativeNumber(
|
||||
options.backoff?.(job.attempts) ?? Math.min(60_000, 1_000 * 2 ** (job.attempts - 1)),
|
||||
@@ -135,6 +164,7 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
|
||||
return {
|
||||
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");
|
||||
|
||||
const maxAttempts = positiveInteger(add.maxAttempts ?? defaultMaxAttempts, "job maxAttempts");
|
||||
@@ -153,6 +183,9 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
);
|
||||
if (existing) return existing as Job<T>;
|
||||
}
|
||||
if ((await store.list()).length >= capacity) {
|
||||
throw new Error(`WRN-QUEUE-CAPACITY: queue capacity of ${capacity} reached`);
|
||||
}
|
||||
|
||||
const createdAt = now();
|
||||
const job: Job<T> = {
|
||||
@@ -181,7 +214,7 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
draining = true;
|
||||
try {
|
||||
const due = await store.due(now(), concurrency);
|
||||
const results = await Promise.all(due.map(runJob));
|
||||
const results = await Promise.all(due.map(beginJob));
|
||||
return results.filter(Boolean).length;
|
||||
} finally {
|
||||
draining = false;
|
||||
@@ -192,6 +225,29 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
return [...deadLetters.values()].map((job) => structuredClone(job));
|
||||
},
|
||||
|
||||
list(name) {
|
||||
return store.list(name);
|
||||
},
|
||||
|
||||
async cancel(id) {
|
||||
const running = active.get(id);
|
||||
if (running) {
|
||||
running.controller.abort();
|
||||
return true;
|
||||
}
|
||||
if (!(await store.get(id))) return false;
|
||||
await store.remove(id);
|
||||
return true;
|
||||
},
|
||||
|
||||
async shutdown(shutdownOptions = {}) {
|
||||
accepting = false;
|
||||
if (shutdownOptions.force) {
|
||||
for (const { controller } of active.values()) controller.abort();
|
||||
}
|
||||
await Promise.allSettled([...active.values()].map(({ promise }) => promise));
|
||||
},
|
||||
|
||||
async retry(id) {
|
||||
const job = deadLetters.get(id);
|
||||
if (!job) return false;
|
||||
|
||||
+104
-16
@@ -25,7 +25,14 @@ export interface Job<T = unknown> {
|
||||
createdAt: number;
|
||||
}
|
||||
|
||||
export type JobHandler<T = unknown> = (job: Job<T>) => void | Promise<void>;
|
||||
export interface JobContext {
|
||||
/** Aborted when an active job is cancelled or the queue is force-stopped. */
|
||||
signal: AbortSignal;
|
||||
/** The same trusted context shape used by HTTP, actions, realtime and webhooks. */
|
||||
execution: ExecutionContext;
|
||||
}
|
||||
|
||||
export type JobHandler<T = unknown> = (job: Job<T>, context: JobContext) => void | Promise<void>;
|
||||
|
||||
export interface AddOptions {
|
||||
/** Delay before the job becomes runnable (ms). */
|
||||
@@ -51,8 +58,11 @@ export interface QueueOptions {
|
||||
onFailed?: (job: Job, error: unknown) => void;
|
||||
/** Maximum jobs executed in one drain. Default: unlimited. */
|
||||
concurrency?: number;
|
||||
/** Maximum queued + active jobs. Adds reject once this limit is reached. */
|
||||
capacity?: number;
|
||||
/** Clock injection (tests). Default Date.now. */
|
||||
now?: () => number;
|
||||
context?: (job: Job, signal: AbortSignal) => ExecutionContext;
|
||||
}
|
||||
|
||||
export interface Queue {
|
||||
@@ -62,10 +72,14 @@ export interface Queue {
|
||||
drain(now?: number): Promise<number>;
|
||||
start(): void;
|
||||
stop(): void;
|
||||
/** Stop accepting work and wait for active handlers (or abort them). */
|
||||
shutdown(options?: { force?: boolean }): Promise<void>;
|
||||
size(): number;
|
||||
get(id: string): Job | undefined;
|
||||
list(name?: string): Job[];
|
||||
cancel(id: string): boolean;
|
||||
failed(): Job[];
|
||||
retry(id: string): Promise<boolean>;
|
||||
}
|
||||
|
||||
export interface JobDefinition<I> {
|
||||
@@ -112,6 +126,7 @@ export function createQueue(options: QueueOptions = {}): Queue {
|
||||
const backoffMs = options.backoffMs ?? 1000;
|
||||
const pollMs = options.pollMs ?? 250;
|
||||
const concurrency = options.concurrency ?? Number.POSITIVE_INFINITY;
|
||||
const capacity = options.capacity ?? Number.POSITIVE_INFINITY;
|
||||
if (!Number.isInteger(defaultMax) || defaultMax < 1)
|
||||
throw new RangeError("queue maxAttempts must be a positive integer");
|
||||
if (!Number.isFinite(backoffMs) || backoffMs < 0)
|
||||
@@ -123,13 +138,18 @@ export function createQueue(options: QueueOptions = {}): Queue {
|
||||
(Number.isInteger(concurrency) && concurrency > 0)
|
||||
))
|
||||
throw new RangeError("queue concurrency must be a positive integer");
|
||||
if (!(capacity === Number.POSITIVE_INFINITY || (Number.isInteger(capacity) && capacity > 0)))
|
||||
throw new RangeError("queue capacity must be a positive integer");
|
||||
const now = options.now ?? Date.now;
|
||||
|
||||
const jobs: Job[] = [];
|
||||
const handlers = new Map<string, JobHandler>();
|
||||
const deadLetters = new Map<string, Job>();
|
||||
const active = new Map<string, { controller: AbortController; promise: Promise<void> }>();
|
||||
let seq = 0;
|
||||
let timer: ReturnType<typeof setInterval> | null = null;
|
||||
let draining = false;
|
||||
let accepting = true;
|
||||
|
||||
async function runJob(job: Job): Promise<void> {
|
||||
const handler = handlers.get(job.name);
|
||||
@@ -137,19 +157,37 @@ export function createQueue(options: QueueOptions = {}): Queue {
|
||||
const idx = jobs.indexOf(job);
|
||||
if (idx >= 0) jobs.splice(idx, 1); // claim it
|
||||
job.attempts++;
|
||||
try {
|
||||
await handler(job);
|
||||
if (job.repeat && job.repeat > 0) {
|
||||
jobs.push({ ...job, attempts: 0, runAt: now() + job.repeat }); // recurring
|
||||
const controller = new AbortController();
|
||||
const execution = (async () => {
|
||||
try {
|
||||
await handler(job, {
|
||||
signal: controller.signal,
|
||||
execution:
|
||||
options.context?.(job, controller.signal) ??
|
||||
createExecutionContext({
|
||||
kind: job.repeat ? "cron" : "queue",
|
||||
signal: controller.signal,
|
||||
metadata: { jobId: job.id, jobName: job.name, attempt: job.attempts },
|
||||
}),
|
||||
});
|
||||
if (job.repeat && job.repeat > 0) {
|
||||
jobs.push({ ...job, attempts: 0, runAt: now() + job.repeat }); // recurring
|
||||
}
|
||||
} catch (error) {
|
||||
if (controller.signal.aborted) return;
|
||||
if (job.attempts < job.maxAttempts) {
|
||||
job.runAt = now() + backoffMs * Math.pow(2, job.attempts - 1); // exponential backoff
|
||||
jobs.push(job);
|
||||
} else {
|
||||
deadLetters.set(job.id, { ...job });
|
||||
await options.onFailed?.(job, error);
|
||||
}
|
||||
} finally {
|
||||
active.delete(job.id);
|
||||
}
|
||||
} catch (error) {
|
||||
if (job.attempts < job.maxAttempts) {
|
||||
job.runAt = now() + backoffMs * Math.pow(2, job.attempts - 1); // exponential backoff
|
||||
jobs.push(job);
|
||||
} else {
|
||||
options.onFailed?.(job, error);
|
||||
}
|
||||
}
|
||||
})();
|
||||
active.set(job.id, { controller, promise: execution });
|
||||
await execution;
|
||||
}
|
||||
|
||||
const drain: Queue["drain"] = async (at) => {
|
||||
@@ -170,6 +208,7 @@ export function createQueue(options: QueueOptions = {}): Queue {
|
||||
|
||||
return {
|
||||
async add(name, data, opts = {}) {
|
||||
if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down");
|
||||
if (!name.trim()) throw new TypeError("queue job name cannot be empty");
|
||||
if (
|
||||
opts.maxAttempts !== undefined &&
|
||||
@@ -186,6 +225,8 @@ export function createQueue(options: QueueOptions = {}): Queue {
|
||||
const existing = jobs.find((job) => job.idempotencyKey === opts.idempotencyKey);
|
||||
if (existing) return existing as Job<typeof data>;
|
||||
}
|
||||
if (jobs.length + active.size >= capacity)
|
||||
throw new Error(`WRN-QUEUE-CAPACITY: queue capacity of ${capacity} reached`);
|
||||
const createdAt = now();
|
||||
const job: Job = {
|
||||
id: `job_${++seq}`,
|
||||
@@ -207,6 +248,7 @@ export function createQueue(options: QueueOptions = {}): Queue {
|
||||
},
|
||||
drain,
|
||||
start() {
|
||||
if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down");
|
||||
if (timer) return;
|
||||
timer = setInterval(() => void drain(), pollMs);
|
||||
},
|
||||
@@ -214,16 +256,62 @@ export function createQueue(options: QueueOptions = {}): Queue {
|
||||
if (timer) clearInterval(timer);
|
||||
timer = null;
|
||||
},
|
||||
async shutdown(shutdownOptions = {}) {
|
||||
accepting = false;
|
||||
if (timer) clearInterval(timer);
|
||||
timer = null;
|
||||
if (shutdownOptions.force) {
|
||||
for (const { controller } of active.values()) controller.abort();
|
||||
}
|
||||
await Promise.allSettled([...active.values()].map(({ promise }) => promise));
|
||||
},
|
||||
size: () => jobs.length,
|
||||
get: (id) => jobs.find((job) => job.id === id),
|
||||
get: (id) => {
|
||||
const job = jobs.find((candidate) => candidate.id === id);
|
||||
return job ? { ...job } : undefined;
|
||||
},
|
||||
list: (name) => jobs.filter((job) => !name || job.name === name).map((job) => ({ ...job })),
|
||||
cancel(id) {
|
||||
const index = jobs.findIndex((job) => job.id === id);
|
||||
if (index < 0) return false;
|
||||
jobs.splice(index, 1);
|
||||
if (index >= 0) {
|
||||
jobs.splice(index, 1);
|
||||
return true;
|
||||
}
|
||||
const running = active.get(id);
|
||||
if (!running) return false;
|
||||
running.controller.abort();
|
||||
return true;
|
||||
},
|
||||
failed: () => [...deadLetters.values()].map((job) => ({ ...job })),
|
||||
async retry(id) {
|
||||
if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down");
|
||||
const job = deadLetters.get(id);
|
||||
if (!job) return false;
|
||||
deadLetters.delete(id);
|
||||
jobs.push({ ...job, attempts: 0, runAt: now() });
|
||||
return true;
|
||||
},
|
||||
};
|
||||
}
|
||||
export { memoryQueueStore, createDurableQueue } from "./durable.ts";
|
||||
export type { QueueStore, DurableQueue, DurableQueueOptions } from "./durable.ts";
|
||||
export { redisQueueStore, postgresQueueStore, POSTGRES_QUEUE_SCHEMA } from "./stores.ts";
|
||||
export type { RedisQueueClient, SqlQueueClient } from "./stores.ts";
|
||||
export {
|
||||
createQueueScheduler,
|
||||
addBatch,
|
||||
queueDashboardSnapshot,
|
||||
renderQueueDashboard,
|
||||
runQueueDaemon,
|
||||
} from "./scheduler.ts";
|
||||
export type { ScheduledJob, QueueScheduler, QueueDashboardSnapshot } from "./scheduler.ts";
|
||||
export { createWorkflowEngine, defineDurableWorkflow, memoryWorkflowStore } from "./workflow.ts";
|
||||
export type {
|
||||
WorkflowDefinition,
|
||||
WorkflowEngine,
|
||||
WorkflowRunContext,
|
||||
WorkflowSnapshot,
|
||||
WorkflowStatus,
|
||||
WorkflowStore,
|
||||
} from "./workflow.ts";
|
||||
import { createExecutionContext, type ExecutionContext } from "@wrnexus/core";
|
||||
|
||||
@@ -0,0 +1,145 @@
|
||||
import type { AddOptions, Job } from "./index.ts";
|
||||
import type { DurableQueue } from "./durable.ts";
|
||||
|
||||
export interface ScheduledJob<T = unknown> {
|
||||
name: string;
|
||||
data: T;
|
||||
everyMs: number;
|
||||
options?: AddOptions;
|
||||
}
|
||||
|
||||
export interface QueueScheduler {
|
||||
start(): void;
|
||||
stop(): void;
|
||||
tick(now?: number): Promise<number>;
|
||||
snapshot(): { running: boolean; schedules: number; nextRuns: Record<string, number> };
|
||||
}
|
||||
|
||||
/** Restart-safe scheduler when used with a durable queue and stable idempotency buckets. */
|
||||
export function createQueueScheduler(
|
||||
queue: DurableQueue,
|
||||
schedules: ScheduledJob[],
|
||||
options: { pollMs?: number; now?: () => number } = {},
|
||||
): QueueScheduler {
|
||||
const now = options.now ?? Date.now;
|
||||
const pollMs = options.pollMs ?? 1000;
|
||||
const nextRuns = new Map(schedules.map((schedule) => [schedule.name, now()]));
|
||||
let timer: ReturnType<typeof setInterval> | null = null;
|
||||
const tick = async (at = now()) => {
|
||||
let added = 0;
|
||||
for (const schedule of schedules) {
|
||||
if (!Number.isFinite(schedule.everyMs) || schedule.everyMs < 1)
|
||||
throw new RangeError(`Schedule '${schedule.name}' everyMs must be positive`);
|
||||
const next = nextRuns.get(schedule.name) ?? at;
|
||||
if (next > at) continue;
|
||||
const bucket = Math.floor(at / schedule.everyMs);
|
||||
await queue.add(schedule.name, schedule.data, {
|
||||
...schedule.options,
|
||||
idempotencyKey: schedule.options?.idempotencyKey ?? `schedule:${schedule.name}:${bucket}`,
|
||||
});
|
||||
nextRuns.set(schedule.name, (bucket + 1) * schedule.everyMs);
|
||||
added++;
|
||||
}
|
||||
return added;
|
||||
};
|
||||
return {
|
||||
start() {
|
||||
if (!timer) timer = setInterval(() => void tick(), pollMs);
|
||||
},
|
||||
stop() {
|
||||
if (timer) clearInterval(timer);
|
||||
timer = null;
|
||||
},
|
||||
tick,
|
||||
snapshot: () => ({
|
||||
running: timer !== null,
|
||||
schedules: schedules.length,
|
||||
nextRuns: Object.fromEntries(nextRuns),
|
||||
}),
|
||||
};
|
||||
}
|
||||
|
||||
export async function addBatch<T>(
|
||||
queue: DurableQueue,
|
||||
name: string,
|
||||
values: T[],
|
||||
options?: AddOptions,
|
||||
): Promise<Job<T>[]> {
|
||||
return Promise.all(
|
||||
values.map((value, index) =>
|
||||
queue.add(name, value, {
|
||||
...options,
|
||||
idempotencyKey: options?.idempotencyKey ? `${options.idempotencyKey}:${index}` : undefined,
|
||||
}),
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
export interface QueueDashboardSnapshot {
|
||||
generatedAt: number;
|
||||
pending: number;
|
||||
failed: number;
|
||||
byName: Record<string, number>;
|
||||
oldestRunAt?: number;
|
||||
}
|
||||
|
||||
export async function queueDashboardSnapshot(queue: DurableQueue): Promise<QueueDashboardSnapshot> {
|
||||
const jobs = await queue.list();
|
||||
return {
|
||||
generatedAt: Date.now(),
|
||||
pending: jobs.length,
|
||||
failed: queue.failed().length,
|
||||
byName: jobs.reduce<Record<string, number>>((counts, job) => {
|
||||
counts[job.name] = (counts[job.name] ?? 0) + 1;
|
||||
return counts;
|
||||
}, {}),
|
||||
oldestRunAt: jobs.length ? Math.min(...jobs.map((job) => job.runAt)) : undefined,
|
||||
};
|
||||
}
|
||||
|
||||
export function renderQueueDashboard(snapshot: QueueDashboardSnapshot): string {
|
||||
const rows = Object.entries(snapshot.byName)
|
||||
.sort(([left], [right]) => left.localeCompare(right))
|
||||
.map(
|
||||
([name, count]) =>
|
||||
`<tr><td>${name.replace(/[&<>"']/g, (c) => ({ "&": "&", "<": "<", ">": ">", '"': """, "'": "'" })[c]!)}</td><td>${count}</td></tr>`,
|
||||
)
|
||||
.join("");
|
||||
return `<!doctype html><html><head><meta charset="utf-8"><title>WRNexus Queue Dashboard</title></head><body><main><h1>Queue dashboard</h1><p>Pending: ${snapshot.pending} · Failed: ${snapshot.failed}</p><table><thead><tr><th>Queue</th><th>Pending</th></tr></thead><tbody>${rows}</tbody></table></main></body></html>`;
|
||||
}
|
||||
|
||||
/** Long-running scheduler/worker loop suitable for a dedicated process or container. */
|
||||
export async function runQueueDaemon(
|
||||
queue: DurableQueue,
|
||||
scheduler: QueueScheduler,
|
||||
options: { signal?: AbortSignal; pollMs?: number; onError?: (error: unknown) => void } = {},
|
||||
): Promise<void> {
|
||||
const pollMs = options.pollMs ?? 250;
|
||||
if (!Number.isInteger(pollMs) || pollMs < 10)
|
||||
throw new RangeError("Daemon pollMs must be at least 10ms");
|
||||
scheduler.start();
|
||||
try {
|
||||
while (!options.signal?.aborted) {
|
||||
try {
|
||||
await scheduler.tick();
|
||||
await queue.drain();
|
||||
} catch (error) {
|
||||
options.onError?.(error);
|
||||
}
|
||||
await new Promise<void>((resolve) => {
|
||||
const timer = setTimeout(resolve, pollMs);
|
||||
options.signal?.addEventListener(
|
||||
"abort",
|
||||
() => {
|
||||
clearTimeout(timer);
|
||||
resolve();
|
||||
},
|
||||
{ once: true },
|
||||
);
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
scheduler.stop();
|
||||
await queue.shutdown();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
import type { Job } from "./index.ts";
|
||||
import type { QueueStore } from "./durable.ts";
|
||||
|
||||
export interface RedisQueueClient {
|
||||
get(key: string): Promise<string | null>;
|
||||
set(key: string, value: string, options?: { NX?: boolean; PX?: number }): Promise<unknown>;
|
||||
del(...keys: string[]): Promise<unknown>;
|
||||
zadd(key: string, score: number, member: string): Promise<unknown>;
|
||||
zrem(key: string, member: string): Promise<unknown>;
|
||||
zrangebyscore(
|
||||
key: string,
|
||||
min: number,
|
||||
max: number,
|
||||
options?: { limit: [number, number] },
|
||||
): Promise<string[]>;
|
||||
smembers(key: string): Promise<string[]>;
|
||||
sadd(key: string, member: string): Promise<unknown>;
|
||||
srem(key: string, member: string): Promise<unknown>;
|
||||
}
|
||||
|
||||
/** Redis-backed queue store using only the common client command surface. */
|
||||
export function redisQueueStore(client: RedisQueueClient, prefix = "wrnexus:queue"): QueueStore {
|
||||
const jobs = `${prefix}:jobs`;
|
||||
const due = `${prefix}:due`;
|
||||
const key = (id: string) => `${prefix}:job:${id}`;
|
||||
return {
|
||||
async put(job) {
|
||||
await client.set(key(job.id), JSON.stringify(job));
|
||||
await client.sadd(jobs, job.id);
|
||||
await client.zadd(due, job.runAt, job.id);
|
||||
await client.del(`${prefix}:lease:${job.id}`);
|
||||
},
|
||||
async get(id) {
|
||||
const value = await client.get(key(id));
|
||||
return value ? (JSON.parse(value) as Job) : null;
|
||||
},
|
||||
async remove(id) {
|
||||
await client.del(key(id), `${prefix}:lease:${id}`);
|
||||
await client.srem(jobs, id);
|
||||
await client.zrem(due, id);
|
||||
},
|
||||
async due(now, limit) {
|
||||
const ids = await client.zrangebyscore(due, 0, now, { limit: [0, limit] });
|
||||
const values = await Promise.all(ids.map((id) => client.get(key(id))));
|
||||
return values
|
||||
.filter((value): value is string => value !== null)
|
||||
.map((value) => JSON.parse(value));
|
||||
},
|
||||
async list(name) {
|
||||
const ids = await client.smembers(jobs);
|
||||
const values = await Promise.all(ids.map((id) => client.get(key(id))));
|
||||
return values
|
||||
.filter((value): value is string => value !== null)
|
||||
.map((value) => JSON.parse(value) as Job)
|
||||
.filter((job) => !name || job.name === name);
|
||||
},
|
||||
async claim(id, worker, leaseUntil) {
|
||||
const ttl = Math.max(1, leaseUntil - Date.now());
|
||||
return Boolean(await client.set(`${prefix}:lease:${id}`, worker, { NX: true, PX: ttl }));
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export interface SqlQueueClient {
|
||||
query<T = Record<string, unknown>>(sql: string, parameters?: unknown[]): Promise<{ rows: T[] }>;
|
||||
}
|
||||
|
||||
/** PostgreSQL store with atomic SKIP LOCKED leasing and JSON payloads. */
|
||||
export function postgresQueueStore(db: SqlQueueClient, table = "wrnexus_jobs"): QueueStore {
|
||||
if (!/^[a-z_][a-z0-9_]*$/i.test(table)) throw new Error("Invalid queue table name");
|
||||
return {
|
||||
async put(job) {
|
||||
await db.query(
|
||||
`INSERT INTO ${table} (id,name,payload,run_at,priority,lease_owner,lease_until) VALUES ($1,$2,$3,$4,$5,NULL,NULL) ON CONFLICT (id) DO UPDATE SET name=$2,payload=$3,run_at=$4,priority=$5,lease_owner=NULL,lease_until=NULL`,
|
||||
[job.id, job.name, JSON.stringify(job), job.runAt, job.priority],
|
||||
);
|
||||
},
|
||||
async get(id) {
|
||||
const result = await db.query<{ payload: Job }>(`SELECT payload FROM ${table} WHERE id=$1`, [
|
||||
id,
|
||||
]);
|
||||
return result.rows[0]?.payload ?? null;
|
||||
},
|
||||
async remove(id) {
|
||||
await db.query(`DELETE FROM ${table} WHERE id=$1`, [id]);
|
||||
},
|
||||
async due(now, limit) {
|
||||
const result = await db.query<{ payload: Job }>(
|
||||
`SELECT payload FROM ${table} WHERE run_at <= $1 AND (lease_until IS NULL OR lease_until < $1) ORDER BY priority DESC, run_at ASC LIMIT $2`,
|
||||
[now, limit],
|
||||
);
|
||||
return result.rows.map((row) => row.payload);
|
||||
},
|
||||
async list(name) {
|
||||
const result = await db.query<{ payload: Job }>(
|
||||
`SELECT payload FROM ${table}${name ? " WHERE name=$1" : ""} ORDER BY run_at ASC`,
|
||||
name ? [name] : [],
|
||||
);
|
||||
return result.rows.map((row) => row.payload);
|
||||
},
|
||||
async claim(id, worker, leaseUntil) {
|
||||
const result = await db.query(
|
||||
`UPDATE ${table} SET lease_owner=$2,lease_until=$3 WHERE id=$1 AND (lease_until IS NULL OR lease_until < $4) RETURNING id`,
|
||||
[id, worker, leaseUntil, Date.now()],
|
||||
);
|
||||
return result.rows.length === 1;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export const POSTGRES_QUEUE_SCHEMA = `CREATE TABLE IF NOT EXISTS wrnexus_jobs (
|
||||
id text PRIMARY KEY, name text NOT NULL, payload jsonb NOT NULL, run_at bigint NOT NULL,
|
||||
priority integer NOT NULL DEFAULT 0, lease_owner text, lease_until bigint
|
||||
); CREATE INDEX IF NOT EXISTS wrnexus_jobs_due ON wrnexus_jobs (run_at, priority DESC);`;
|
||||
@@ -0,0 +1,199 @@
|
||||
export type WorkflowStatus =
|
||||
"pending" | "running" | "waiting-approval" | "completed" | "failed" | "cancelled";
|
||||
export interface WorkflowStep<I = unknown, O = unknown> {
|
||||
name: string;
|
||||
dependsOn?: string[];
|
||||
approval?: boolean;
|
||||
run(input: I, context: WorkflowRunContext): O | Promise<O>;
|
||||
}
|
||||
export interface WorkflowRunContext {
|
||||
workflowId: string;
|
||||
step: string;
|
||||
results: Readonly<Record<string, unknown>>;
|
||||
signal: AbortSignal;
|
||||
progress(value: number, message?: string): void;
|
||||
}
|
||||
export interface WorkflowSnapshot {
|
||||
id: string;
|
||||
name: string;
|
||||
status: WorkflowStatus;
|
||||
input: unknown;
|
||||
results: Record<string, unknown>;
|
||||
completed: string[];
|
||||
waitingFor?: string;
|
||||
progress: number;
|
||||
message?: string;
|
||||
error?: string;
|
||||
updatedAt: number;
|
||||
}
|
||||
export interface WorkflowStore {
|
||||
get(id: string): Promise<WorkflowSnapshot | null>;
|
||||
put(snapshot: WorkflowSnapshot): Promise<void>;
|
||||
list(): Promise<WorkflowSnapshot[]>;
|
||||
}
|
||||
|
||||
export function memoryWorkflowStore(): WorkflowStore {
|
||||
const values = new Map<string, WorkflowSnapshot>();
|
||||
return {
|
||||
async get(id) {
|
||||
const value = values.get(id);
|
||||
return value ? structuredClone(value) : null;
|
||||
},
|
||||
async put(value) {
|
||||
values.set(value.id, structuredClone(value));
|
||||
},
|
||||
async list() {
|
||||
return [...values.values()].map((value) => structuredClone(value));
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export interface WorkflowDefinition<I = unknown> {
|
||||
name: string;
|
||||
steps: WorkflowStep<any, any>[];
|
||||
/** Compile-time input marker; definitions do not store runtime input values. */
|
||||
readonly __input?: I;
|
||||
}
|
||||
export interface WorkflowEngine {
|
||||
start<I>(definition: WorkflowDefinition<I>, input: I, id?: string): Promise<WorkflowSnapshot>;
|
||||
resume<I>(definition: WorkflowDefinition<I>, id: string): Promise<WorkflowSnapshot>;
|
||||
approve<I>(
|
||||
definition: WorkflowDefinition<I>,
|
||||
id: string,
|
||||
step: string,
|
||||
actor: string,
|
||||
): Promise<WorkflowSnapshot>;
|
||||
cancel(id: string): Promise<boolean>;
|
||||
get(id: string): Promise<WorkflowSnapshot | null>;
|
||||
list(): Promise<WorkflowSnapshot[]>;
|
||||
}
|
||||
|
||||
function validate(definition: WorkflowDefinition): void {
|
||||
const names = new Set(definition.steps.map((step) => step.name));
|
||||
if (names.size !== definition.steps.length) throw new Error("WRN-WORKFLOW-DUPLICATE-STEP");
|
||||
for (const step of definition.steps)
|
||||
for (const dependency of step.dependsOn ?? [])
|
||||
if (!names.has(dependency))
|
||||
throw new Error(
|
||||
`WRN-WORKFLOW-DEPENDENCY: '${step.name}' depends on missing '${dependency}'.`,
|
||||
);
|
||||
const visit = (name: string, path: Set<string>): void => {
|
||||
if (path.has(name)) throw new Error(`WRN-WORKFLOW-CYCLE: ${[...path, name].join(" -> ")}`);
|
||||
const next = new Set(path).add(name);
|
||||
const step = definition.steps.find((value) => value.name === name)!;
|
||||
for (const dependency of step.dependsOn ?? []) visit(dependency, next);
|
||||
};
|
||||
for (const step of definition.steps) visit(step.name, new Set());
|
||||
}
|
||||
|
||||
export function createWorkflowEngine(store: WorkflowStore = memoryWorkflowStore()): WorkflowEngine {
|
||||
const controllers = new Map<string, AbortController>();
|
||||
async function execute<I>(
|
||||
definition: WorkflowDefinition<I>,
|
||||
snapshot: WorkflowSnapshot,
|
||||
): Promise<WorkflowSnapshot> {
|
||||
validate(definition);
|
||||
const controller = new AbortController();
|
||||
controllers.set(snapshot.id, controller);
|
||||
snapshot.status = "running";
|
||||
await store.put(snapshot);
|
||||
try {
|
||||
while (snapshot.completed.length < definition.steps.length) {
|
||||
const ready = definition.steps.filter(
|
||||
(step) =>
|
||||
!snapshot.completed.includes(step.name) &&
|
||||
(step.dependsOn ?? []).every((dependency) => snapshot.completed.includes(dependency)),
|
||||
);
|
||||
if (!ready.length) throw new Error("WRN-WORKFLOW-BLOCKED: no runnable steps.");
|
||||
const step = ready[0]!;
|
||||
if (step.approval && snapshot.waitingFor !== `approved:${step.name}`) {
|
||||
snapshot.status = "waiting-approval";
|
||||
snapshot.waitingFor = step.name;
|
||||
snapshot.updatedAt = Date.now();
|
||||
await store.put(snapshot);
|
||||
return structuredClone(snapshot);
|
||||
}
|
||||
snapshot.waitingFor = undefined;
|
||||
const dependencies = step.dependsOn ?? [];
|
||||
const value =
|
||||
dependencies.length === 1
|
||||
? snapshot.results[dependencies[0]!]
|
||||
: dependencies.length
|
||||
? Object.fromEntries(dependencies.map((name) => [name, snapshot.results[name]]))
|
||||
: snapshot.input;
|
||||
snapshot.results[step.name] = await step.run(value, {
|
||||
workflowId: snapshot.id,
|
||||
step: step.name,
|
||||
results: snapshot.results,
|
||||
signal: controller.signal,
|
||||
progress(value, message) {
|
||||
snapshot.progress = Math.max(0, Math.min(100, value));
|
||||
snapshot.message = message;
|
||||
snapshot.updatedAt = Date.now();
|
||||
void store.put(snapshot);
|
||||
},
|
||||
});
|
||||
snapshot.completed.push(step.name);
|
||||
snapshot.progress = Math.round((snapshot.completed.length / definition.steps.length) * 100);
|
||||
snapshot.updatedAt = Date.now();
|
||||
await store.put(snapshot);
|
||||
}
|
||||
snapshot.status = "completed";
|
||||
snapshot.progress = 100;
|
||||
} catch (error) {
|
||||
snapshot.status = controller.signal.aborted ? "cancelled" : "failed";
|
||||
snapshot.error = error instanceof Error ? error.message : String(error);
|
||||
} finally {
|
||||
snapshot.updatedAt = Date.now();
|
||||
controllers.delete(snapshot.id);
|
||||
await store.put(snapshot);
|
||||
}
|
||||
return structuredClone(snapshot);
|
||||
}
|
||||
return {
|
||||
async start(definition, input, id = `workflow-${crypto.randomUUID()}`) {
|
||||
if (await store.get(id)) throw new Error(`WRN-WORKFLOW-ID: '${id}' already exists.`);
|
||||
return execute(definition, {
|
||||
id,
|
||||
name: definition.name,
|
||||
status: "pending",
|
||||
input,
|
||||
results: {},
|
||||
completed: [],
|
||||
progress: 0,
|
||||
updatedAt: Date.now(),
|
||||
});
|
||||
},
|
||||
async resume(definition, id) {
|
||||
const snapshot = await store.get(id);
|
||||
if (!snapshot) throw new Error(`WRN-WORKFLOW-NOT-FOUND: '${id}'.`);
|
||||
if (["completed", "cancelled"].includes(snapshot.status)) return snapshot;
|
||||
return execute(definition, snapshot);
|
||||
},
|
||||
async approve(definition, id, step, actor) {
|
||||
const snapshot = await store.get(id);
|
||||
if (!snapshot || snapshot.status !== "waiting-approval" || snapshot.waitingFor !== step)
|
||||
throw new Error(`WRN-WORKFLOW-APPROVAL: '${step}' is not awaiting approval.`);
|
||||
snapshot.waitingFor = `approved:${step}`;
|
||||
snapshot.results[`${step}:approval`] = { actor, approvedAt: Date.now() };
|
||||
await store.put(snapshot);
|
||||
return execute(definition, snapshot);
|
||||
},
|
||||
async cancel(id) {
|
||||
const snapshot = await store.get(id);
|
||||
if (!snapshot || ["completed", "cancelled"].includes(snapshot.status)) return false;
|
||||
controllers.get(id)?.abort();
|
||||
snapshot.status = "cancelled";
|
||||
snapshot.updatedAt = Date.now();
|
||||
await store.put(snapshot);
|
||||
return true;
|
||||
},
|
||||
get: (id) => store.get(id),
|
||||
list: () => store.list(),
|
||||
};
|
||||
}
|
||||
|
||||
export function defineDurableWorkflow<I>(definition: WorkflowDefinition<I>): WorkflowDefinition<I> {
|
||||
validate(definition);
|
||||
return definition;
|
||||
}
|
||||
Reference in New Issue
Block a user