import { Cause, Context, Effect, FiberMap, Layer, ManagedRuntime } from "effect"; export type WorkflowDefectHandler = (error: unknown) => void; export interface WorkflowScheduler { /** Starts a workflow immediately, interrupting and replacing the previous workflow at the key. */ start( key: Key, workflow: Effect.Effect, onDefect: WorkflowDefectHandler, ): void; /** Interrupts the current workflow at the key and waits for its fiber to settle. */ remove(key: Key): Promise; /** Interrupts every retained workflow and closes the scheduler scope. */ close(): Promise; } interface WorkflowSupervisorService { readonly start: ( key: unknown, workflow: Effect.Effect, ) => void; readonly remove: (key: unknown) => Effect.Effect; } class WorkflowSupervisor extends Context.Service< WorkflowSupervisor, WorkflowSupervisorService >()("@zachwill/pi-orchestrate/WorkflowSupervisor") {} const workflowSupervisorLayer = Layer.effect( WorkflowSupervisor, Effect.gen(function* () { const fibers = yield* FiberMap.make(); const run = yield* FiberMap.runtime(fibers)(); return WorkflowSupervisor.of({ start(key, workflow) { run(key, workflow); }, remove: (key) => FiberMap.remove(fibers, key), }); }), ); class EffectWorkflowScheduler implements WorkflowScheduler { private readonly managedRuntime = ManagedRuntime.make(workflowSupervisorLayer); private readonly supervisor = this.managedRuntime.runSync(WorkflowSupervisor); private closePromise: Promise | undefined; start( key: Key, workflow: Effect.Effect, onDefect: WorkflowDefectHandler, ): void { const supervised = workflow.pipe( Effect.catchCause((cause) => { if (!Cause.hasInterruptsOnly(cause)) { return Effect.sync(() => { try { onDefect(Cause.squash(cause)); } catch { // Defect reporting must not become another unsupervised defect. } }); } return Effect.void; }), ); this.supervisor.start(key, supervised); } remove(key: Key): Promise { return this.managedRuntime.runPromise(this.supervisor.remove(key)); } close(): Promise { this.closePromise ??= this.managedRuntime.dispose(); return this.closePromise; } } export function createWorkflowScheduler(): WorkflowScheduler { return new EffectWorkflowScheduler(); }