Files
agentic-dev/packages/core-events/src/payload-jobs-event-bus.test.ts
Danijel Martinek ef451de437 feat(core-events): PayloadJobsEventBus (fan-out via IJobQueue)
TDD: red (test-only), then green. PayloadJobsEventBus validates before
enqueueing, names tasks __events.<event>.<consumer> deterministically, and
enqueues one task per subscriber via Promise.all. 3 new tests (11 total).

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-08 11:57:48 +02:00

52 lines
1.9 KiB
TypeScript

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.<event>.<consumer>`", 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);
});
});