diff --git a/packages/jobs/src/worker.test.ts b/packages/jobs/src/worker.test.ts index 893e238..76a3020 100644 --- a/packages/jobs/src/worker.test.ts +++ b/packages/jobs/src/worker.test.ts @@ -6,6 +6,7 @@ function createMockBoss() { return { start: vi.fn().mockResolvedValue(undefined), stop: vi.fn().mockResolvedValue(undefined), + createQueue: vi.fn().mockResolvedValue(undefined), schedule: vi.fn().mockResolvedValue(undefined), work: vi.fn().mockResolvedValue("worker-id"), }; @@ -22,6 +23,16 @@ describe("startWorker", () => { expect(boss.start).toHaveBeenCalled(); }); + it("creates queues before registering handlers", async () => { + const boss = createMockBoss(); + const jobs: JobDefinition[] = [{ name: "test-job", handler: vi.fn() }]; + + await startWorker(boss as never, jobs); + + expect(boss.createQueue).toHaveBeenCalledWith("test-job"); + expect(boss.createQueue).toHaveBeenCalledBefore(boss.work); + }); + it("registers job handlers", async () => { const boss = createMockBoss(); const handler = vi.fn(); diff --git a/packages/jobs/src/worker.ts b/packages/jobs/src/worker.ts index 63a90c2..ee56806 100644 --- a/packages/jobs/src/worker.ts +++ b/packages/jobs/src/worker.ts @@ -8,6 +8,8 @@ export async function startWorker( await boss.start(); for (const job of jobs) { + await boss.createQueue(job.name); + if (job.cron) { await boss.schedule(job.name, job.cron, undefined, { retryLimit: job.retryLimit,