import type { InstanceStatus as CloudflareInstanceStatus, Workflow as CloudflareWorkflow, WorkflowInstance as CloudflareWorkflowInstance, WorkflowInstanceCreateOptions as CloudflareWorkflowInstanceCreateOptions, } from "@cloudflare/workers-types"; import { Context, Data, Effect, Option, Schema as S } from "effect"; import * as Binding from "./Binding"; import type * as RpcDefinition from "./RpcDefinition"; const expectedWorkflow = "Workflow binding with create(), createBatch(), and get()"; export type WorkflowInstanceCreateOptions = Omit< CloudflareWorkflowInstanceCreateOptions, "params" >; export type WorkflowInstanceCreateBatchOptions = ReadonlyArray< { readonly payload: Payload } & WorkflowInstanceCreateOptions >; export interface WorkflowInstanceRestartOptions { readonly from?: { readonly name?: string; readonly count?: number; readonly type?: string; }; } export type WorkflowInstanceStatusName = CloudflareInstanceStatus["status"]; export interface WorkflowInstanceStatus { readonly status: WorkflowInstanceStatusName; readonly output: Option.Option; readonly error: Option.Option<{ readonly name: string; readonly message: string; }>; } export interface WorkflowInstance { readonly raw: CloudflareWorkflowInstance; readonly id: string; readonly pause: Effect.Effect; readonly resume: Effect.Effect; readonly terminate: Effect.Effect; readonly restart: ( options?: WorkflowInstanceRestartOptions, ) => Effect.Effect; readonly status: Effect.Effect< WorkflowInstanceStatus, WorkflowOperationError | WorkflowResultDecodeError >; readonly sendEvent: (event: WorkflowInstanceEvent) => Effect.Effect; } export interface WorkflowInstanceEvent { readonly type: string; readonly payload: unknown; } export interface WorkflowBindingDefinition< Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, > { /** Binding name as configured in `wrangler.jsonc`. */ readonly binding: string; /** Codec used to encode payloads passed to `Workflow.create`. */ readonly payload: Payload; /** Codec used to decode completed workflow status output. */ readonly result: Result; } export interface WorkflowBindingClient< Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, > { readonly create: ( payload: S.Schema.Type, options?: WorkflowInstanceCreateOptions>, ) => Effect.Effect< WorkflowInstance>, WorkflowOperationError | S.SchemaError >; readonly createBatch: ( batch: WorkflowInstanceCreateBatchOptions, S.Codec.Encoded>, ) => Effect.Effect< ReadonlyArray>>, WorkflowOperationError | S.SchemaError >; readonly get: ( instanceId: string, ) => Effect.Effect>, WorkflowOperationError>; readonly unsafeRaw: Effect.Effect>>; } export class WorkflowOperationError extends Data.TaggedError("WorkflowOperationError")<{ readonly binding: string; readonly operation: string; readonly cause: unknown; }> {} export class WorkflowResultDecodeError extends Data.TaggedError("WorkflowResultDecodeError")<{ readonly binding: string; readonly instanceId: string; readonly cause: unknown; }> {} const workflowError = (binding: string, operation: string, cause: unknown) => new WorkflowOperationError({ binding, operation, cause }); const tryWorkflowPromise = ( binding: string, operation: string, evaluate: () => Promise, ): Effect.Effect => Effect.tryPromise({ try: evaluate, catch: (cause) => workflowError(binding, operation, cause), }); export const isWorkflow = (value: unknown): value is CloudflareWorkflow => { if (typeof value !== "object" || value === null) { return false; } const resource = value as Record; return ( typeof resource.create === "function" && typeof resource.createBatch === "function" && typeof resource.get === "function" ); }; export const makeClient = < Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, >( definition: WorkflowBindingDefinition, ): (( workflow: CloudflareWorkflow>, ) => WorkflowBindingClient) => { type PayloadValue = S.Schema.Type; type EncodedPayload = S.Codec.Encoded; type ResultValue = S.Schema.Type; const encodePayload = S.encodeEffect(definition.payload); const decodeResult = S.decodeUnknownEffect(definition.result); const wrapInstance = (raw: CloudflareWorkflowInstance): WorkflowInstance => { const operation = (name: string, evaluate: () => Promise) => tryWorkflowPromise(definition.binding, name, evaluate); return { raw, id: raw.id, pause: operation("pause", () => raw.pause()), resume: operation("resume", () => raw.resume()), terminate: operation("terminate", () => raw.terminate()), restart: (options) => operation("restart", () => (raw as { restart(options?: WorkflowInstanceRestartOptions): Promise }).restart( options, ), ), status: operation("status", () => raw.status()).pipe( Effect.flatMap((status) => Effect.gen(function* () { const output = status.output === undefined ? Option.none() : Option.some( yield* decodeResult(status.output).pipe( Effect.mapError( (cause) => new WorkflowResultDecodeError({ binding: definition.binding, instanceId: raw.id, cause, }), ), ), ); return { status: status.status, output, error: status.error === undefined ? Option.none() : Option.some(status.error), }; }), ), ), sendEvent: (event) => operation("sendEvent", () => raw.sendEvent(event)), }; }; return (workflow) => ({ create: Effect.fnUntraced(function* ( payload: PayloadValue, options?: WorkflowInstanceCreateOptions, ) { const encoded = yield* encodePayload(payload); const raw = yield* tryWorkflowPromise(definition.binding, "create", () => workflow.create({ ...options, params: encoded }), ); return wrapInstance(raw); }), createBatch: Effect.fnUntraced(function* ( batch: WorkflowInstanceCreateBatchOptions, ) { const encodedBatch: Array> = []; for (const item of batch) { const { payload, ...options } = item; encodedBatch.push({ ...options, params: yield* encodePayload(payload), }); } const rawInstances = yield* tryWorkflowPromise(definition.binding, "createBatch", () => workflow.createBatch(encodedBatch), ); return rawInstances.map(wrapInstance); }), get: (instanceId) => tryWorkflowPromise(definition.binding, "get", () => workflow.get(instanceId)).pipe( Effect.map(wrapInstance), ), unsafeRaw: Effect.succeed(workflow), }); }; export const layer = < Self, Payload extends RpcDefinition.ServiceFreeSchema, Result extends RpcDefinition.ServiceFreeSchema, >( tag: Context.Service>, definition: WorkflowBindingDefinition, ) => Binding.layer( tag, definition.binding, (value): value is CloudflareWorkflow> => isWorkflow>(value), makeClient(definition), { expected: expectedWorkflow }, );