import { DurableObject as CloudflareDurableObject } from "cloudflare:workers"; import { ConfigProvider, Effect, Layer, ManagedRuntime, type Context, type Scope } from "effect"; import type { Schema as S } from "effect"; import { NativeRequest } from "./Worker"; import { WorkerConfig, WorkerEnvironment, type WorkerEnv } from "./Environment"; import { DurableObjectState, fromDurableObjectState } from "./DurableObjectState"; import { fromWebSocket, type DurableWebSocket } from "./DurableObjectWebSocket"; import type * as Binding from "./Binding"; import * as DurableObjectDefinition from "./DurableObjectDefinition"; import type * as DurableObjectNamespace from "./DurableObjectNamespace"; import type * as Rpc from "./Rpc"; import * as CloudflareClock from "./internal/Clock"; import * as Entrypoint from "./internal/Entrypoint"; const reservedMethodNames = new Set([ "constructor", "dup", "fetch", "alarm", "webSocketMessage", "webSocketClose", "webSocketError", ]); type RuntimeContext = DurableObjectState | WorkerEnvironment | ROut; type HandlerContext = RuntimeContext | Scope.Scope; type FetchContext = HandlerContext | NativeRequest; type RunOptions = { readonly eventLayer?: boolean; }; const RunSymbol = Symbol.for("effect-cf/DurableObject/run"); /** * Effect type for Durable Object lifecycle and RPC handlers. */ export type DurableObjectHandler = Effect.Effect< A, unknown, HandlerContext >; /** * Shape of Durable Object RPC handlers passed to {@link make}. */ export type DurableObjectRpc = Record< string, (...args: Array) => DurableObjectHandler >; export type DurableObjectRpcShape, 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 CloudflareDurableObject ? never : Key extends string ? [Api[Key]] extends [never] ? never : Api[Key] extends (...args: Array) => Promise ? Key : never : never]: Api[Key] extends (...args: infer Args) => Promise ? (...args: Args) => DurableObjectHandler : never; }; /** * Options for creating a Durable Object class backed by Effect handlers. */ export interface DurableObjectOptions< RRuntime, REvent = never, EventLayerError = never, Rpc extends DurableObjectRpc = Record, > { /** * Layer provided around each Cloudflare event handled by this Durable Object. * * The layer is built inside the event's Effect scope and finalized when the * event effect completes. It is not applied to `initialize`, which is an * instance-load lifecycle hook rather than a platform event. */ readonly eventLayer?: Layer.Layer< REvent, EventLayerError, DurableObjectState | WorkerEnvironment | RRuntime >; /** * Effect run when Cloudflare loads this Durable Object instance into memory. * * Use `DurableObjectState.blockConcurrencyWhile` inside this hook when * incoming events should wait for setup to finish. Cloudflare may construct * the same Durable Object id again after eviction or restart; use Durable * Object storage if work must happen only once per id. */ readonly initialize?: Effect.Effect>; /** Optional RPC methods exposed as Durable Object instance methods. */ readonly rpc?: Rpc; /** Optional fetch handler for HTTP/WebSocket requests. */ readonly fetch?: Effect.Effect>; /** * Optional logical alarm processing effect. * * This runs before `alarm` and should be built with helpers such as * `DurableObjectAlarm.processDue(...)` so the reusable scheduler stays inside * the Durable Object's single managed runtime boundary. */ readonly alarms?: Effect.Effect>; readonly alarm?: ( alarmInfo?: globalThis.AlarmInvocationInfo, ) => Effect.Effect>; readonly webSocketMessage?: ( socket: DurableWebSocket, message: string | ArrayBuffer, ) => Effect.Effect>; readonly webSocketClose?: ( socket: DurableWebSocket, code: number, reason: string, wasClean: boolean, ) => Effect.Effect>; readonly webSocketError?: ( socket: DurableWebSocket, error: unknown, ) => Effect.Effect>; } /** * Cloudflare `DurableObject` constructor produced by {@link make}. */ export type DurableObjectClass, ROut> = new ( state: globalThis.DurableObjectState, env: WorkerEnv, ) => CloudflareDurableObject & DurableObjectRpcShape; /** * Creates a Durable Object class backed by a single managed Effect runtime. */ export const make = < ROut, LayerError, REvent = never, EventLayerError = never, const Rpc extends DurableObjectRpc = Record, >( layer: Layer.Layer, options: DurableObjectOptions = {}, ): DurableObjectClass => { class EffectDurableObject extends CloudflareDurableObject { readonly runtime: ManagedRuntime.ManagedRuntime, LayerError>; constructor(state: globalThis.DurableObjectState, env: WorkerEnv) { super(state, env); const services = Layer.mergeAll( CloudflareClock.layer, ConfigProvider.layer(Effect.succeed(WorkerConfig.providerFromEnv(env))), Layer.succeed(DurableObjectState, fromDurableObjectState(state)), Layer.succeed(WorkerEnvironment, env), ); const runtimeLayer = Entrypoint.provideEntrypointServices(layer, services); this.runtime = ManagedRuntime.make(runtimeLayer); const initialize = options.initialize; if (initialize !== undefined) { state.waitUntil(this[RunSymbol](initialize, { eventLayer: false })); } } [RunSymbol]( effect: Effect.Effect>, 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>, ), ); } fetch(request: Request): Promise { const fetchHandler = options.fetch; if (fetchHandler === undefined) { return Promise.resolve(new Response("Not Found", { status: 404 })); } return this[RunSymbol](Effect.provideService(fetchHandler, NativeRequest, request)); } alarm(alarmInfo?: globalThis.AlarmInvocationInfo): Promise | void { const logicalAlarms = options.alarms?.pipe(Effect.asVoid); const rawAlarm = options.alarm?.(alarmInfo); if (logicalAlarms !== undefined && rawAlarm !== undefined) { return this[RunSymbol]( Effect.gen(function* () { yield* logicalAlarms; yield* rawAlarm; }), ); } if (logicalAlarms !== undefined) { return this[RunSymbol](logicalAlarms); } if (rawAlarm !== undefined) { return this[RunSymbol](rawAlarm); } } webSocketMessage(socket: WebSocket, message: string | ArrayBuffer): Promise | void { if (options.webSocketMessage !== undefined) { return this[RunSymbol](options.webSocketMessage(fromWebSocket(socket), message)); } } webSocketClose( socket: WebSocket, code: number, reason: string, wasClean: boolean, ): Promise | void { if (options.webSocketClose !== undefined) { return this[RunSymbol]( options.webSocketClose(fromWebSocket(socket), code, reason, wasClean), ); } } webSocketError(socket: WebSocket, error: unknown): Promise | void { if (options.webSocketError !== undefined) { return this[RunSymbol](options.webSocketError(fromWebSocket(socket), error)); } } } Entrypoint.defineEntrypointRpcMethods( "Durable Object", EffectDurableObject.prototype, options.rpc, reservedMethodNames, (self, effect) => self[RunSymbol](effect), ); return Entrypoint.assumeEntrypointClass>( EffectDurableObject, ); }; 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 ReservedMethodName = DurableObjectDefinition.ReservedMethodName; 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 LayerOptions = DurableObjectDefinition.LayerOptions; 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 ) => DurableObjectHandler>; }; export interface Options< ROut, Self extends Definition.Any, REvent = never, EventLayerError = never, > extends Omit< DurableObjectOptions>, "rpc" > { readonly rpc: Handlers; } export type TagClass = Context.ServiceClass< Self, Id, DurableObjectNamespace.DurableObjectNamespaceEffectClient< Api>, Definition > > & DurableObjectNamespace.DurableObjectNamespaceStaticClient< Self, Api>, Definition > & { readonly id: Id; readonly methods: MethodsShape; readonly make: ( layer: Layer.Layer, options: Options, REvent, EventLayerError>, ) => DurableObjectClass>, 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 = DurableObjectDefinition.Tag as unknown as TagFactory; export const method = DurableObjectDefinition.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 = DurableObjectDefinition.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"], > = DurableObjectHandler>;