Compare commits

..
Author SHA1 Message Date
Kit Langton b171a295df fix(cli): adopt the login shell environment in the background service
A managed service elected by a GUI client (the desktop app, an editor, a
login item) inherits launchd's environment: PATH is /usr/bin:/bin:/usr/sbin:/sbin
and nothing exported from .zshrc or .bashrc exists. Every process the server
spawns extends process.env, so stdio MCP servers using #!/usr/bin/env node,
npx, uvx, formatters, and hooks fail with "Connection closed" while the bash
tool, which prefers the client's terminal environment, keeps working.

Resolve the user's interactive login shell once at service startup when the
inherited environment did not come from a shell, and adopt it before anything
reads process.env.
2026-09-16 11:57:09 -07:00
16 changed files with 325 additions and 308 deletions
@@ -73,6 +73,11 @@ const driver = (options: Options, body: string): WebSocketChannelDriver => {
)
if (event.type === "error") {
terminal = true
yield* OpenResponses.decodeKnownErrorEvent(event).pipe(
Effect.mapError((cause) =>
ProviderShared.eventError(options.id, `${options.name} returned a malformed error event`, frame, cause),
),
)
return {
type: "provider-failure",
error: OpenResponses.providerFailure(event, `${options.name} stream error`, frame),
@@ -108,7 +108,7 @@ const incremental = (
return input.slice(baseline.length)
}
const code = (event: OpenResponses.Event) => OpenResponses.errorDetail(event).code
const code = (event: OpenResponses.Event) => event.code || event.error?.code || event.response?.error?.code || undefined
const rejected = (
observation: Extract<ChannelObservation, { readonly type: "provider-failure" }>,
+94 -28
View File
@@ -1,4 +1,4 @@
import { Effect, Option, Schema } from "effect"
import { Effect, Option, Schema, SchemaGetter } from "effect"
import type { Content } from "@opencode/schema/tool"
import { HttpTransport } from "../route/transport/index.js"
import { Protocol } from "../route/protocol.js"
@@ -333,13 +333,53 @@ export const StreamItem = Schema.StructWithRest(
export type StreamItem = Schema.Schema.Type<typeof StreamItem>
export type OutputItem = StreamItem & { readonly id: string }
// Responses-compatible providers put error details at the top level, under `error`, or under
// `response.error`, and gateways reshape them freely: strings, numeric codes, extra fields. Those
// fields decode as opaque values and `errorDetail` reads them defensively, so an error frame can
// only fail on invalid JSON and otherwise always classifies with the raw body as the fallback.
// Responses-compatible providers put streaming error details at the top level or
// under `error`, and response failures under `response.error`. Accept all three shapes.
// https://www.openresponses.org/specification
const asText = (value: unknown) =>
typeof value === "string" && value.length > 0 ? value : typeof value === "number" ? String(value) : undefined
const OpenResponsesErrorObject = Schema.Struct({
type: optionalNull(Schema.String),
code: optionalNull(Schema.String),
message: optionalNull(Schema.String),
param: optionalNull(Schema.String),
})
const OpenResponsesErrorPayload = Schema.Union([Schema.String, OpenResponsesErrorObject]).pipe(
Schema.decodeTo(OpenResponsesErrorObject, {
decode: SchemaGetter.transform((error) => (typeof error === "string" ? { message: error } : error)),
encode: SchemaGetter.passthrough(),
}),
)
type OpenResponsesErrorPayload = Schema.Schema.Type<typeof OpenResponsesErrorPayload>
const WebSocketErrorHeader = Schema.Union([Schema.String, Schema.Number, Schema.Boolean])
export const WebSocketErrorEvent = Schema.StructWithRest(
Schema.Struct({
type: Schema.tag("error"),
status: Schema.optional(Schema.Number),
status_code: Schema.optional(Schema.Number),
code: optionalNull(Schema.String),
message: Schema.optional(Schema.String),
param: optionalNull(Schema.String),
error: optionalNull(OpenResponsesErrorPayload),
headers: Schema.optional(Schema.Record(Schema.String, WebSocketErrorHeader)),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
)
const decodeWebSocketErrorEvent = Schema.decodeUnknownEffect(WebSocketErrorEvent)
export const decodeKnownErrorEvent = (event: Event) =>
decodeWebSocketErrorEvent({
...event,
status: typeof event.status === "number" ? event.status : undefined,
status_code: typeof event.status_code === "number" ? event.status_code : undefined,
headers: ProviderShared.isRecord(event.headers)
? Object.fromEntries(
Object.entries(event.headers).filter(
(entry): entry is [string, string | number | boolean] =>
typeof entry[1] === "string" || typeof entry[1] === "number" || typeof entry[1] === "boolean",
),
)
: undefined,
})
export const Event = Schema.StructWithRest(
Schema.Struct({
@@ -360,18 +400,31 @@ export const Event = Schema.StructWithRest(
incomplete_details: optionalNull(Schema.Struct({ reason: Schema.optional(Schema.String) })),
output: Schema.optional(Schema.Array(StreamItem)),
usage: optionalNull(OpenResponsesUsage),
error: Schema.optional(Schema.Unknown),
error: optionalNull(OpenResponsesErrorPayload),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
),
),
code: Schema.optional(Schema.Unknown),
message: Schema.optional(Schema.Unknown),
error: Schema.optional(Schema.Unknown),
code: optionalNull(Schema.String),
message: Schema.optional(Schema.String),
param: optionalNull(Schema.String),
error: optionalNull(OpenResponsesErrorPayload),
status: Schema.optional(Schema.Unknown),
status_code: Schema.optional(Schema.Unknown),
headers: Schema.optional(Schema.Unknown),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
).pipe(
Schema.decode({
decode: SchemaGetter.transform((event) => {
if (event.type !== "error" || event.error != null) return event
const { code, message, param, ...rest } = event
if (code === undefined && message === undefined && param === undefined) return event
// Flat errors (for example, Meta's) can also arrive through generic Responses endpoints.
return { ...rest, error: { code, message, param } }
}),
encode: SchemaGetter.passthrough(),
}),
)
export type Event = Schema.Schema.Type<typeof Event>
export type NormalizedEvent = Event & { readonly item?: OutputItem | null }
@@ -380,15 +433,16 @@ const decodeEventValue = Schema.decodeUnknownEffect(Event)
const decodeFrame = Schema.decodeUnknownEffect(ProviderShared.Json)
/**
* Decodes one WebSocket frame. Some providers and gateways answer a rejected `response.create` with a bare
* `{ "error": ... }` envelope and no event type; that reads as an error event so it classifies instead of
* failing decoding.
* Decodes one WebSocket frame. xAI answers a rejected `response.create` with `{ "error": { "message", "type" } }` and no
* event type; that envelope reads as an error event so the failure classifies instead of failing decoding.
*/
export const decodeChannelEvent = (frame: string) =>
decodeFrame(frame).pipe(
Effect.flatMap((value) =>
decodeEventValue(
ProviderShared.isRecord(value) && value.type === undefined && value.error != null
ProviderShared.isRecord(value) &&
value.type === undefined &&
(typeof value.error === "string" || ProviderShared.isRecord(value.error))
? { ...value, type: "error" }
: value,
),
@@ -1368,21 +1422,22 @@ const onResponseFinish = Effect.fn("OpenResponses.onResponseFinish")(function* (
return [{ ...current, lifecycle }, events] satisfies StepResult
})
/** Error code and message from wherever the frame put them; top-level fields win over nested ones. */
export const errorDetail = (event: Event) => {
const raw = event.error ?? event.response?.error
const nested = typeof raw === "string" ? { message: raw } : ProviderShared.isRecord(raw) ? raw : undefined
return {
message: asText(event.message) ?? asText(nested?.message),
code: asText(event.code) ?? asText(nested?.code),
}
// Build the prettiest summary available from whatever the provider supplied.
// When both code and message are present, prefix the code so consumers see
// the failure mode (e.g. `rate_limit_exceeded: Slow down`) instead of just
// the bare message — production rate limits and context-length failures used
// to be indistinguishable from generic stream drops. Returns undefined when
// the payload carries no usable summary.
const providerErrorMessage = (event: Event, nested: OpenResponsesErrorPayload | undefined): string | undefined => {
const message = event.message || nested?.message || undefined
const code = event.code || nested?.code || undefined
if (message && code) return `${code}: ${message}`
return message || code
}
// Prefix the code when both are present (`rate_limit_exceeded: Slow down`) so the failure mode is
// visible; fall back to the raw frame rather than a generic message when neither decodes.
export const providerFailure = (event: Event, fallback: string, body = ProviderShared.encodeJson(event)) => {
const detail = errorDetail(event)
const summary = detail.message && detail.code ? `${detail.code}: ${detail.message}` : (detail.message ?? detail.code)
const nested = event.error ?? event.response?.error ?? undefined
const summary = providerErrorMessage(event, nested)
const message = summary ?? (body === "{}" ? fallback : body)
const status =
typeof event.status === "number"
@@ -1465,7 +1520,18 @@ export const step = (state: ParserState, event: NormalizedEvent) => {
if (event.type === "response.output_item.done") return onOutputItemDone(state, event.item)
if (event.type === "response.completed" || event.type === "response.incomplete") return onResponseFinish(state, event)
if (event.type === "response.failed") return providerFailure(event, `${state.name} response failed`)
if (event.type === "error") return providerFailure(event, `${state.name} stream error`)
if (event.type === "error")
return decodeKnownErrorEvent(event).pipe(
Effect.mapError((cause) =>
ProviderShared.eventError(
state.id,
`${state.name} returned a malformed error event`,
ProviderShared.encodeJson(event),
cause,
),
),
Effect.flatMap(() => providerFailure(event, `${state.name} stream error`)),
)
return Effect.succeed<StepResult>([state, NO_EVENTS])
}
@@ -11,78 +11,66 @@ import { sseEvents } from "../lib/sse.js"
const decodeEvent = Schema.decodeUnknownEffect(OpenResponses.protocol.stream.event)
it.effect("decodes error frames verbatim in shared SSE and WebSocket decoding", () =>
it.effect("normalizes flat errors in shared SSE and WebSocket decoding", () =>
Effect.gen(function* () {
const frame = {
type: "error",
sequence_number: 4,
code: "server_shutting_down",
message: "Server is shutting down. Please retry your request.",
param: null,
}
for (const decode of [decodeEvent, OpenResponses.decodeChannelEvent]) {
for (const frame of [
{ type: "error", sequence_number: 4, code: "server_shutting_down", message: "Shutting down", param: null },
const event = yield* decode(JSON.stringify(frame))
expect(event).toEqual({
type: "error",
sequence_number: 4,
error: { code: frame.code, message: frame.message, param: null },
})
for (const unchanged of [
event,
{ type: "error" },
{ type: "error", error: "Gateway failed" },
{ type: "error", error: { code: 429, message: "slow down" } },
{ type: "error", error: 42 },
{ type: "error", code: 500, message: ["not", "a", "string"] },
{ type: "response.failed", response: { id: "resp_failed", error: "Gateway failed" } },
{ type: "response.failed", response: { id: "resp_failed", error: ["weird"] } },
{
type: "response.failed",
response: { id: "resp_failed", error: { code: "server_error", message: "Internal server error" } },
},
{ type: "response.output_text.delta", item_id: "msg_text", delta: "Hello" },
]) {
expect(yield* decode(JSON.stringify(frame))).toEqual(frame)
expect(yield* decode(JSON.stringify(unchanged))).toEqual(unchanged)
}
}
}),
)
it.effect("reads bare WebSocket error envelopes as error events", () =>
it.effect("continues to normalize untyped xAI WebSocket errors", () =>
Effect.gen(function* () {
const message = "gRPC error: Response with id=resp_missing not found"
for (const error of [{ type: "api_error", message }, message, 42]) {
expect(yield* OpenResponses.decodeChannelEvent(JSON.stringify({ error }))).toEqual({ type: "error", error })
}
for (const frame of [{ error: null }, { message }]) {
expect(yield* OpenResponses.decodeChannelEvent(JSON.stringify(frame)).pipe(Effect.flip)).toBeDefined()
for (const error of [{ type: "api_error", message }, message]) {
expect(yield* OpenResponses.decodeChannelEvent(JSON.stringify({ error }))).toEqual({
type: "error",
error: typeof error === "string" ? { message } : error,
})
}
}),
)
it.effect("extracts error details from every shape and falls back to the raw frame", () =>
it.effect("normalizes string errors in shared SSE and WebSocket decoding", () =>
Effect.gen(function* () {
const cases: Array<[frame: Record<string, unknown>, message: string, tag: string]> = [
[
{ type: "error", code: "server_shutting_down", message: "Shutting down" },
"server_shutting_down: Shutting down",
"UnknownProvider",
],
[{ type: "error", error: "Gateway failed" }, "Gateway failed", "UnknownProvider"],
[{ type: "error", error: { code: 429, message: "slow down" } }, "429: slow down", "UnknownProvider"],
[{ type: "error", error: { message: "slow down" }, status: 429 }, "slow down", "RateLimit"],
[{ type: "error", code: 500, message: ["not", "a", "string"] }, "500", "UnknownProvider"],
[
{ type: "response.failed", response: { id: "resp_failed", error: "Gateway failed" } },
"Gateway failed",
"UnknownProvider",
],
]
for (const [frame, message, tag] of cases) {
const event = yield* OpenResponses.decodeChannelEvent(JSON.stringify(frame))
const error = OpenResponses.providerFailure(event, "fallback", JSON.stringify(frame))
expect(error.message).toBe(message)
expect(error.reason._tag).toBe(tag)
expect(error.reason.body).toBe(JSON.stringify(frame))
for (const decode of [decodeEvent, OpenResponses.decodeChannelEvent]) {
expect(yield* decode(JSON.stringify({ type: "error", error: "Gateway failed" }))).toEqual({
type: "error",
error: { message: "Gateway failed" },
})
expect(
yield* decode(
JSON.stringify({ type: "response.failed", response: { id: "resp_failed", error: "Gateway failed" } }),
),
).toEqual({
type: "response.failed",
response: { id: "resp_failed", error: { message: "Gateway failed" } },
})
}
for (const frame of [
{ type: "error", error: 42 },
{ type: "response.failed", response: { id: "resp_failed", error: ["weird"] } },
]) {
const event = yield* OpenResponses.decodeChannelEvent(JSON.stringify(frame))
const error = OpenResponses.providerFailure(event, "fallback", JSON.stringify(frame))
expect(error.message).toBe(JSON.stringify(frame))
expect(error.reason._tag).toBe("UnknownProvider")
}
expect(OpenResponses.providerFailure({ type: "error" }, "fallback", "{}").message).toBe("fallback")
expect(OpenResponses.providerFailure({ type: "error" }, "fallback", "{}").reason._tag).toBe("ProviderInternal")
}),
)
+5 -2
View File
@@ -13,6 +13,7 @@ import { HttpServer } from "effect/unstable/http"
import { Env } from "./env"
import { ServiceConfig } from "./services/service-config"
import { ServiceRegistration } from "./services/service-registration"
import { ShellEnvironment } from "./shell-environment"
import { Updater } from "./services/updater"
import { WebUi } from "./services/web-ui"
import { databasePath } from "./database-path"
@@ -28,6 +29,9 @@ export type Options = {
// The process effect lives until server shutdown; tracing it would parent every request to one process-lifetime trace.
export const run = Effect.fnUntraced(function* (options: Options) {
// A managed service may have been started by a GUI client with launchd's environment. Adopt the
// login shell's before anything reads process.env, including the OPENCODE_* settings below.
if (options.mode === "service") yield* ShellEnvironment.adopt()
return yield* processEffect(options).pipe(
Effect.provide(
LayerNode.compile(LayerNode.group([Global.node, AppProcess.node]), {
@@ -38,9 +42,8 @@ export const run = Effect.fnUntraced(function* (options: Options) {
],
}),
),
Effect.provide(NodeServices.layer),
)
})
}, Effect.provide(NodeServices.layer))
const processEffect = Effect.fnUntraced(function* (options: Options) {
const inherited = process.env.OPENCODE_PTY_HANDOFF
+80
View File
@@ -0,0 +1,80 @@
export * as ShellEnvironment from "./shell-environment"
import { randomUUID } from "node:crypto"
import { Effect } from "effect"
import { ChildProcess } from "effect/unstable/process"
import { ChildProcessSpawner } from "effect/unstable/process/ChildProcessSpawner"
// Set on the probe shell so an rc file that itself starts opencode cannot recurse into another probe.
export const RESOLVING = "OPENCODE_RESOLVING_SHELL_ENVIRONMENT"
// Shell bookkeeping that describes the probe process, not the user's environment.
const TRANSIENT = new Set([RESOLVING, "SHLVL", "PWD", "OLDPWD", "_"])
export type Variables = Readonly<Record<string, string | undefined>>
/**
* Resolve the environment of the user's interactive login shell when the current environment did
* not come from one.
*
* A background service elected by a GUI client (the desktop app, an editor, a login item) inherits
* launchd's environment: PATH is `/usr/bin:/bin:/usr/sbin:/sbin` and nothing exported from
* `.zshrc` or `.bashrc` exists. Every process the server spawns extends `process.env`, so stdio MCP
* servers using `#!/usr/bin/env node`, `npx`, `uvx`, formatters, and hooks fail with no useful error
* while the bash tool, which prefers the client's terminal environment, keeps working. Editors such
* as VS Code and Zed resolve the login shell at startup for the same reason.
*
* Returns `undefined` when there is nothing to adopt: Windows, an environment that already passed
* through a shell (`SHLVL`), a probe already in progress, or a shell that fails to report.
*/
export const resolve = Effect.fnUntraced(function* (env: Variables = process.env) {
if (process.platform === "win32") return undefined
if (env[RESOLVING] !== undefined || env.SHLVL !== undefined) return undefined
const spawner = yield* ChildProcessSpawner
const marker = randomUUID()
// Interactive and login: most PATH edits live in `.zshrc`/`.bashrc`, which only interactive shells read.
// `env -0` keeps values containing newlines intact.
const output = yield* spawner
.string(
ChildProcess.make(env.SHELL || "/bin/sh", ["-ilc", `printf '%s' '${marker}'; env -0; printf '%s' '${marker}'`], {
env: { ...env, [RESOLVING]: "1" },
stdin: "ignore",
// Interactive shells without a terminal warn on stderr; nothing reads it, so never let it fill.
stderr: "ignore",
}),
)
.pipe(
Effect.timeout("10 seconds"),
Effect.tapError((error) => Effect.logWarning("shell environment unavailable", { error })),
Effect.orElseSucceed(() => undefined),
)
if (output === undefined) return undefined
const start = output.indexOf(marker)
const end = output.lastIndexOf(marker)
if (start === -1 || end === start) {
yield* Effect.logWarning("shell environment unavailable", { shell: env.SHELL, reason: "missing markers" })
return undefined
}
return Object.fromEntries(
output
.slice(start + marker.length, end)
.split("\0")
.flatMap((entry) => {
const separator = entry.indexOf("=")
if (separator <= 0) return []
const key = entry.slice(0, separator)
return TRANSIENT.has(key) ? [] : [[key, entry.slice(separator + 1)] as const]
}),
)
})
/** Apply the resolved login-shell environment to this process before anything reads `process.env`. */
export const adopt = Effect.fnUntraced(function* () {
const variables = yield* resolve()
if (variables === undefined) return
Object.assign(process.env, variables)
yield* Effect.logInfo("shell environment adopted", {
shell: process.env.SHELL,
variables: Object.keys(variables).length,
})
})
@@ -0,0 +1,68 @@
import { NodeServices } from "@effect/platform-node"
import { afterAll, beforeAll, expect, test } from "bun:test"
import { Effect } from "effect"
import fs from "node:fs/promises"
import os from "node:os"
import path from "node:path"
import { ShellEnvironment } from "../src/shell-environment"
// A stand-in login shell: it records every invocation, applies "rc file" exports, then runs the
// probe script exactly as `$SHELL -ilc <script>` would.
let root: string
let shell: string
let invocations: string
beforeAll(async () => {
root = await fs.mkdtemp(path.join(os.tmpdir(), "opencode-shell-env-"))
shell = path.join(root, "login-shell")
invocations = path.join(root, "invocations")
await fs.writeFile(
shell,
[
"#!/bin/sh",
`printf '%s\\n' "$1" >> '${invocations}'`,
`[ "$1" = -ilc ] || exit 64`,
'export PATH="/rc/bin:$PATH"',
"export FROM_RC=yes",
"export MULTILINE='first",
"second'",
'export SAW_GUARD="${OPENCODE_RESOLVING_SHELL_ENVIRONMENT:-unset}"',
"export SHLVL=1 PWD=/rc OLDPWD=/rc",
'eval "$2"',
].join("\n"),
{ mode: 0o755 },
)
})
afterAll(() => fs.rm(root, { recursive: true, force: true }))
const skip = process.platform === "win32"
const launchd = { SHELL: "", HOME: "", PATH: "/usr/bin:/bin:/usr/sbin:/sbin" }
const resolve = (env: Record<string, string | undefined>) =>
Effect.runPromise(
ShellEnvironment.resolve({ ...launchd, ...env, SHELL: env.SHELL ?? shell, HOME: root }).pipe(
Effect.provide(NodeServices.layer),
),
)
test.skipIf(skip)("adopts the login shell's exports over a launchd environment", async () => {
const variables = await resolve({})
expect(variables?.PATH).toBe("/rc/bin:/usr/bin:/bin:/usr/sbin:/sbin")
expect(variables?.FROM_RC).toBe("yes")
expect(variables?.MULTILINE).toBe("first\nsecond")
expect(variables?.SAW_GUARD).toBe("1")
for (const key of [ShellEnvironment.RESOLVING, "SHLVL", "PWD", "OLDPWD", "_"])
expect(variables).not.toHaveProperty(key)
})
test.skipIf(skip)("does not probe when the environment already came from a shell or a probe", async () => {
const before = await fs.readFile(invocations, "utf8").catch(() => "")
expect(await resolve({ SHLVL: "1" })).toBeUndefined()
expect(await resolve({ [ShellEnvironment.RESOLVING]: "1" })).toBeUndefined()
expect(await fs.readFile(invocations, "utf8").catch(() => "")).toBe(before)
})
test.skipIf(skip)("falls back to the inherited environment when the shell cannot report", async () => {
expect(await resolve({ SHELL: "/bin/false" })).toBeUndefined()
expect(await resolve({ SHELL: path.join(root, "missing-shell") })).toBeUndefined()
})
-3
View File
@@ -57,9 +57,6 @@ export const Plugin = define({
const hook = (event: SessionHooks["context"]) =>
Effect.gen(function* () {
const session = yield* ctx.session.get({ sessionID: event.sessionID }).pipe(Effect.orDie)
if (session.parentID) return
const active = sessions.get(event.sessionID)
const settings = yield* loadSettings()
if (!settings) {
+9 -18
View File
@@ -320,24 +320,15 @@ export const layer = Layer.effect(
// which transport actually carries the request, so both hook families are always offered.
const webSocket =
input.webSocket === "session" && model.transport === "websocket"
? transport.bind(session.id, {
handshake: (connect) =>
hooks
.trigger("session", "experimental.ws.handshake", {
...scope,
url: connect.url,
headers: connect.headers,
})
.pipe(Effect.map((event) => ({ url: event.url, headers: event.headers }))),
send: (frame) =>
hooks
.trigger("session", "experimental.ws.send", { ...scope, frame })
.pipe(Effect.map((event) => event.frame)),
receive: (frame) =>
hooks
.trigger("session", "experimental.ws.receive", { ...scope, frame })
.pipe(Effect.map((event) => event.frame)),
})
? transport.bind(session.id, (connect) =>
hooks
.trigger("session", "experimental.ws.handshake", {
...scope,
url: connect.url,
headers: connect.headers,
})
.pipe(Effect.map((event) => ({ url: event.url, headers: event.headers }))),
)
: undefined
return {
+13 -21
View File
@@ -59,18 +59,11 @@ export interface Handshake {
readonly headers: Record<string, string>
}
/**
* Per-exchange taps. `handshake` runs before the connection is selected; `send` sees each outbound
* frame after the driver builds it; `receive` sees each inbound frame before the driver observes it.
*/
export interface Interceptor {
readonly handshake?: (connect: Handshake) => Effect.Effect<Handshake>
readonly send?: (frame: string) => Effect.Effect<string>
readonly receive?: (frame: string) => Effect.Effect<string>
}
export interface Interface {
readonly bind: (sessionID: SessionSchema.ID, interceptor?: Interceptor) => WebSocketChannelExecutor
readonly bind: (
sessionID: SessionSchema.ID,
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
) => WebSocketChannelExecutor
readonly close: (sessionID: SessionSchema.ID) => Effect.Effect<void>
readonly closeAll: Effect.Effect<void>
}
@@ -285,7 +278,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
const start = Effect.fn("SessionModelTransport.start")(function* (
owner: State,
input: WebSocketChannelExchange,
interceptor?: Interceptor,
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
) {
if (owner.closed)
return yield* transportError("Session WebSocket owner is closed", {
@@ -295,8 +288,8 @@ export const makeLayer = (connector: WebSocketConnector) =>
delivery: "not-sent",
})
if (owner.httpFallback) return fallback(input)
const selected = interceptor?.handshake
? yield* interceptor.handshake({ url: input.connect.url, headers: { ...input.connect.headers } })
const selected = handshake
? yield* handshake({ url: input.connect.url, headers: { ...input.connect.headers } })
: undefined
const exchange: WebSocketChannelExchange = selected
? { ...input, connect: { ...input.connect, url: selected.url, headers: Headers.fromInput(selected.headers) } }
@@ -361,9 +354,6 @@ export const makeLayer = (connector: WebSocketConnector) =>
Effect.onInterrupt(() => closeChannel(owner, channel)),
)
if (create.mode === "full") channel.checkpoint = undefined
const message = interceptor?.send
? yield* interceptor.send(create.message).pipe(Effect.onInterrupt(() => closeChannel(owner, channel)))
: create.message
yield* Effect.logDebug("session websocket sending", {
sessionTransport: "websocket",
phase: "send",
@@ -374,7 +364,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
delivery: "send-attempted",
}
channel.active = active
const sent = yield* channel.connection.sendText(message).pipe(
const sent = yield* channel.connection.sendText(create.message).pipe(
Effect.withSpan("SessionModelTransport.send"),
Effect.onInterrupt(() => closeChannel(owner, channel)),
Effect.result,
@@ -415,7 +405,6 @@ export const makeLayer = (connector: WebSocketConnector) =>
}),
),
}),
Stream.mapEffect((frame) => (interceptor?.receive ? interceptor.receive(frame) : Effect.succeed(frame))),
Stream.mapEffect((frame) => exchange.driver.observe(create, frame)),
Stream.tap((observation) =>
Effect.sync(() => {
@@ -493,7 +482,10 @@ export const makeLayer = (connector: WebSocketConnector) =>
return { frames, complete, http: channel.connection.http }
})
const bind = (sessionID: SessionSchema.ID, interceptor?: Interceptor): WebSocketChannelExecutor => ({
const bind = (
sessionID: SessionSchema.ID,
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
): WebSocketChannelExecutor => ({
execute: (exchange) => {
const owner = state(sessionID)
let execution: WebSocketChannelExecution | undefined
@@ -503,7 +495,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
},
frames: Stream.unwrap(
Effect.acquireRelease(owner.lock.take(1), () => owner.lock.release(1), { interruptible: true }).pipe(
Effect.andThen(start(owner, exchange, interceptor)),
Effect.andThen(start(owner, exchange, handshake)),
Effect.tap((started) =>
Effect.sync(() => {
execution = started
@@ -80,7 +80,7 @@ describe("SessionModelRequest HTTP hooks", () => {
}).pipe(Effect.provideService(SessionModelTransport.Service, transport)),
)
it.effect("offers the WebSocket executor alongside HTTP hooks and routes the WebSocket hooks", () =>
it.effect("offers the WebSocket executor alongside HTTP hooks and routes the handshake hook", () =>
Effect.gen(function* () {
const hooks = yield* PluginHooks.Service
const seen: string[] = []
@@ -92,31 +92,13 @@ describe("SessionModelRequest HTTP hooks", () => {
delete event.headers["api-key"]
}),
)
yield* hooks.register("session", "experimental.ws.send", (event) =>
Effect.sync(() => {
seen.push(`send:${event.kind}:${event.frame}`)
event.frame = `${event.frame}+plugin`
}),
)
yield* hooks.register("session", "experimental.ws.receive", (event) =>
Effect.sync(() => {
seen.push(`receive:${event.kind}:${event.frame}`)
event.frame = event.frame.toUpperCase()
}),
)
const bound: Array<{ url: string; headers: Record<string, string> }> = []
const frames: string[] = []
const websocketTransport = SessionModelTransport.Service.of({
bind: (_sessionID, interceptor) => ({
bind: (_sessionID, handshake) => ({
execute: () =>
Effect.gen(function* () {
if (!interceptor?.handshake || !interceptor.send || !interceptor.receive)
throw new Error("Expected a full WebSocket interceptor")
bound.push(
yield* interceptor.handshake({ url: "wss://example.test/v1/responses", headers: { "api-key": "k" } }),
)
frames.push(yield* interceptor.send("create"))
frames.push(yield* interceptor.receive("created"))
if (!handshake) throw new Error("Expected a handshake interceptor")
bound.push(yield* handshake({ url: "wss://example.test/v1/responses", headers: { "api-key": "k" } }))
return { frames: Stream.empty, complete: Effect.void }
}),
}),
@@ -145,12 +127,7 @@ describe("SessionModelRequest HTTP hooks", () => {
expect(prepared.options.webSocket).toBeDefined()
yield* prepared.options.webSocket!.execute({} as never)
expect(bound).toEqual([{ url: "wss://example.test/v1/responses", headers: { authorization: "Bearer minted" } }])
expect(frames).toEqual(["create+plugin", "CREATED"])
expect(seen).toEqual([
"handshake:primary:wss://example.test/v1/responses",
"send:primary:create",
"receive:primary:created",
])
expect(seen).toEqual(["handshake:primary:wss://example.test/v1/responses"])
}),
)
})
@@ -178,13 +178,12 @@ describe("SessionModelTransport", () => {
fixture.connector,
Effect.gen(function* () {
const transport = yield* SessionModelTransport.Service
const executor = transport.bind(session, {
handshake: (connect) =>
Effect.succeed({
url: connect.url,
headers: { ...connect.headers, authorization: `Bearer ${tokens.shift()}` },
}),
})
const executor = transport.bind(session, (connect) =>
Effect.succeed({
url: connect.url,
headers: { ...connect.headers, authorization: `Bearer ${tokens.shift()}` },
}),
)
yield* collect(executor, exchange("first", { headers: { "api-key": "k" } }))
yield* collect(executor, exchange("second", { headers: { "api-key": "k" } }))
yield* collect(executor, exchange("third", { headers: { "api-key": "k" } }))
@@ -197,36 +196,6 @@ describe("SessionModelTransport", () => {
)
})
test("sends the frame the send tap returns and observes the frame the receive tap returns", async () => {
const fixture = automatic()
const seen: Array<{ tap: "send" | "receive"; frame: string }> = []
await run(
fixture.connector,
Effect.gen(function* () {
const transport = yield* SessionModelTransport.Service
const executor = transport.bind(session, {
send: (frame) => {
seen.push({ tap: "send", frame })
return Effect.succeed(`${frame}:rewritten`)
},
receive: (frame) => {
seen.push({ tap: "receive", frame })
return Effect.succeed(`${frame}:observed`)
},
})
const frames = yield* collect(executor, exchange("first"))
// The wire carries the rewritten outbound frame; the driver sees the rewritten inbound frame.
expect(fixture.connections.map((item) => item.sent)).toEqual([["first:rewritten"]])
expect(frames).toEqual(["completed:first:rewritten:observed"])
expect(seen).toEqual([
{ tap: "send", frame: "first" },
{ tap: "receive", frame: "completed:first:rewritten" },
])
}),
)
})
test("does not carry a checkpoint across physical connection rotation", async () => {
const fixture = automatic()
const checkpoints: Array<unknown> = []
-27
View File
@@ -99,31 +99,6 @@ export interface SessionWebSocketHandshake {
headers: Record<string, string>
}
/**
* Outbound frame about to be written to the Session's socket, after the provider driver has built
* it. Replacing `frame` sends the replacement verbatim; the driver still tracks state from the
* provider's replies, so a rewrite that changes protocol meaning is on the plugin. Experimental.
*/
export interface SessionWebSocketSend {
readonly sessionID: Session.ID
readonly agent: Agent.ID
readonly model: Model.Ref
readonly kind: SessionRequestKind
frame: string
}
/**
* Inbound frame read from the Session's socket, before the provider driver observes it. Replacing
* `frame` hands the replacement to the driver verbatim. Experimental.
*/
export interface SessionWebSocketReceive {
readonly sessionID: Session.ID
readonly agent: Agent.ID
readonly model: Model.Ref
readonly kind: SessionRequestKind
frame: string
}
export type SessionRetryDecision = { retry: false } | { retry: true; delay: number }
export interface SessionRetry {
@@ -145,8 +120,6 @@ export interface SessionHooks {
readonly "http.request": SessionHttpRequest
readonly "http.response": SessionHttpResponse
readonly "experimental.ws.handshake": SessionWebSocketHandshake
readonly "experimental.ws.send": SessionWebSocketSend
readonly "experimental.ws.receive": SessionWebSocketReceive
readonly retry: SessionRetry
}
-27
View File
@@ -99,31 +99,6 @@ export interface SessionWebSocketHandshake {
headers: Record<string, string>
}
/**
* Outbound frame about to be written to the Session's socket, after the provider driver has built
* it. Replacing `frame` sends the replacement verbatim; the driver still tracks state from the
* provider's replies, so a rewrite that changes protocol meaning is on the plugin. Experimental.
*/
export interface SessionWebSocketSend {
readonly sessionID: Session.ID
readonly agent: Agent.ID
readonly model: Model.Ref
readonly kind: SessionRequestKind
frame: string
}
/**
* Inbound frame read from the Session's socket, before the provider driver observes it. Replacing
* `frame` hands the replacement to the driver verbatim. Experimental.
*/
export interface SessionWebSocketReceive {
readonly sessionID: Session.ID
readonly agent: Agent.ID
readonly model: Model.Ref
readonly kind: SessionRequestKind
frame: string
}
export type SessionRetryDecision = { retry: false } | { retry: true; delay: number }
export interface SessionRetry {
@@ -145,8 +120,6 @@ export interface SessionHooks {
readonly "http.request": SessionHttpRequest
readonly "http.response": SessionHttpResponse
readonly "experimental.ws.handshake": SessionWebSocketHandshake
readonly "experimental.ws.send": SessionWebSocketSend
readonly "experimental.ws.receive": SessionWebSocketReceive
readonly retry: SessionRetry
}
@@ -1277,26 +1277,6 @@ effect: (ctx) =>
}),
```
`experimental.ws.send` and `experimental.ws.receive` expose the frames themselves: `send` runs after the provider
driver builds an outbound frame, `receive` runs on each inbound frame before the driver observes it. Whatever `frame` holds when the hook returns is what crosses the wire or reaches the driver;
OpenCode does not validate it.
```ts
effect: (ctx) =>
Effect.gen(function* () {
yield* ctx.session.hook(
"experimental.ws.send",
(event) =>
Effect.sync(() => {
const body = JSON.parse(event.frame)
if (body.type === "response.create") body.metadata = { ...body.metadata, session: event.sessionID }
event.frame = JSON.stringify(body)
}),
{ providerID: "openai" },
)
}),
```
Override the retry decision for a provider failure or replace its delay in milliseconds. The hook runs after OpenCode
classifies the failure and proposes its policy, but before any retry is scheduled. It does not expose how OpenCode
internally performs the next attempt.
@@ -1337,9 +1317,6 @@ interface SessionHooks {
readonly "model.request": SessionModelRequest
readonly "http.request": SessionHttpRequest
readonly "http.response": SessionHttpResponse
readonly "experimental.ws.handshake": SessionWebSocketHandshake
readonly "experimental.ws.send": SessionWebSocketSend
readonly "experimental.ws.receive": SessionWebSocketReceive
readonly retry: SessionRetry
}
@@ -1410,31 +1410,7 @@ await ctx.session.hook(
)
```
`experimental.ws.send` and `experimental.ws.receive` expose the frames themselves, the WebSocket counterpart of
editing an HTTP request or response body. `send` runs after the provider driver builds an outbound frame and before it
is written; `receive` runs on each inbound frame before the driver observes it. Both carry the frame as a string and
send whatever `frame` holds when the hook returns.
OpenCode does not validate rewritten frames. The driver tracks state from the provider's replies, so a rewrite that
changes protocol meaning is the plugin's responsibility, just as a rewritten HTTP body is.
```ts
await ctx.session.hook(
"experimental.ws.send",
(event) => {
const body = JSON.parse(event.frame)
if (body.type === "response.create") body.metadata = { ...body.metadata, session: event.sessionID }
event.frame = JSON.stringify(body)
},
{ providerID: "openai" },
)
await ctx.session.hook("experimental.ws.receive", (event) => {
if (event.frame.includes('"type":"error"')) console.error(event.frame)
})
```
These hooks are experimental and their names or shapes may change.
This hook is experimental and its name or shape may change.
#### Retry policy
@@ -1482,8 +1458,6 @@ interface SessionHooks {
"http.request": SessionHttpRequestHook
"http.response": SessionHttpResponseHook
"experimental.ws.handshake": SessionWebSocketHandshakeHook
"experimental.ws.send": SessionWebSocketSendHook
"experimental.ws.receive": SessionWebSocketReceiveHook
retry: SessionRetryHook
}
@@ -1496,22 +1470,6 @@ interface SessionWebSocketHandshakeHook {
headers: Record<string, string>
}
interface SessionWebSocketSendHook {
readonly sessionID: string
readonly agent: string
readonly model: { providerID: string; id: string; variant?: string }
readonly kind: "primary" | "compaction" | "title" | "generate"
frame: string
}
interface SessionWebSocketReceiveHook {
readonly sessionID: string
readonly agent: string
readonly model: { providerID: string; id: string; variant?: string }
readonly kind: "primary" | "compaction" | "title" | "generate"
frame: string
}
type RetryDecision = { retry: false } | { retry: true; delay: number }
interface SessionRetryHook {