diff --git a/packages/cfworkers/test/mq.test.ts b/packages/cfworkers/test/mq.test.ts index a46124493..8ca715918 100644 --- a/packages/cfworkers/test/mq.test.ts +++ b/packages/cfworkers/test/mq.test.ts @@ -229,6 +229,87 @@ describe("WorkersMessageQueue", () => { sendBatchSpy.mockRestore(); }); + describe("message wrapper round trip", () => { + it("round-trips a message through enqueue() and processMessage()", async () => { + const sendSpy = vi + .spyOn(env.Q1, "send") + .mockImplementation(async () => {}); + const getSpy = vi.spyOn(env.KV1, "get").mockResolvedValue(null); + const putSpy = vi.spyOn(env.KV1, "put").mockResolvedValue(undefined); + const deleteSpy = vi.spyOn(env.KV1, "delete").mockResolvedValue(undefined); + + const queue = new WorkersMessageQueue(env.Q1, { orderingKv: env.KV1 }); + const payload = { foo: 1 }; + await queue.enqueue(payload, { orderingKey: "key1" }); + + expect(sendSpy).toHaveBeenCalledTimes(1); + + const requestMessage = sendSpy.mock.calls[0][0]; + + const result = await queue.processMessage(requestMessage); + + expect(result.shouldProcess).toBe(true); + expect(result.message).toEqual(payload); + expect(result.release).toBeDefined(); + + expect(getSpy).toHaveBeenCalledWith("__fedify_ordering_key1"); + expect(putSpy).toHaveBeenCalledWith( + "__fedify_ordering_key1", + expect.any(String), + { expirationTtl: 60 }, + ); + + // Release the lock + await result.release!(); + expect(deleteSpy).toHaveBeenCalledWith("__fedify_ordering_key1"); + + sendSpy.mockRestore(); + getSpy.mockRestore(); + putSpy.mockRestore(); + deleteSpy.mockRestore(); + }); + + it("round-trips a message through enqueueMany() and processMessage()", async () => { + const sendBatchSpy = vi + .spyOn(env.Q1, "sendBatch") + .mockImplementation(async () => {}); + const getSpy = vi.spyOn(env.KV1, "get").mockResolvedValue(null); + const putSpy = vi.spyOn(env.KV1, "put").mockResolvedValue(undefined); + const deleteSpy = vi.spyOn(env.KV1, "delete").mockResolvedValue(undefined); + + const queue = new WorkersMessageQueue(env.Q1, { orderingKv: env.KV1 }); + const payloads = [{ a: 1 }, { b: 2 }]; + await queue.enqueueMany(payloads, { orderingKey: "batch-key" }); + + expect(sendBatchSpy).toHaveBeenCalledTimes(1); + + const requests = sendBatchSpy.mock.calls[0][0]; + const firstRequestMessage = requests[0].body; + + const result = await queue.processMessage(firstRequestMessage); + + expect(result.shouldProcess).toBe(true); + expect(result.message).toEqual(payloads[0]); + expect(result.release).toBeDefined(); + + expect(getSpy).toHaveBeenCalledWith("__fedify_ordering_batch-key"); + expect(putSpy).toHaveBeenCalledWith( + "__fedify_ordering_batch-key", + expect.any(String), + { expirationTtl: 60 }, + ); + + // Release the lock + await result.release!(); + expect(deleteSpy).toHaveBeenCalledWith("__fedify_ordering_batch-key"); + + sendBatchSpy.mockRestore(); + getSpy.mockRestore(); + putSpy.mockRestore(); + deleteSpy.mockRestore(); + }); + }); + it("listen() throws TypeError", () => { const queue = new WorkersMessageQueue(env.Q1); expect(() => queue.listen(() => {})).toThrow(TypeError);