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(); const sets = new Map>(); const scores = new Map>(); 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); });