import { describe, it, expect, vi } from "vitest"; import { z } from "zod"; import { defineEvent } from "@/event-descriptor"; import { PayloadJobsEventBus } from "@/payload-jobs-event-bus"; import type { IJobQueue } from "@repo/core-shared/jobs"; const evt = defineEvent("auth.user.signed-up", z.object({ userId: z.string() }).strict()); function recordingQueue(): IJobQueue & { enqueued: { taskSlug: string; input: unknown }[] } { const enqueued: { taskSlug: string; input: unknown }[] = []; const q: IJobQueue = { async enqueue(taskSlug, input) { enqueued.push({ taskSlug, input }); return { jobId: `recording-${enqueued.length}` }; }, }; return Object.assign(q, { enqueued }); } describe("PayloadJobsEventBus", () => { it("validates the payload before enqueueing", async () => { const queue = recordingQueue(); const bus = new PayloadJobsEventBus(queue); bus.subscribe(evt, "marketing-pages", vi.fn()); await expect( bus.publish(evt, { userId: 42 } as unknown as { userId: string }), ).rejects.toThrow(); expect(queue.enqueued).toHaveLength(0); }); it("enqueues one task per subscriber, naming `__events..`", async () => { const queue = recordingQueue(); const bus = new PayloadJobsEventBus(queue); bus.subscribe(evt, "marketing-pages", vi.fn()); bus.subscribe(evt, "blog", vi.fn()); await bus.publish(evt, { userId: "u1" }); expect(queue.enqueued).toHaveLength(2); expect(queue.enqueued.map((e) => e.taskSlug).sort()).toEqual([ "__events.auth.user.signed-up.blog", "__events.auth.user.signed-up.marketing-pages", ]); expect(queue.enqueued[0]!.input).toEqual({ userId: "u1" }); }); it("enqueues nothing when no subscribers are registered", async () => { const queue = recordingQueue(); const bus = new PayloadJobsEventBus(queue); await bus.publish(evt, { userId: "u1" }); expect(queue.enqueued).toHaveLength(0); }); });