feat: centralize application framework primitives
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@wrnexus/queue",
|
||||
"version": "0.8.8",
|
||||
"version": "0.8.9",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "src/index.ts",
|
||||
|
||||
@@ -14,6 +14,8 @@ export interface QueueStore {
|
||||
remove(id: string): Promise<void>;
|
||||
due(now: number, limit: number): Promise<Job[]>;
|
||||
list(name?: string): Promise<Job[]>;
|
||||
findByIdempotencyKey?(name: string, key: string): Promise<Job | null>;
|
||||
size?(): Promise<number>;
|
||||
claim?(id: string, worker: string, leaseUntil: number): Promise<boolean>;
|
||||
}
|
||||
|
||||
@@ -48,6 +50,15 @@ export function memoryQueueStore(): QueueStore {
|
||||
.filter((job) => !name || job.name === name)
|
||||
.map((job) => structuredClone(job));
|
||||
},
|
||||
async findByIdempotencyKey(name, key) {
|
||||
const job = [...jobs.values()].find(
|
||||
(candidate) => candidate.name === name && candidate.idempotencyKey === key,
|
||||
);
|
||||
return job ? structuredClone(job) : null;
|
||||
},
|
||||
async size() {
|
||||
return jobs.size;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
@@ -73,6 +84,10 @@ export interface DurableQueue {
|
||||
list(name?: string): Promise<Job[]>;
|
||||
failed(): Job[];
|
||||
retry(id: string): Promise<boolean>;
|
||||
addBatch<T>(
|
||||
name: string,
|
||||
entries: readonly { data: T; options?: AddOptions }[],
|
||||
): Promise<Job<T>[]>;
|
||||
}
|
||||
|
||||
function positiveInteger(value: number, label: string): number {
|
||||
@@ -178,12 +193,12 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
}
|
||||
|
||||
if (add.idempotencyKey) {
|
||||
const existing = (await store.list(name)).find(
|
||||
(job) => job.idempotencyKey === add.idempotencyKey,
|
||||
);
|
||||
const existing = store.findByIdempotencyKey
|
||||
? await store.findByIdempotencyKey(name, add.idempotencyKey)
|
||||
: (await store.list(name)).find((job) => job.idempotencyKey === add.idempotencyKey);
|
||||
if (existing) return existing as Job<T>;
|
||||
}
|
||||
if ((await store.list()).length >= capacity) {
|
||||
if ((store.size ? await store.size() : (await store.list()).length) >= capacity) {
|
||||
throw new Error(`WRN-QUEUE-CAPACITY: queue capacity of ${capacity} reached`);
|
||||
}
|
||||
|
||||
@@ -204,6 +219,12 @@ export function createDurableQueue(options: DurableQueueOptions = {}): DurableQu
|
||||
return structuredClone(job);
|
||||
},
|
||||
|
||||
async addBatch<T>(name: string, entries: readonly { data: T; options?: AddOptions }[]) {
|
||||
const output: Job<T>[] = [];
|
||||
for (const entry of entries) output.push(await this.add(name, entry.data, entry.options));
|
||||
return output;
|
||||
},
|
||||
|
||||
process(name, handler) {
|
||||
if (!name.trim()) throw new TypeError("queue worker name cannot be empty");
|
||||
handlers.set(name, handler as JobHandler);
|
||||
|
||||
@@ -92,6 +92,9 @@ export function defineJob<I>(definition: JobDefinition<I>): JobDefinition<I> {
|
||||
return definition;
|
||||
}
|
||||
|
||||
export { defineWorker, runWorker } from "./worker.ts";
|
||||
export type { WorkerDefinition } from "./worker.ts";
|
||||
|
||||
export interface WorkflowStep<I, O> {
|
||||
name: string;
|
||||
run(input: I): O | Promise<O>;
|
||||
|
||||
@@ -54,6 +54,18 @@ export function redisQueueStore(client: RedisQueueClient, prefix = "wrnexus:queu
|
||||
.map((value) => JSON.parse(value) as Job)
|
||||
.filter((job) => !name || job.name === name);
|
||||
},
|
||||
async findByIdempotencyKey(name, idempotencyKey) {
|
||||
const ids = await client.smembers(jobs);
|
||||
const values = await Promise.all(ids.map((id) => client.get(key(id))));
|
||||
const found = values
|
||||
.filter((value): value is string => value !== null)
|
||||
.map((value) => JSON.parse(value) as Job)
|
||||
.find((job) => job.name === name && job.idempotencyKey === idempotencyKey);
|
||||
return found ?? null;
|
||||
},
|
||||
async size() {
|
||||
return (await client.smembers(jobs)).length;
|
||||
},
|
||||
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 }));
|
||||
@@ -98,6 +110,19 @@ export function postgresQueueStore(db: SqlQueueClient, table = "wrnexus_jobs"):
|
||||
);
|
||||
return result.rows.map((row) => row.payload);
|
||||
},
|
||||
async findByIdempotencyKey(name, key) {
|
||||
const result = await db.query<{ payload: Job }>(
|
||||
`SELECT payload FROM ${table} WHERE name=$1 AND payload->>'idempotencyKey'=$2 LIMIT 1`,
|
||||
[name, key],
|
||||
);
|
||||
return result.rows[0]?.payload ?? null;
|
||||
},
|
||||
async size() {
|
||||
const result = await db.query<{ total: number | string }>(
|
||||
`SELECT COUNT(*) AS total FROM ${table}`,
|
||||
);
|
||||
return Number(result.rows[0]?.total ?? 0);
|
||||
},
|
||||
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`,
|
||||
|
||||
@@ -0,0 +1,36 @@
|
||||
import type { JobHandler } from "./index.ts";
|
||||
import type { DurableQueue } from "./durable.ts";
|
||||
|
||||
export interface WorkerDefinition<T> {
|
||||
name: string;
|
||||
handler: JobHandler<T>;
|
||||
onStart?: () => void | Promise<void>;
|
||||
onStop?: () => void | Promise<void>;
|
||||
}
|
||||
|
||||
export function defineWorker<T>(definition: WorkerDefinition<T>): WorkerDefinition<T> {
|
||||
if (!definition.name.trim()) throw new TypeError("worker name cannot be empty");
|
||||
return Object.freeze({ ...definition });
|
||||
}
|
||||
|
||||
/** Explicit process lifecycle for queue workers; safe to start and stop repeatedly. */
|
||||
export function runWorker<T>(queue: DurableQueue, definition: WorkerDefinition<T>) {
|
||||
let started = false;
|
||||
return {
|
||||
async start() {
|
||||
if (started) return;
|
||||
started = true;
|
||||
queue.process(definition.name, definition.handler);
|
||||
await definition.onStart?.();
|
||||
},
|
||||
async stop(options?: { force?: boolean }) {
|
||||
if (!started) return;
|
||||
started = false;
|
||||
await queue.shutdown(options);
|
||||
await definition.onStop?.();
|
||||
},
|
||||
get running() {
|
||||
return started;
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -5,8 +5,27 @@ import {
|
||||
cronToInterval,
|
||||
defineWorkflow,
|
||||
memoryQueueStore,
|
||||
defineWorker,
|
||||
runWorker,
|
||||
} from "../src/index.ts";
|
||||
|
||||
test("durable queues batch enqueue and worker lifecycle are idempotent", async () => {
|
||||
const queue = createDurableQueue();
|
||||
const seen: number[] = [];
|
||||
const worker = runWorker(
|
||||
queue,
|
||||
defineWorker<number>({ name: "number", handler: (job) => void seen.push(job.data) }),
|
||||
);
|
||||
await worker.start();
|
||||
await worker.start();
|
||||
const jobs = await queue.addBatch("number", [{ data: 1 }, { data: 2 }]);
|
||||
expect(jobs).toHaveLength(2);
|
||||
await queue.drain();
|
||||
expect(seen).toEqual([1, 2]);
|
||||
await worker.stop();
|
||||
expect(worker.running).toBe(false);
|
||||
});
|
||||
|
||||
test("processes a job", async () => {
|
||||
const queue = createQueue();
|
||||
const done: string[] = [];
|
||||
|
||||
Reference in New Issue
Block a user