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, // local sessions,
// 4. runs the prompt queue until the footer closes. // 4. runs the prompt queue until the footer closes.
import { createOpencodeClient } from "@opencode-ai/sdk/v2" import { createOpencodeClient } from "@opencode-ai/sdk/v2"
import { Effect } from "effect"
import { Flag } from "@opencode-ai/core/flag/flag" import { Flag } from "@opencode-ai/core/flag/flag"
import { AppRuntime } from "@/effect/app-runtime"
import { createRunDemo } from "./demo" import { createRunDemo } from "./demo"
import { resolveDiffStyle, resolveFooterKeybinds, resolveModelInfo, resolveSessionInfo } from "./runtime.boot" import { resolveDiffStyle, resolveFooterKeybinds, resolveModelInfo, resolveSessionInfo } from "./runtime.boot"
import { createRuntimeLifecycle } from "./runtime.lifecycle" import { createRuntimeLifecycle } from "./runtime.lifecycle"
@ -70,9 +72,10 @@ type RunLocalInput = {
demo?: RunInput["demo"] demo?: RunInput["demo"]
} }
type StreamModule = Awaited<typeof import("./stream.transport")>
type StreamState = { type StreamState = {
mod: Awaited<typeof import("./stream.transport")> mod: StreamModule
handle: Awaited<ReturnType<Awaited<typeof import("./stream.transport")>["createSessionTransport"]>> handle: import("./stream.transport").SessionTransport
} }
type ResolvedSession = { type ResolvedSession = {
@ -485,7 +488,12 @@ async function runInteractiveRuntime(input: RunRuntimeInput): Promise<void> {
throw new Error("runtime closed") 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, sdk: ctx.sdk,
directory: ctx.directory, directory: ctx.directory,
sessionID: state.sessionID, sessionID: state.sessionID,
@ -494,6 +502,13 @@ async function runInteractiveRuntime(input: RunRuntimeInput): Promise<void> {
footer, footer,
trace: log, 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) { if (footer.isClosed) {
await handle.close() await handle.close()
throw new Error("runtime closed") throw new Error("runtime closed")

View file

@ -17,7 +17,7 @@
// delayed idle from an older turn cannot complete a newer busy turn. // delayed idle from an older turn cannot complete a newer busy turn.
import type { Event, GlobalEvent, OpencodeClient } from "@opencode-ai/sdk/v2" import type { Event, GlobalEvent, OpencodeClient } from "@opencode-ai/sdk/v2"
import { Context, Deferred, Effect, Exit, Layer, Scope, Stream } from "effect" import { Context, Deferred, Effect, Exit, Layer, Scope, Stream } from "effect"
import { makeRuntime } from "@/effect/run-service" import { AppRuntime } from "@/effect/app-runtime"
import { import {
blockerStatus, blockerStatus,
bootstrapSessionData, bootstrapSessionData,
@ -61,7 +61,7 @@ type Trace = {
write(type: string, data?: unknown): void write(type: string, data?: unknown): void
} }
type StreamInput = { export type StreamInput = {
sdk: OpencodeClient sdk: OpencodeClient
directory?: string directory?: string
sessionID: string sessionID: string
@ -89,12 +89,26 @@ export type SessionTurnInput = {
signal?: AbortSignal 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 = { export type SessionTransport = {
runPromptTurn(input: SessionTurnInput): Promise<void> runPromptTurn(input: SessionTurnInput): Promise<void>
selectSubagent(sessionID: string | undefined): void selectSubagent(sessionID: string | undefined): void
close(): Promise<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 = { type State = {
data: SessionData data: SessionData
subagent: SubagentData subagent: SubagentData
@ -107,14 +121,6 @@ type State = {
blockers: Map<string, number> 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 { function sid(event: Event): string | undefined {
if (event.type === "message.updated") { if (event.type === "message.updated") {
return event.properties.sessionID return event.properties.sessionID
@ -369,11 +375,7 @@ function traceTabs(trace: Trace | undefined, prev: FooterSubagentTab[], next: Fo
} }
} }
function createLayer(input: StreamInput) { const make = Effect.fn("SessionTransport.make")(function* (input: StreamInput) {
return Layer.fresh(
Layer.effect(
Service,
Effect.gen(function* () {
const scope = yield* Scope.make() const scope = yield* Scope.make()
const abort = yield* Scope.provide(scope)( const abort = yield* Scope.provide(scope)(
Effect.acquireRelease( Effect.acquireRelease(
@ -402,7 +404,6 @@ function createLayer(input: StreamInput) {
} }
input.signal?.addEventListener("abort", halt, { once: true }) input.signal?.addEventListener("abort", halt, { once: true })
yield* Effect.addFinalizer(() => closeScope())
const events = yield* Scope.provide(scope)( const events = yield* Scope.provide(scope)(
Effect.acquireRelease( Effect.acquireRelease(
@ -1033,15 +1034,17 @@ function createLayer(input: StreamInput) {
yield* closeScope() yield* closeScope()
}) })
return Service.of({ const handle: TransportHandle = {
runPromptTurn, runPromptTurn,
selectSubagent, selectSubagent,
close, 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. // 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 // The transport is single-turn: only one runPromptTurn() call can be active
// at a time. The prompt queue enforces this from above. // 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> { export async function createSessionTransport(input: StreamInput): Promise<SessionTransport> {
const runtime = makeRuntime(Service, createLayer(input)) const handle = await AppRuntime.runPromise(
await runtime.runPromise(() => Effect.void) Effect.gen(function* () {
const svc = yield* Service
return yield* svc.make(input)
}).pipe(Effect.provide(layer)),
)
return { return {
runPromptTurn: (next) => runtime.runPromise((svc) => svc.runPromptTurn(next)), runPromptTurn: (next) => AppRuntime.runPromise(handle.runPromptTurn(next)),
selectSubagent: (sessionID) => runtime.runSync((svc) => svc.selectSubagent(sessionID)), selectSubagent: (sessionID) => AppRuntime.runSync(handle.selectSubagent(sessionID)),
close: () => runtime.runPromise((svc) => svc.close()), close: () => AppRuntime.runPromise(handle.close()),
} }
} }
export * as SessionTransport from "./stream.transport"