import { WorkerEntrypoint as CloudflareWorkerEntrypoint } from "cloudflare:workers"; import { Cause, ConfigProvider, Context, Effect, Layer, ManagedRuntime, type Scope } from "effect"; import type { Schema as S } from "effect"; import { HttpServerRequest, HttpServerResponse } from "effect/unstable/http"; import type * as Binding from "./Binding"; import { WorkerConfig, WorkerEnvironment, type WorkerEnv } from "./Environment"; import { fromMessage, fromMessageBatch, type QueueHandler } from "./Queue"; import type * as Rpc from "./Rpc"; import type * as RpcDefinition from "./RpcDefinition"; import type * as ServiceBinding from "./ServiceBinding"; import * as WorkerDefinition from "./WorkerDefinition"; import * as CloudflareClock from "./internal/Clock"; import * as Entrypoint from "./internal/Entrypoint"; import { fromExecutionContext, type RunWaitUntilEffect } from "./internal/WorkerContext"; export class ExecutionContext extends Context.Service< ExecutionContext, globalThis.ExecutionContext >()("effect-cf/ExecutionContext") {} export interface WorkerContextWaitUntilOptions { readonly mode?: "observe" | "propagate"; readonly onFailure?: (cause: Cause.Cause) => Effect.Effect; } export interface WorkerContextService { readonly raw: globalThis.ExecutionContext; waitUntil( effect: Effect.Effect, options?: WorkerContextWaitUntilOptions, ): Effect.Effect; waitUntilPropagating( effect: Effect.Effect, options?: Omit, "mode">, ): Effect.Effect; readonly passThroughOnException: Effect.Effect; } export class WorkerContext extends Context.Service()( "effect-cf/WorkerContext", ) {} export class NativeRequest extends Context.Service()( "effect-cf/NativeRequest", ) {} export const isWebSocketUpgrade = (request: Request): boolean => request.headers.get("Upgrade")?.toLowerCase() === "websocket"; export type ReservedMethodName = | RpcDefinition.ReservedMethodName | "fetch" | "connect" | "queue" | "scheduled" | "tail" | "tailStream" | "test" | "trace" | "alarm" | "webSocketMessage" | "webSocketClose" | "webSocketError"; const reservedMethodNames = new Set([ "constructor", "dup", "fetch", "connect", "queue", "scheduled", "tail", "tailStream", "test", "trace", "alarm", "webSocketMessage", "webSocketClose", "webSocketError", ]); type WorkerBaseContext = ExecutionContext | WorkerContext | WorkerEnvironment | ROut; type WorkerFetchContext = | WorkerBaseContext | NativeRequest | HttpServerRequest.HttpServerRequest | Scope.Scope; type WorkerRpcContext = WorkerBaseContext | Scope.Scope; type RuntimeContext = WorkerBaseContext; type RunOptions = { readonly eventLayer?: boolean; }; const RunSymbol = Symbol.for("effect-cf/Worker/run"); export type WorkerFetchSuccess = Response | HttpServerResponse.HttpServerResponse; export type WorkerHandler = Effect.Effect< A, unknown, WorkerFetchContext >; export type WorkerRpcHandler = Effect.Effect>; export type WorkerRpc = Record) => WorkerRpcHandler>; export type WorkerRpcShape, ROut> = { readonly [Key in keyof Rpc]: Rpc[Key] extends ( ...args: infer Args ) => Effect.Effect> ? (...args: Args) => Promise : never; }; export type RpcHandlers = { readonly [Key in keyof Api as Key extends keyof CloudflareWorkerEntrypoint ? never : Key extends string ? Key extends ReservedMethodName ? never : [Api[Key]] extends [never] ? never : Api[Key] extends (...args: Array) => Promise ? Key : never : never]: Api[Key] extends (...args: infer Args) => Promise ? (...args: Args) => WorkerRpcHandler : never; }; export interface WorkerOptions< RRuntime, REvent = never, EventLayerError = never, Rpc extends WorkerRpc = Record, > { /** * Layer provided around each Cloudflare event handled by this Worker. * * The layer is built inside the event's Effect scope and finalized when the * event effect completes. Use this for event-scoped resources such as OTLP * trace/log exporters that should flush at event completion. */ readonly eventLayer?: Layer.Layer>; readonly fetch?: Effect.Effect< WorkerFetchSuccess, unknown, WorkerFetchContext >; readonly queue?: QueueHandler; readonly rpc?: Rpc; } export type FetchWorkerOptions = Omit< WorkerOptions>, "rpc" > & { readonly rpc?: never; }; export type WorkerClass, ROut> = new ( ctx: globalThis.ExecutionContext, env: WorkerEnv, ) => CloudflareWorkerEntrypoint & { fetch(request: Request): Promise; queue(batch: globalThis.MessageBatch): Promise; } & WorkerRpcShape; export interface FetchHandler { readonly fetch: ( request: Request, env: Env, ctx: globalThis.ExecutionContext, ) => Promise; } export const renderHttpResponse = ( effect: Effect.Effect, ): Effect.Effect => Effect.flatMap(effect, (response) => Effect.map(Effect.context(), (context) => HttpServerResponse.toWeb(response, { context }), ), ); const renderFetchSuccess = ( effect: Effect.Effect, ): Effect.Effect => Effect.flatMap(effect, (response) => response instanceof Response ? Effect.succeed(response) : Effect.map(Effect.context(), (context) => HttpServerResponse.toWeb(response, { context }), ), ); const isWorkerOptions = >( options: WorkerOptions | WorkerHandler, ): options is WorkerOptions => typeof options === "object" && options !== null && ("eventLayer" in options || "fetch" in options || "queue" in options || "rpc" in options); export function make( layer: Layer.Layer, fetch: WorkerHandler, ): WorkerClass, ROut>; export function make< ROut, LayerError, REvent = never, EventLayerError = never, const Rpc extends WorkerRpc = Record, >( layer: Layer.Layer, options: WorkerOptions, ): WorkerClass; export function make< ROut, LayerError, REvent = never, EventLayerError = never, const Rpc extends WorkerRpc = Record, >( layer: Layer.Layer, optionsOrFetch: WorkerOptions | WorkerHandler, ): WorkerClass { const options = isWorkerOptions(optionsOrFetch) ? optionsOrFetch : ({ fetch: optionsOrFetch } as WorkerOptions); class EffectWorker extends CloudflareWorkerEntrypoint { readonly runtime: ManagedRuntime.ManagedRuntime, LayerError>; constructor(ctx: globalThis.ExecutionContext, env: WorkerEnv) { super(ctx, env); let runWaitUntilEffect: RunWaitUntilEffect = () => Promise.reject(new Error("WorkerContext runtime is not initialized")); const services = Layer.mergeAll( CloudflareClock.layer, Layer.succeed(ExecutionContext, ctx), ConfigProvider.layer(Effect.succeed(WorkerConfig.providerFromEnv(env))), Layer.succeed( WorkerContext, fromExecutionContext(ctx, (effect) => runWaitUntilEffect(effect)), ), Layer.succeed(WorkerEnvironment, env), ) as Layer.Layer; const runtimeLayer = Entrypoint.provideEntrypointServices< ROut, LayerError, ExecutionContext | WorkerContext | WorkerEnvironment >(layer, services); this.runtime = ManagedRuntime.make(runtimeLayer); runWaitUntilEffect = (effect: Effect.Effect) => this.runtime.runPromiseExit(effect as Effect.Effect>); } [RunSymbol]( effect: Effect.Effect | REvent | Scope.Scope>, runOptions: RunOptions = {}, ): Promise { const effectWithEventLayer = runOptions.eventLayer === false || options.eventLayer === undefined ? effect : effect.pipe(Effect.provide(options.eventLayer, { local: true })); return this.runtime.runPromise( Effect.scoped( effectWithEventLayer as Effect.Effect< A, E | EventLayerError, RuntimeContext | Scope.Scope >, ), ); } fetch(request: Request): Promise { const fetchHandler = options.fetch; if (fetchHandler === undefined) { return Promise.resolve(new Response("Not Found", { status: 404 })); } const requestServices = Layer.mergeAll( Layer.succeed(NativeRequest, request), Layer.succeed(HttpServerRequest.HttpServerRequest, HttpServerRequest.fromWeb(request)), ); return this[RunSymbol]( renderFetchSuccess(fetchHandler).pipe(Effect.provide(requestServices)) as Effect.Effect< Response, unknown, RuntimeContext | REvent | Scope.Scope >, ); } queue(batch: globalThis.MessageBatch): Promise { const queueHandler = options.queue; if (queueHandler === undefined) { return Promise.resolve(); } const messages = batch.messages.map((message) => fromMessage(message, message.body)); return this[RunSymbol](queueHandler(fromMessageBatch(batch, messages))); } } Entrypoint.defineEntrypointRpcMethods( "Worker", EffectWorker.prototype, options.rpc, reservedMethodNames, (self, effect) => self[RunSymbol](effect), ); return Entrypoint.assumeEntrypointClass>(EffectWorker); } export const makeFetchHandler = < ROut, LayerError, REvent = never, EventLayerError = never, Env extends WorkerEnv = WorkerEnv, >( layer: Layer.Layer, options: FetchWorkerOptions, ): FetchHandler => { const WorkerClass = make(layer, options); return { fetch: (request, env, ctx) => Promise.resolve(new WorkerClass(ctx, env).fetch(request)), }; }; export type ServiceFreeSchema = S.Codec; export interface Method< Args extends ReadonlyArray = ReadonlyArray, Success extends ServiceFreeSchema = ServiceFreeSchema, > { readonly args: Args; readonly success: Success; } export namespace Method { export type Any = Method, ServiceFreeSchema>; type ArgsFromSchemas> = Args extends readonly [] ? [] : Args extends readonly [ infer Head extends ServiceFreeSchema, ...infer Tail extends ReadonlyArray, ] ? [S.Schema.Type, ...ArgsFromSchemas] : Array>; export type Args = ArgsFromSchemas; export type Success = S.Schema.Type; } export type Methods = Record; export type NoReservedMethods = Extract extends never ? MethodsShape : never; export interface Definition { readonly id: Id; readonly methods: MethodsShape; } export namespace Definition { export type Any = Definition; } export type ServerApi = { readonly [Key in keyof Self["methods"]]: ( ...args: Method.Args ) => Promise>; }; export type Api = Rpc.Provider, ReservedMethodName>; export type Handlers = { readonly [Key in keyof Self["methods"]]: ( ...args: Method.Args ) => WorkerRpcHandler>; }; export interface Options< ROut, Self extends Definition.Any, REvent = never, EventLayerError = never, > extends Omit>, "rpc"> { readonly rpc: Handlers; } export type LayerOptions = WorkerDefinition.LayerOptions; export type TagClass = Context.ServiceClass< Self, Id, ServiceBinding.ServiceBindingEffectClient< Api>, Definition > > & ServiceBinding.ServiceBindingStaticClient< Self, Api>, Definition > & { readonly id: Id; readonly methods: MethodsShape; readonly make: ( layer: Layer.Layer, options: Options, REvent, EventLayerError>, ) => WorkerClass>, ROut | REvent>; readonly layer: ( options: LayerOptions, ) => Layer.Layer< Self, Binding.BindingNotFoundError | Binding.BindingValidationError, WorkerEnvironment >; }; export type TagFactory = () => ( id: Id, methods: MethodsShape & NoReservedMethods, ) => TagClass; export const Tag = WorkerDefinition.Tag as unknown as TagFactory; export const method = WorkerDefinition.method as { (definition: { readonly success: Success; }): Method; < const Args extends ReadonlyArray, Success extends ServiceFreeSchema, >(definition: { readonly args: Args; readonly success: Success; }): Method; }; export const implement = WorkerDefinition.implement as unknown as < ROut, const Self extends Definition.Any, >( _definition: Self, handlers: Handlers, ) => Handlers; export type HandlerEffect< ROut, Self extends Definition.Any, Key extends keyof Self["methods"], > = WorkerRpcHandler>;