Compare commits

...
Sign in to create a new pull request.

1 commit

Author SHA1 Message Date
Kit Langton
248a45801d effect(cli): convert stream.transport to SessionTransport.Service
Replace the per-input `Layer.fresh` + `makeRuntime` facade with a proper
`SessionTransport.Service` exposing a `make(input)` method that returns the
transport handle. The layer is now static; per-input state lives inside the
handle returned by `svc.make(...)`. Update the runtime.ts caller to yield
the service inside `Effect.gen` and wrap the resulting `TransportHandle` in
the existing Promise-shaped `SessionTransport` for the Promise-based
prompt-queue code. `createSessionTransport(...)` is kept as a thin
backward-compatible wrapper for tests.
2026-05-13 19:41:54 -04:00
2 changed files with 644 additions and 616 deletions

View file

@ -13,7 +13,9 @@
// local sessions,
// 4. runs the prompt queue until the footer closes.
import { createOpencodeClient } from "@opencode-ai/sdk/v2"
import { Effect } from "effect"
import { Flag } from "@opencode-ai/core/flag/flag"
import { AppRuntime } from "@/effect/app-runtime"
import { createRunDemo } from "./demo"
import { resolveDiffStyle, resolveFooterKeybinds, resolveModelInfo, resolveSessionInfo } from "./runtime.boot"
import { createRuntimeLifecycle } from "./runtime.lifecycle"
@ -70,9 +72,10 @@ type RunLocalInput = {
demo?: RunInput["demo"]
}
type StreamModule = Awaited<typeof import("./stream.transport")>
type StreamState = {
mod: Awaited<typeof import("./stream.transport")>
handle: Awaited<ReturnType<Awaited<typeof import("./stream.transport")>["createSessionTransport"]>>
mod: StreamModule
handle: import("./stream.transport").SessionTransport
}
type ResolvedSession = {
@ -485,7 +488,12 @@ async function runInteractiveRuntime(input: RunRuntimeInput): Promise<void> {
throw new Error("runtime closed")
}
const handle = await mod.createSessionTransport({
// Yield SessionTransport.Service inside Effect.gen and call svc.make.
// The handle holds its own internal scope; teardown happens via handle.close().
const inner = await AppRuntime.runPromise(
Effect.gen(function* () {
const svc = yield* mod.SessionTransport.Service
return yield* svc.make({
sdk: ctx.sdk,
directory: ctx.directory,
sessionID: state.sessionID,
@ -494,6 +502,13 @@ async function runInteractiveRuntime(input: RunRuntimeInput): Promise<void> {
footer,
trace: log,
})
}).pipe(Effect.provide(mod.SessionTransport.layer)),
)
const handle: import("./stream.transport").SessionTransport = {
runPromptTurn: (next) => AppRuntime.runPromise(inner.runPromptTurn(next)),
selectSubagent: (sessionID) => AppRuntime.runSync(inner.selectSubagent(sessionID)),
close: () => AppRuntime.runPromise(inner.close()),
}
if (footer.isClosed) {
await handle.close()
throw new Error("runtime closed")

View file

@ -17,7 +17,7 @@
// delayed idle from an older turn cannot complete a newer busy turn.
import type { Event, GlobalEvent, OpencodeClient } from "@opencode-ai/sdk/v2"
import { Context, Deferred, Effect, Exit, Layer, Scope, Stream } from "effect"
import { makeRuntime } from "@/effect/run-service"
import { AppRuntime } from "@/effect/app-runtime"
import {
blockerStatus,
bootstrapSessionData,
@ -61,7 +61,7 @@ type Trace = {
write(type: string, data?: unknown): void
}
type StreamInput = {
export type StreamInput = {
sdk: OpencodeClient
directory?: string
sessionID: string
@ -89,12 +89,26 @@ export type SessionTurnInput = {
signal?: AbortSignal
}
export type TransportHandle = {
readonly runPromptTurn: (input: SessionTurnInput) => Effect.Effect<void, unknown>
readonly selectSubagent: (sessionID: string | undefined) => Effect.Effect<void>
readonly close: () => Effect.Effect<void>
}
/** Backward-compatible promise-shaped handle used by Promise-based callers. */
export type SessionTransport = {
runPromptTurn(input: SessionTurnInput): Promise<void>
selectSubagent(sessionID: string | undefined): void
close(): Promise<void>
}
export interface Interface {
/** Build a transport handle for the given input. Holds its own internal scope; tear down via `handle.close()`. */
readonly make: (input: StreamInput) => Effect.Effect<TransportHandle>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/SessionTransport") {}
type State = {
data: SessionData
subagent: SubagentData
@ -107,14 +121,6 @@ type State = {
blockers: Map<string, number>
}
type TransportService = {
readonly runPromptTurn: (input: SessionTurnInput) => Effect.Effect<void, unknown>
readonly selectSubagent: (sessionID: string | undefined) => Effect.Effect<void>
readonly close: () => Effect.Effect<void>
}
class Service extends Context.Service<Service, TransportService>()("@opencode/RunStreamTransport") {}
function sid(event: Event): string | undefined {
if (event.type === "message.updated") {
return event.properties.sessionID
@ -369,11 +375,7 @@ function traceTabs(trace: Trace | undefined, prev: FooterSubagentTab[], next: Fo
}
}
function createLayer(input: StreamInput) {
return Layer.fresh(
Layer.effect(
Service,
Effect.gen(function* () {
const make = Effect.fn("SessionTransport.make")(function* (input: StreamInput) {
const scope = yield* Scope.make()
const abort = yield* Scope.provide(scope)(
Effect.acquireRelease(
@ -402,7 +404,6 @@ function createLayer(input: StreamInput) {
}
input.signal?.addEventListener("abort", halt, { once: true })
yield* Effect.addFinalizer(() => closeScope())
const events = yield* Scope.provide(scope)(
Effect.acquireRelease(
@ -1033,15 +1034,17 @@ function createLayer(input: StreamInput) {
yield* closeScope()
})
return Service.of({
const handle: TransportHandle = {
runPromptTurn,
selectSubagent,
close,
})
}),
),
)
}
}
return handle
})
export const layer = Layer.succeed(Service, Service.of({ make }))
export const defaultLayer = layer
// Opens an SDK event subscription and returns a SessionTransport.
//
@ -1052,13 +1055,23 @@ function createLayer(input: StreamInput) {
//
// The transport is single-turn: only one runPromptTurn() call can be active
// at a time. The prompt queue enforces this from above.
//
// Promise-shaped wrapper for Promise-based callers (CLI runtime, tests). New
// Effect-style callers should yield `SessionTransport.Service` and call
// `svc.make(input)` directly inside their own scope.
export async function createSessionTransport(input: StreamInput): Promise<SessionTransport> {
const runtime = makeRuntime(Service, createLayer(input))
await runtime.runPromise(() => Effect.void)
const handle = await AppRuntime.runPromise(
Effect.gen(function* () {
const svc = yield* Service
return yield* svc.make(input)
}).pipe(Effect.provide(layer)),
)
return {
runPromptTurn: (next) => runtime.runPromise((svc) => svc.runPromptTurn(next)),
selectSubagent: (sessionID) => runtime.runSync((svc) => svc.selectSubagent(sessionID)),
close: () => runtime.runPromise((svc) => svc.close()),
runPromptTurn: (next) => AppRuntime.runPromise(handle.runPromptTurn(next)),
selectSubagent: (sessionID) => AppRuntime.runSync(handle.selectSubagent(sessionID)),
close: () => AppRuntime.runPromise(handle.close()),
}
}
export * as SessionTransport from "./stream.transport"