import { Context, Effect, Schema as S, type Layer, type Scope } from "effect"; import * as Binding from "./Binding"; import type { WorkerEnvironment } from "./Environment"; import * as QueueEntrypoint from "./Queue"; import * as QueueBinding from "./QueueBinding"; import type { ExecutionContext, WorkerContext } from "./Worker"; import * as WorkerEntrypoint from "./Worker"; import type * as RpcDefinition from "./RpcDefinition"; export interface Definition< Id extends string = string, Message extends RpcDefinition.ServiceFreeSchema = RpcDefinition.ServiceFreeSchema, > { readonly id: Id; readonly message: Message; } export namespace Definition { export type Any = Definition; } export type Handler = ( batch: QueueEntrypoint.QueueBatch>, ) => Effect.Effect< void, unknown, ExecutionContext | WorkerEntrypoint.WorkerContext | WorkerEnvironment | Scope.Scope | ROut >; export interface Options extends Omit< WorkerEntrypoint.WorkerOptions>, "queue" | "rpc" > { readonly queue: Handler; readonly rpc?: never; } export type LayerOptions = { readonly binding: string; }; export interface TagClass< Self, Id extends string, Message extends RpcDefinition.ServiceFreeSchema, > extends Context.ServiceClass> { readonly id: Id; readonly message: Message; readonly make: ( layer: Layer.Layer, options: Options>, ) => WorkerEntrypoint.WorkerClass, ROut>; readonly layer: ( options: LayerOptions, ) => Layer.Layer< Self, Binding.BindingNotFoundError | Binding.BindingValidationError, WorkerEnvironment >; readonly send: ( message: S.Schema.Type, options?: QueueBinding.QueueSendOptions, ) => Effect.Effect; readonly sendBatch: ( messages: Iterable>>, options?: QueueBinding.QueueSendBatchOptions, ) => Effect.Effect; readonly metrics: () => Effect.Effect< QueueBinding.QueueMetrics, QueueBinding.QueueOperationError, Self >; readonly unsafeRaw: () => Effect.Effect>, never, Self>; } const makeDefinition = ( id: Id, definition: { readonly message: Message }, ) => { type SelfDefinition = Definition; const queueDefinition: SelfDefinition = { id, message: definition.message, }; return Object.assign(queueDefinition, { make: ( layer: Layer.Layer, options: Options, ) => WorkerEntrypoint.make(layer, { ...options, queue: wrapHandler(queueDefinition, options.queue), }), }); }; export const make = ( id: Id, definition: { readonly message: Message }, ) => Tag>()(id, definition); export const Tag = () => ( id: Id, definition: { readonly message: Message }, ) => { const queueDefinition = makeDefinition(id, definition); const tag = Context.Service>()(id); const layer = (binding: LayerOptions) => QueueBinding.layer(tag, { ...binding, message: definition.message, }); const send = Effect.fnUntraced(function* ( message: S.Schema.Type, options?: QueueBinding.QueueSendOptions, ) { const queue = yield* tag; yield* queue.send(message, options); }); const sendBatch = Effect.fnUntraced(function* ( messages: Iterable>>, options?: QueueBinding.QueueSendBatchOptions, ) { const queue = yield* tag; yield* queue.sendBatch(messages, options); }); const metrics = Effect.fnUntraced(function* () { const queue = yield* tag; return yield* queue.metrics(); }); const unsafeRaw = Effect.fnUntraced(function* () { const queue = yield* tag; return yield* queue.unsafeRaw; }); return Object.assign(tag, { id: queueDefinition.id, message: queueDefinition.message, make: queueDefinition.make, layer, send, sendBatch, metrics, unsafeRaw, }) as TagClass; }; export const Queue = Tag; const wrapHandler = ( definition: Self, handler: Handler, ): QueueEntrypoint.QueueHandler => { const decodeBody = S.decodeUnknownEffect(definition.message); return (batch) => Effect.gen(function* () { const decoded = yield* QueueEntrypoint.decodeBatch(batch.raw, decodeBody); yield* handler(decoded); }); }; export const implement = ( _definition: Self, handler: Handler, ): Handler => handler;