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
>;