File
Blob: test/util/queue.ts
| 1 | import { createExecutionContext, createMessageBatch, getQueueResult } from "cloudflare:test"; |
| 2 | import { env as testEnv } from "cloudflare:workers"; |
| 3 | import worker from "@/worker/index"; |
| 4 | |
| 5 | export type QueueRunResult = { |
| 6 | acked: boolean; |
| 7 | retried: boolean; |
| 8 | }; |
| 9 | |
| 10 | // `getQueueResult()` returns more state than most tests need. Keep the |
| 11 | // asserted subset named here so the helper stays typed without reaching for |
| 12 | // casts when the generated Worker types do not expose `FetcherQueueResult`. |
| 13 | type QueueResultState = { |
| 14 | ackAll: boolean; |
| 15 | explicitAcks: string[]; |
| 16 | retryBatch: { |
| 17 | retry: boolean; |
| 18 | }; |
| 19 | retryMessages: Array<{ |
| 20 | msgId: string; |
| 21 | }>; |
| 22 | }; |
| 23 | |
| 24 | function createQueueMetrics(): MessageBatchMetrics { |
| 25 | return { |
| 26 | backlogBytes: 0, |
| 27 | backlogCount: 0, |
| 28 | }; |
| 29 | } |
| 30 | |
| 31 | function createMessageBatchMetadata(): MessageBatchMetadata { |
| 32 | return { |
| 33 | metrics: createQueueMetrics(), |
| 34 | }; |
| 35 | } |
| 36 | |
| 37 | /** |
| 38 | * Wrangler's Queue test/runtime types include delivery metadata on both |
| 39 | * batches and send responses. Tests that stub queue behavior do not care |
| 40 | * about those live backlog values, but they still need to provide the same |
| 41 | * shape so mocked bindings stay honest with the Worker API. |
| 42 | */ |
| 43 | export function createQueueSendResponse(): QueueSendResponse { |
| 44 | return { |
| 45 | metadata: createMessageBatchMetadata(), |
| 46 | }; |
| 47 | } |
| 48 | |
| 49 | /** |
| 50 | * Run a single repo task queue message through the handler and return |
| 51 | * whether it was acked or retried. Uses the real test `env` by default; |
| 52 | * pass `overrideEnv` for tests that stub bindings. The Cloudflare Queue |
| 53 | * helpers own ack/retry tracking and wait for queue `ctx.waitUntil()` work. |
| 54 | */ |
| 55 | export async function runQueueMessage(body: unknown, overrideEnv?: Env): Promise<QueueRunResult> { |
| 56 | const messageId = "queue-1"; |
| 57 | const messages = [ |
| 58 | { |
| 59 | id: messageId, |
| 60 | timestamp: new Date(), |
| 61 | attempts: 1, |
| 62 | body, |
| 63 | }, |
| 64 | ]; |
| 65 | const batch = createMessageBatch("git-on-cloudflare-repo-maint", messages); |
| 66 | |
| 67 | const ctx = createExecutionContext(); |
| 68 | await worker.queue(batch, overrideEnv ?? testEnv, ctx); |
| 69 | const result: QueueResultState = await getQueueResult(batch, ctx); |
| 70 | |
| 71 | return { |
| 72 | acked: result.ackAll || result.explicitAcks.includes(messageId), |
| 73 | retried: |
| 74 | result.retryBatch.retry || |
| 75 | result.retryMessages.some((message) => message.msgId === messageId), |
| 76 | }; |
| 77 | } |