import { Context, Effect, Schema as S, type Layer } from "effect"; import * as Binding from "./Binding"; import type { WorkerEnvironment } from "./Environment"; import type { ExecutionContext, WorkerContext } from "./Worker"; import type * as RpcDefinition from "./RpcDefinition"; import * as WorkflowBinding from "./WorkflowBinding"; import * as WorkflowEntrypoint from "./Workflow"; export interface Definition< Id extends string = string, Payload extends RpcDefinition.ServiceFreeSchema = RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema = RpcDefinition.ServiceFreeSchema, > { readonly id: Id; readonly payload: Payload; readonly result: Result; } export namespace Definition { export type Any = Definition< string, RpcDefinition.ServiceFreeSchema, RpcDefinition.ServiceFreeSchema >; } export type Handler = ( payload: S.Schema.Type, ) => Effect.Effect< S.Schema.Type, unknown, WorkflowEntrypoint.WorkflowRunContext >; export interface Options { readonly run: Handler; } export type LayerOptions = { readonly binding: string; }; export interface TagClass< Self, Id extends string, Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, > extends Context.ServiceClass> { readonly id: Id; readonly payload: Payload; readonly result: Result; readonly make: ( layer: Layer.Layer, options: Options>, ) => WorkflowEntrypoint.WorkflowClass, S.Codec.Encoded, ROut>; readonly layer: ( options: LayerOptions, ) => Layer.Layer< Self, Binding.BindingNotFoundError | Binding.BindingValidationError, WorkerEnvironment >; readonly create: ( payload: S.Schema.Type, options?: WorkflowBinding.WorkflowInstanceCreateOptions>, ) => Effect.Effect< WorkflowBinding.WorkflowInstance>, WorkflowBinding.WorkflowOperationError | S.SchemaError, Self >; readonly createBatch: ( batch: WorkflowBinding.WorkflowInstanceCreateBatchOptions< S.Schema.Type, S.Codec.Encoded >, ) => Effect.Effect< ReadonlyArray>>, WorkflowBinding.WorkflowOperationError | S.SchemaError, Self >; readonly get: ( instanceId: string, ) => Effect.Effect< WorkflowBinding.WorkflowInstance>, WorkflowBinding.WorkflowOperationError, Self >; readonly unsafeRaw: () => Effect.Effect< globalThis.Workflow>, never, Self >; } const makeDefinition = < Id extends string, Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, >( id: Id, definition: { readonly payload: Payload; readonly result: Result; }, ) => { type SelfDefinition = Definition; const workflowDefinition: SelfDefinition = { id, payload: definition.payload, result: definition.result, }; return Object.assign(workflowDefinition, { make: ( layer: Layer.Layer, options: Options, ) => WorkflowEntrypoint.make(layer, { run: wrapHandler(workflowDefinition, options.run), }), }); }; export const make = < Id extends string, Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, >( id: Id, definition: { readonly payload: Payload; readonly result: Result; }, ) => Tag>()(id, definition); export const Tag = () => < Id extends string, Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, >( id: Id, definition: { readonly payload: Payload; readonly result: Result; }, ) => { const workflowDefinition = makeDefinition(id, definition); const tag = Context.Service>()(id); const layer = (binding: LayerOptions) => WorkflowBinding.layer(tag, { ...binding, payload: definition.payload, result: definition.result, }); const create = Effect.fnUntraced(function* ( payload: S.Schema.Type, options?: WorkflowBinding.WorkflowInstanceCreateOptions>, ) { const workflow = yield* tag; return yield* workflow.create(payload, options); }); const createBatch = Effect.fnUntraced(function* ( batch: WorkflowBinding.WorkflowInstanceCreateBatchOptions< S.Schema.Type, S.Codec.Encoded >, ) { const workflow = yield* tag; return yield* workflow.createBatch(batch); }); const get = Effect.fnUntraced(function* (instanceId: string) { const workflow = yield* tag; return yield* workflow.get(instanceId); }); const unsafeRaw = Effect.fnUntraced(function* () { const workflow = yield* tag; return yield* workflow.unsafeRaw; }); return Object.assign(tag, { id: workflowDefinition.id, payload: workflowDefinition.payload, result: workflowDefinition.result, make: workflowDefinition.make, layer, create, createBatch, get, unsafeRaw, }) as TagClass; }; export const Workflow = Tag; const wrapHandler = ( definition: Self, handler: Handler, ): WorkflowEntrypoint.WorkflowHandler< ROut, S.Codec.Encoded, S.Codec.Encoded > => { const decodePayload = S.decodeUnknownEffect(definition.payload); const encodeResult = S.encodeEffect(definition.result); return (payload) => Effect.gen(function* () { const decodedPayload = yield* decodePayload(payload); const event = yield* WorkflowEntrypoint.WorkflowEvent; const decodedEvent = { ...event, payload: decodedPayload, } as WorkflowEntrypoint.WorkflowEventService>; const result = yield* handler(decodedPayload as S.Schema.Type).pipe( Effect.provideService(WorkflowEntrypoint.WorkflowEvent, decodedEvent), ); return yield* encodeResult(result as S.Schema.Type); }); }; export const implement = ( _definition: Self, handler: Handler, ): Handler => handler;