import { Context, Effect, type Layer } from "effect"; import type { Schema as S } from "effect"; import * as Binding from "./Binding"; import type { WorkerEnvironment } from "./Environment"; import * as WorkerEntrypoint from "./Worker"; import type { WorkerRpcHandler } from "./Worker"; import type * as Rpc from "./Rpc"; import * as RpcDefinition from "./RpcDefinition"; import * as ServiceBinding from "./ServiceBinding"; 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>; type EncodedArgsFromSchemas> = Args extends readonly [] ? [] : Args extends readonly [ infer Head extends ServiceFreeSchema, ...infer Tail extends ReadonlyArray, ] ? [S.Codec.Encoded, ...EncodedArgsFromSchemas] : Array>; export type Args = ArgsFromSchemas; export type EncodedArgs = EncodedArgsFromSchemas; export type Success = S.Schema.Type; export type EncodedSuccess = S.Codec.Encoded; } export type Methods = Record; /** * RPC contract for a Worker service. * * Create with {@link make} and reuse to type both worker implementations and * service bindings in other workers. */ export interface Definition { readonly id: Id; readonly methods: MethodsShape; } export namespace Definition { export type Any = Definition; } export type ReservedMethodName = WorkerEntrypoint.ReservedMethodName; export type NoReservedMethods = Extract extends never ? MethodsShape : never; const reservedMethodNames = new Set([ "constructor", "dup", "fetch", "connect", "queue", "scheduled", "tail", "tailStream", "test", "trace", "alarm", "webSocketMessage", "webSocketClose", "webSocketError", ]); /** * Promise-based client API derived from a {@link Definition}. */ export type ServerApi = { readonly [Key in keyof Self["methods"]]: ( ...args: Method.Args ) => Promise>; }; export type Api = Rpc.Provider, ReservedMethodName>; /** * Effect handlers for each RPC method in a worker definition. */ export type Handlers = { readonly [Key in keyof Self["methods"]]: ( ...args: Method.Args ) => WorkerRpcHandler>; }; type BoundaryHandlers = { readonly [Key in keyof Self["methods"]]: ( ...args: Array ) => WorkerRpcHandler>; }; /** * Worker constructor options for a specific RPC definition. */ export interface Options< ROut, Self extends Definition.Any, REvent = never, EventLayerError = never, > extends Omit< WorkerEntrypoint.WorkerOptions>, "rpc" > { readonly rpc: Handlers; } export type LayerOptions = { readonly binding: string; }; 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< ROut, LayerError, WorkerEntrypoint.ExecutionContext | WorkerEntrypoint.WorkerContext | WorkerEnvironment >, options: Options, REvent, EventLayerError>, ) => WorkerEntrypoint.WorkerClass< Handlers>, ROut | REvent >; readonly layer: ( options: LayerOptions, ) => Layer.Layer< Self, Binding.BindingNotFoundError | Binding.BindingValidationError, WorkerEnvironment >; }; /** * Defines a single RPC method schema in a worker definition. */ export const method = RpcDefinition.method as { (definition: { readonly success: Success; }): Method; < const Args extends ReadonlyArray, Success extends ServiceFreeSchema, >(definition: { readonly args: Args; readonly success: Success; }): Method; }; /** * Creates a typed worker RPC definition plus helpers for implementation and bindings. * * @example * ```ts * const CounterWorker = WorkerDefinition.make("CounterWorker", { * increment: WorkerDefinition.method({ * args: [Schema.Number], * success: Schema.Number, * }), * }); * ``` */ const makeDefinition = ( id: Id, methods: MethodsShape & NoReservedMethods, ) => { type SelfDefinition = Definition; RpcDefinition.assertNoReservedMethods("Worker", methods, reservedMethodNames); const definition: SelfDefinition = RpcDefinition.make(id, methods); return Object.assign(definition, { make: ( layer: Layer.Layer< ROut, LayerError, WorkerEntrypoint.ExecutionContext | WorkerEntrypoint.WorkerContext | WorkerEnvironment >, options: Options, ) => WorkerEntrypoint.make(layer, { ...options, rpc: wrapHandlers(definition, options.rpc), }), }); }; export const make = ( id: Id, methods: MethodsShape & NoReservedMethods, ) => Tag>()( id, methods as MethodsShape & NoReservedMethods, ); export const Tag = () => ( id: Id, methods: MethodsShape & NoReservedMethods, ) => { const definition = makeDefinition(id, methods); type SelfDefinition = Definition; type ClientApi = Api; const tag = Context.Service< Self, ServiceBinding.ServiceBindingEffectClient >()(id); const bindingDefinition = (binding: LayerOptions) => ({ ...binding, definition, }); const layer = (binding: LayerOptions) => ServiceBinding.layer(tag, bindingDefinition(binding)); const fetch = (input: RequestInfo | URL, init?: RequestInit) => Effect.gen(function* () { const service = yield* tag; return yield* service.fetch(input, init); }); const rpc = ( method: Method, ...args: ClientApi[Method] extends (...args: infer Args) => unknown ? Args : never ) => Effect.gen(function* () { const service = yield* tag; return yield* service.rpc(method as never, ...(args as never)); }); const call = ( method: Method, ...args: ClientApi[Method] extends (...args: infer Args) => unknown ? Args : never ) => Effect.gen(function* () { const service = yield* tag; return yield* service.call(method as never, ...(args as never)); }); const scopedCall = ( method: Method, ...args: ClientApi[Method] extends (...args: infer Args) => unknown ? Args : never ) => Effect.gen(function* () { const service = yield* tag; return yield* service.scopedCall(method as never, ...(args as never)); }); const directMethods = ServiceBinding.makeDirectMethods( definition, call as never, ); return Object.assign(tag, directMethods, { id: definition.id, methods: definition.methods, make: definition.make, layer, fetch, rpc, call, scopedCall, }) as unknown as TagClass; }; export const Worker = Tag; const wrapHandlers = ( definition: Self, handlers: Handlers, ): BoundaryHandlers => { const wrapped = {} as Record; for (const key of Object.keys(definition.methods) as Array< RpcDefinition.Definition.MethodNames >) { const handler = handlers[key]; wrapped[key] = (...args: Array) => Effect.gen(function* () { const decodedArgs = yield* RpcDefinition.decodeArgs(definition, key, args); const value = yield* handler(...decodedArgs); return yield* RpcDefinition.encodeSuccess(definition, key, value); }); } return wrapped as BoundaryHandlers; }; /** * Helper for implementing handlers with the exact method shape of a definition. */ export const implement = ( _definition: Self, handlers: Handlers, ): Handlers => handlers; /** * Convenience alias for a single worker RPC handler Effect. */ export type HandlerEffect< ROut, Self extends Definition.Any, Key extends keyof Self["methods"], > = WorkerRpcHandler>;