133 lines
4.0 KiB
TypeScript
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);
|
|
});
|