import { LLMClient, LLMEvent, type LLMClientShape, type LLMClientService, type LLMError, type LLMRequest, } from "@opencode-ai/llm" import { Cause, Deferred, Effect, Exit, Fiber, Layer, Option, Queue, Ref, Stream } from "effect" class RunnerEndedError extends Error { constructor() { super("Runner completed before the next LLM request") } } class RunnerScenarioActiveError extends Error { constructor() { super("RunnerScenario.run is already active") } } class RunnerLLMRequestError extends Error { constructor(readonly id: number) { super(`Runner LLM request #${id} was not handled by the interaction`) } } class RunnerLLMResponseClosedError extends Error { constructor(readonly id: number) { super(`Runner LLM request #${id} is no longer awaiting a response`) } } export interface RunnerLLMCall { readonly request: LLMRequest readonly respond: { readonly stop: () => Effect.Effect readonly text: (text: string, options?: { readonly id?: string }) => Effect.Effect readonly toolCall: ( name: string, input: unknown, options?: { readonly id?: string; readonly providerExecuted?: boolean }, ) => Effect.Effect readonly events: (...events: ReadonlyArray) => Effect.Effect readonly stream: (stream: Stream.Stream) => Effect.Effect readonly fail: (error: LLMError) => Effect.Effect } } interface CallState { readonly generation: number readonly status: "pending" | "responded" | "closed" } interface State { readonly accepting: boolean readonly active: number | undefined readonly nextGeneration: number readonly nextID: number readonly calls: ReadonlyMap readonly requests: ReadonlyArray } type Registration = | { readonly _tag: "Registered"; readonly id: number; readonly generation: number } | { readonly _tag: "Unexpected"; readonly error: Error } interface PendingCall { readonly _tag: "Call" readonly generation: number readonly id: number readonly call: RunnerLLMCall } interface RunnerEnd { readonly _tag: "End" readonly generation: number } interface RunnerLLMInternal { readonly llm: RunnerLLM readonly begin: Effect.Effect readonly finishInteraction: (generation: number) => Effect.Effect readonly end: (generation: number) => Effect.Effect } export interface RunnerLLM { readonly layer: Layer.Layer readonly next: () => Effect.Effect readonly requests: Effect.Effect> } const events = (...events: ReadonlyArray) => Stream.fromIterable(events) const stop = () => events( LLMEvent.stepStart({ index: 0 }), LLMEvent.stepFinish({ index: 0, reason: "stop" }), LLMEvent.finish({ reason: "stop" }), ) const text = (value: string, options?: { readonly id?: string }) => { const id = options?.id ?? "text" return events( LLMEvent.stepStart({ index: 0 }), LLMEvent.textStart({ id }), LLMEvent.textDelta({ id, text: value }), LLMEvent.textEnd({ id }), LLMEvent.stepFinish({ index: 0, reason: "stop" }), LLMEvent.finish({ reason: "stop" }), ) } const toolCall = ( name: string, input: unknown, options?: { readonly id?: string; readonly providerExecuted?: boolean }, ) => { const id = options?.id ?? `call-${name}` return events( LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ id, name, input, providerExecuted: options?.providerExecuted }), LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), LLMEvent.finish({ reason: "tool-calls" }), ) } const makeRunnerLLM = Effect.gen(function* () { const calls = yield* Queue.unbounded() const state = yield* Ref.make({ accepting: false, active: undefined, nextGeneration: 1, nextID: 1, calls: new Map(), requests: [], }) const respond = ( id: number, response: Deferred.Deferred>, stream: Stream.Stream, ) => Effect.uninterruptible( Effect.gen(function* () { const pending = yield* Ref.modify(state, (current) => { const call = current.calls.get(id) if (call?.status !== "pending") return [false, current] return [true, { ...current, calls: new Map(current.calls).set(id, { ...call, status: "responded" }) }] }) if (!pending) return yield* Effect.fail(new RunnerLLMResponseClosedError(id)) yield* Deferred.succeed(response, stream) }), ) const stream = (request: LLMRequest) => Stream.unwrap( Effect.uninterruptibleMask((restore) => Effect.gen(function* () { const response = yield* Deferred.make>() const registered = yield* Ref.modify(state, (current) => { if (!current.accepting || current.active === undefined) { const error = new RunnerLLMRequestError(current.nextID) return [ { _tag: "Unexpected" as const, error }, { ...current, nextID: current.nextID + 1, requests: [...current.requests, request] }, ] } return [ { _tag: "Registered" as const, id: current.nextID, generation: current.active }, { ...current, nextID: current.nextID + 1, calls: new Map(current.calls).set(current.nextID, { generation: current.active, status: "pending", }), requests: [...current.requests, request], }, ] }) if (registered._tag === "Unexpected") return yield* Effect.fail(registered.error) const reply = (stream: Stream.Stream) => respond(registered.id, response, stream) yield* Queue.offer(calls, { _tag: "Call", generation: registered.generation, id: registered.id, call: { request, respond: { stop: () => reply(stop()), text: (value, options) => reply(text(value, options)), toolCall: (name, input, options) => reply(toolCall(name, input, options)), events: (...input) => reply(events(...input)), stream: reply, fail: (error) => reply(Stream.fail(error)), }, }, }) return yield* restore(Deferred.await(response)).pipe( Effect.onInterrupt(() => Ref.update(state, (current) => { const call = current.calls.get(registered.id) if (call?.status !== "pending") return current return { ...current, calls: new Map(current.calls).set(registered.id, { ...call, status: "closed" }) } }), ), ) }), ), ) const next = (): Effect.Effect => Effect.gen(function* () { const generation = (yield* Ref.get(state)).active if (generation === undefined) return yield* Effect.fail(new RunnerEndedError()) const pending = yield* Queue.take(calls) if (pending.generation !== generation) return yield* next() if (pending._tag === "End") return yield* Effect.fail(new RunnerEndedError()) if ((yield* Ref.get(state)).calls.get(pending.id)?.status === "pending") return pending.call return yield* next() }) return { llm: { layer: Layer.succeed( LLMClient.Service, LLMClient.Service.of({ prepare: () => Effect.die("RunnerLLM.prepare is not implemented"), stream: stream as LLMClientShape["stream"], generate: () => Effect.die("RunnerLLM.generate is not implemented"), }), ), next, requests: Ref.get(state).pipe(Effect.map((state) => state.requests)), }, begin: Ref.modify>(state, (current) => current.active !== undefined ? [Effect.fail(new RunnerScenarioActiveError()), current] : [ Effect.succeed(current.nextGeneration), { ...current, accepting: true, active: current.nextGeneration, nextGeneration: current.nextGeneration + 1, }, ], ).pipe(Effect.flatten), finishInteraction: (generation) => Effect.gen(function* () { const unanswered = yield* Ref.modify(state, (current) => [ [...current.calls].find(([, call]) => call.generation === generation && call.status === "pending")?.[0], current.active === generation ? { ...current, accepting: false } : current, ]) if (unanswered !== undefined) return yield* Effect.fail(new RunnerLLMRequestError(unanswered)) }), end: (generation) => Effect.gen(function* () { yield* Ref.update(state, (current) => ({ ...current, accepting: current.active === generation ? false : current.accepting, active: current.active === generation ? undefined : current.active, calls: new Map( [...current.calls].map(([id, call]) => [ id, call.generation === generation && call.status === "pending" ? { ...call, status: "closed" as const } : call, ]), ), })) yield* Queue.offer(calls, { _tag: "End", generation }) }), } satisfies RunnerLLMInternal }) export interface RunnerScenario { readonly llm: RunnerLLM readonly run: ( interaction: () => Effect.gen.Return, ) => Effect.Effect } export namespace RunnerScenario { export const make = ( start: (llm: RunnerLLM) => Effect.Effect, ): Effect.Effect> => Effect.gen(function* () { const internal = yield* makeRunnerLLM const run = ( interaction: () => Effect.gen.Return, ): Effect.Effect => Effect.uninterruptibleMask((restore) => Effect.gen(function* () { const generation = yield* internal.begin const runner = yield* start(internal.llm).pipe(Effect.ensuring(internal.end(generation)), Effect.forkChild) return yield* restore( Effect.gen(function* () { const interactionExit = yield* Effect.gen(interaction).pipe(Effect.exit) if (Exit.isSuccess(interactionExit)) { yield* internal.finishInteraction(generation) yield* Fiber.join(runner) return interactionExit.value } const error = Cause.findErrorOption(interactionExit.cause) if (Option.isSome(error) && error.value instanceof RunnerEndedError) yield* Fiber.join(runner) return yield* Effect.failCause(interactionExit.cause) }), ).pipe(Effect.ensuring(Fiber.interrupt(runner).pipe(Effect.andThen(internal.end(generation))))) }), ) return { llm: internal.llm, run, } }) }