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 },
);