Compare commits
1 commit
dev
...
effect/ses
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
248a45801d |
2 changed files with 644 additions and 616 deletions
|
|
@ -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")
|
||||||
|
|
|
||||||
|
|
@ -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"
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue