import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { InstanceState } from "@/effect/instance-state" import { SessionV1 } from "@opencode-ai/core/v1/session" import { Runner } from "@/effect/runner" import { BackgroundJob } from "@/background/job" import { Effect, Latch, Layer, Scope, Context } from "effect" import { Session } from "./session" import { SessionID } from "./schema" import { SessionStatus } from "./status" export interface Interface { readonly assertNotBusy: (sessionID: SessionID) => Effect.Effect readonly cancel: (sessionID: SessionID) => Effect.Effect readonly ensureRunning: ( sessionID: SessionID, onInterrupt: Effect.Effect, work: Effect.Effect, ) => Effect.Effect readonly startShell: ( sessionID: SessionID, onInterrupt: Effect.Effect, work: Effect.Effect, ready?: Latch.Latch, ) => Effect.Effect } export class Service extends Context.Service()("@opencode/SessionRunState") {} export const layer = Layer.effect( Service, Effect.gen(function* () { const background = yield* BackgroundJob.Service const status = yield* SessionStatus.Service const state = yield* InstanceState.make( Effect.fn("SessionRunState.state")(function* () { const scope = yield* Scope.Scope const runners = new Map>() yield* Effect.addFinalizer( Effect.fnUntraced(function* () { yield* Effect.forEach(runners.values(), (runner) => runner.cancel, { concurrency: "unbounded", discard: true, }) runners.clear() }), ) return { runners, scope } }), ) const runner = Effect.fn("SessionRunState.runner")(function* ( sessionID: SessionID, onInterrupt: Effect.Effect, ) { const data = yield* InstanceState.get(state) const existing = data.runners.get(sessionID) if (existing) return existing const next = Runner.make(data.scope, { onIdle: Effect.gen(function* () { data.runners.delete(sessionID) yield* status.set(sessionID, { type: "idle" }) }), onBusy: status.set(sessionID, { type: "busy" }), onInterrupt, }) data.runners.set(sessionID, next) return next }) const assertNotBusy = Effect.fn("SessionRunState.assertNotBusy")(function* (sessionID: SessionID) { const data = yield* InstanceState.get(state) const existing = data.runners.get(sessionID) if (existing?.busy) yield* busyError(sessionID) }) const cancel = Effect.fn("SessionRunState.cancel")(function* (sessionID: SessionID) { yield* cancelBackgroundJobs(background, sessionID) const data = yield* InstanceState.get(state) const existing = data.runners.get(sessionID) if (!existing) { yield* status.set(sessionID, { type: "idle" }) return } yield* existing.cancel }) const ensureRunning = Effect.fn("SessionRunState.ensureRunning")(function* ( sessionID: SessionID, onInterrupt: Effect.Effect, work: Effect.Effect, ) { return yield* (yield* runner(sessionID, onInterrupt)).ensureRunning(work) }) const startShell = Effect.fn("SessionRunState.startShell")(function* ( sessionID: SessionID, onInterrupt: Effect.Effect, work: Effect.Effect, ready?: Latch.Latch, ) { return yield* (yield* runner(sessionID, onInterrupt)) .startShell(work, ready) .pipe(Effect.catchTag("RunnerBusy", () => Effect.fail(busyError(sessionID)))) }) return Service.of({ assertNotBusy, cancel, ensureRunning, startShell }) }), ) export const defaultLayer = layer.pipe( Layer.provide(BackgroundJob.defaultLayer), Layer.provide(SessionStatus.defaultLayer), ) const cancelBackgroundJobs = Effect.fn("SessionRunState.cancelBackgroundJobs")(function* ( background: BackgroundJob.Interface, sessionID: SessionID, ) { const jobs = yield* background.list() const pending = new Set([sessionID]) const cancelled = new Set() const matches = (job: BackgroundJob.Info) => { if (job.status !== "running") return false if (cancelled.has(job.id)) return false if (pending.has(job.id)) return true if (typeof job.metadata?.sessionId === "string" && pending.has(job.metadata.sessionId)) return true return typeof job.metadata?.parentSessionId === "string" && pending.has(job.metadata.parentSessionId) } let batch = jobs.filter(matches) while (batch.length > 0) { yield* Effect.forEach( batch, (job) => background.cancel(job.id).pipe( Effect.tap(() => Effect.sync(() => { cancelled.add(job.id) pending.add(job.id) if (typeof job.metadata?.sessionId === "string") pending.add(job.metadata.sessionId) }), ), ), { concurrency: "unbounded", discard: true }, ) batch = jobs.filter(matches) } }) function busyError(sessionID: SessionID) { return new Session.BusyError({ sessionID }) } export const node = LayerNode.make(layer, [BackgroundJob.node, SessionStatus.node]) export * as SessionRunState from "./run-state"