import { Data, Effect, type Scope } from "effect"; import type { ExecutionContext, WorkerContext } from "./Worker"; import type { WorkerEnvironment } from "./Environment"; import * as QueueDefinition from "./QueueDefinition"; export interface QueueMessage { readonly raw: globalThis.Message; readonly id: string; readonly timestamp: Date; readonly body: Body; readonly attempts: number; readonly ack: Effect.Effect; readonly retry: (options?: globalThis.QueueRetryOptions) => Effect.Effect; } export interface QueueBatch { readonly raw: globalThis.MessageBatch; readonly messages: ReadonlyArray>; readonly queue: string; readonly metadata: globalThis.MessageBatchMetadata; readonly ackAll: Effect.Effect; readonly retryAll: (options?: globalThis.QueueRetryOptions) => Effect.Effect; } export class QueueMessageDecodeError extends Data.TaggedError("QueueMessageDecodeError")<{ readonly queue: string; readonly messageId: string; readonly index: number; readonly cause: unknown; }> {} type RuntimeContext = ExecutionContext | WorkerContext | WorkerEnvironment | ROut; type QueueHandlerContext = RuntimeContext | Scope.Scope; export type QueueHandler = ( batch: QueueBatch, ) => Effect.Effect>; export interface QueueOptions { readonly queue: QueueHandler; } export const fromMessage = (message: globalThis.Message, body: Body) => ({ raw: message, id: message.id, timestamp: message.timestamp, body, attempts: message.attempts, ack: Effect.sync(() => message.ack()), retry: (options?: globalThis.QueueRetryOptions) => Effect.sync(() => message.retry(options)), }); export const fromMessageBatch = ( batch: globalThis.MessageBatch, messages: ReadonlyArray>, ): QueueBatch => ({ raw: batch, messages, queue: batch.queue, metadata: batch.metadata, ackAll: Effect.sync(() => batch.ackAll()), retryAll: (options?: globalThis.QueueRetryOptions) => Effect.sync(() => batch.retryAll(options)), }); export const decodeBatch = ( batch: globalThis.MessageBatch, decodeBody: (body: unknown) => Effect.Effect, ): Effect.Effect, QueueMessageDecodeError> => Effect.gen(function* () { const messages: Array> = []; for (let index = 0; index < batch.messages.length; index++) { const message = batch.messages[index]; const body = yield* decodeBody(message.body).pipe( Effect.mapError( (cause) => new QueueMessageDecodeError({ queue: batch.queue, messageId: message.id, index, cause, }), ), ); messages.push(fromMessage(message, body)); } return fromMessageBatch(batch, messages); }); export type Definition< Id extends string = string, Message extends QueueDefinition.Definition.Any["message"] = QueueDefinition.Definition.Any["message"], > = QueueDefinition.Definition; export namespace Definition { export type Any = QueueDefinition.Definition.Any; } export type LayerOptions = QueueDefinition.LayerOptions; export type TagClass< Self, Id extends string, Message extends QueueDefinition.Definition.Any["message"], > = QueueDefinition.TagClass; export const Tag: () => < Id extends string, Message extends QueueDefinition.Definition.Any["message"], >( id: Id, definition: { readonly message: Message }, ) => TagClass = QueueDefinition.Tag; export const implement = QueueDefinition.implement; export type Handler = QueueDefinition.Handler< ROut, Self >; export type Options = QueueDefinition.Options< ROut, Self >;