Files
Clintchiz 586a6db8ff
Quality / quality (ubuntu-latest) (push) Failing after 21s
Quality / quality (windows-latest) (push) Canceled after 0s
release: WRNexusJS 0.8.0
2026-08-02 23:18:51 +05:30

133 lines
4.0 KiB
TypeScript

import { expect, test } from "bun:test";
import {
addBatch,
createDurableQueue,
createQueueScheduler,
memoryQueueStore,
postgresQueueStore,
queueDashboardSnapshot,
redisQueueStore,
runQueueDaemon,
type Job,
type RedisQueueClient,
} from "../src/index.ts";
const job: Job = {
id: "one",
name: "email",
data: { to: "a@b.test" },
attempts: 0,
maxAttempts: 3,
runAt: 10,
priority: 2,
createdAt: 1,
};
test("Redis queue store persists, orders and atomically leases jobs", async () => {
const strings = new Map<string, string>();
const sets = new Map<string, Set<string>>();
const scores = new Map<string, Map<string, number>>();
const client: RedisQueueClient = {
async get(key) {
return strings.get(key) ?? null;
},
async set(key, value, options) {
if (options?.NX && strings.has(key)) return null;
strings.set(key, value);
return "OK";
},
async del(...keys) {
for (const key of keys) strings.delete(key);
return keys.length;
},
async zadd(key, score, member) {
(scores.get(key) ?? scores.set(key, new Map()).get(key)!).set(member, score);
return 1;
},
async zrem(key, member) {
return scores.get(key)?.delete(member) ? 1 : 0;
},
async zrangebyscore(key, min, max, options) {
return [...(scores.get(key) ?? [])]
.filter(([, score]) => score >= min && score <= max)
.sort((a, b) => a[1] - b[1])
.slice(...options!.limit)
.map(([id]) => id);
},
async smembers(key) {
return [...(sets.get(key) ?? [])];
},
async sadd(key, member) {
(sets.get(key) ?? sets.set(key, new Set()).get(key)!).add(member);
return 1;
},
async srem(key, member) {
return sets.get(key)?.delete(member) ? 1 : 0;
},
};
const store = redisQueueStore(client);
await store.put(job);
expect(await store.due(10, 5)).toEqual([job]);
expect(await store.claim!(job.id, "a", Date.now() + 1000)).toBeTrue();
expect(await store.claim!(job.id, "b", Date.now() + 1000)).toBeFalse();
await store.remove(job.id);
expect(await store.list()).toEqual([]);
});
test("standalone daemon schedules, drains and shuts down on abort", async () => {
const queue = createDurableQueue();
let processed = 0;
queue.process("pulse", async () => {
processed++;
});
const scheduler = createQueueScheduler(queue, [{ name: "pulse", data: {}, everyMs: 10 }], {
pollMs: 10,
});
const controller = new AbortController();
setTimeout(() => controller.abort(), 35);
await runQueueDaemon(queue, scheduler, { signal: controller.signal, pollMs: 10 });
expect(processed).toBeGreaterThan(0);
expect(scheduler.snapshot().running).toBeFalse();
});
test("PostgreSQL queue store parameterizes values and validates identifiers", async () => {
const calls: Array<{ sql: string; args?: unknown[] }> = [];
const store = postgresQueueStore({
async query(sql, args) {
calls.push({ sql, args });
return { rows: [] };
},
});
await store.put(job);
expect(calls[0]!.sql).toContain("ON CONFLICT");
expect(calls[0]!.args?.[0]).toBe("one");
expect(() =>
postgresQueueStore(
{
async query() {
return { rows: [] };
},
},
"jobs; DROP TABLE jobs",
),
).toThrow("Invalid queue table");
});
test("scheduler, batch and dashboard provide operational queue APIs", async () => {
let now = 1_000;
const queue = createDurableQueue({ store: memoryQueueStore(), now: () => now });
await addBatch(queue, "email", [{ id: 1 }, { id: 2 }], { idempotencyKey: "emails" });
const scheduler = createQueueScheduler(queue, [{ name: "cleanup", data: {}, everyMs: 500 }], {
now: () => now,
});
expect(await scheduler.tick()).toBe(1);
expect(await scheduler.tick()).toBe(0);
now = 1_500;
expect(await scheduler.tick()).toBe(1);
expect(await queueDashboardSnapshot(queue)).toMatchObject({
pending: 4,
byName: { email: 2, cleanup: 2 },
});
expect(scheduler.snapshot().schedules).toBe(1);
});