import type { MessageSendRequest as CloudflareMessageSendRequest, Queue as CloudflareQueue, QueueMetrics as CloudflareQueueMetrics, QueueSendBatchOptions as CloudflareQueueSendBatchOptions, QueueSendBatchResponse as CloudflareQueueSendBatchResponse, QueueSendOptions as CloudflareQueueSendOptions, QueueSendResponse as CloudflareQueueSendResponse, } from "@cloudflare/workers-types"; import { Context, Data, Effect, Schema as S } from "effect"; import * as Binding from "./Binding"; import type * as RpcDefinition from "./RpcDefinition"; export type QueueSendOptions = CloudflareQueueSendOptions; export type QueueSendResponse = CloudflareQueueSendResponse; export type QueueSendBatchOptions = CloudflareQueueSendBatchOptions; export type QueueSendBatchResponse = CloudflareQueueSendBatchResponse; export type QueueMetrics = CloudflareQueueMetrics; export type MessageSendRequest = CloudflareMessageSendRequest; export type QueueProducer = Pick, "send"> & Partial, "sendBatch" | "metrics">>; const expectedQueueProducer = "Queue producer binding with send(); optional sendBatch()/metrics()"; export interface QueueBindingDefinition { /** Binding name as configured in `wrangler.jsonc`. */ readonly binding: string; /** Codec used to encode messages before sending them to Cloudflare Queues. */ readonly message: Message; } export interface QueueBindingClient { readonly send: ( message: S.Schema.Type, options?: QueueSendOptions, ) => Effect.Effect; readonly sendBatch: ( messages: Iterable>>, options?: QueueSendBatchOptions, ) => Effect.Effect; readonly metrics: () => Effect.Effect; readonly unsafeRaw: Effect.Effect>>; } export class QueueOperationError extends Data.TaggedError("QueueOperationError")<{ readonly binding: string; readonly operation: string; readonly cause: unknown; readonly message: string; }> {} const tryQueuePromise = ( binding: string, operation: string, evaluate: () => Promise, ): Effect.Effect => Effect.tryPromise({ try: evaluate, catch: (cause) => new QueueOperationError({ binding, operation, cause, message: `Cloudflare queue binding "${binding}" operation "${operation}" failed`, }), }); export const isQueue = (value: unknown): value is QueueProducer => { if (typeof value !== "object" || value === null) { return false; } const resource = value as Record; return typeof resource.send === "function"; }; export const makeClient = ( definition: QueueBindingDefinition, ): ((queue: QueueProducer>) => QueueBindingClient) => { type Body = S.Schema.Type; type EncodedBody = S.Codec.Encoded; const encodeMessage = S.encodeEffect(definition.message); return (queue) => ({ send: Effect.fnUntraced(function* (message: Body, options?: QueueSendOptions) { const encoded = yield* encodeMessage(message); yield* tryQueuePromise(definition.binding, "send", () => queue.send(encoded, options)); }), sendBatch: Effect.fnUntraced(function* ( messages: Iterable>, options?: QueueSendBatchOptions, ) { const encodedMessages: Array> = []; for (const message of messages) { encodedMessages.push({ ...message, body: yield* encodeMessage(message.body), }); } const sendBatch = queue.sendBatch; if (sendBatch !== undefined) { yield* tryQueuePromise(definition.binding, "sendBatch", () => sendBatch.call(queue, encodedMessages, options), ); return; } yield* tryQueuePromise(definition.binding, "sendBatch", async () => { for (const message of encodedMessages) { await queue.send(message.body, { contentType: message.contentType, delaySeconds: message.delaySeconds ?? options?.delaySeconds, }); } }); }), metrics: () => { const metrics = queue.metrics; if (metrics === undefined) { return Effect.fail( new QueueOperationError({ binding: definition.binding, operation: "metrics", cause: new Error(`Queue binding "${definition.binding}" does not provide metrics()`), message: `Cloudflare queue binding "${definition.binding}" does not provide metrics()`, }), ); } return tryQueuePromise(definition.binding, "metrics", () => metrics.call(queue)); }, unsafeRaw: Effect.succeed(queue as CloudflareQueue), }); }; export const layer = ( tag: Context.Service>, definition: QueueBindingDefinition, ) => Binding.layer( tag, definition.binding, (value): value is QueueProducer> => isQueue>(value), makeClient(definition), { expected: expectedQueueProducer }, );