mirror of
https://github.com/anomalyco/opencode.git
synced 2026-09-17 14:26:26 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
93eca9b573 |
@@ -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" }>,
|
||||
|
||||
@@ -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")
|
||||
}),
|
||||
)
|
||||
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -100,6 +100,7 @@ export const layer = (options?: Options) =>
|
||||
const recoverShell = Effect.fnUntraced(function* (
|
||||
background: Job.Background,
|
||||
recovery: Extract<Job.Recovery, { kind: "shell" }>,
|
||||
suspended: ReadonlySet<SessionSchema.ID>,
|
||||
) {
|
||||
const state = background.status === "running" ? "cancelled" : background.status
|
||||
const text =
|
||||
@@ -123,9 +124,7 @@ export const layer = (options?: Options) =>
|
||||
state,
|
||||
text,
|
||||
}),
|
||||
// Restart notices must not revive idle owners of long-lived shells.
|
||||
// Interrupted executions resume separately after their notices are admitted.
|
||||
resume: false,
|
||||
...(suspended.has(recovery.sessionID) ? { resume: false } : {}),
|
||||
})
|
||||
.pipe(
|
||||
Effect.catchTag("Session.NotFoundError", () => Effect.void),
|
||||
@@ -209,7 +208,7 @@ export const layer = (options?: Options) =>
|
||||
if ((yield* jobs.get(background.id))?.status === "running") return
|
||||
const recovery = background.recovery
|
||||
yield* recovery.kind === "shell"
|
||||
? recoverShell(background, recovery)
|
||||
? recoverShell(background, recovery, suspended)
|
||||
: recoverSubagent(background, recovery, suspended)
|
||||
}),
|
||||
{ discard: true },
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -379,7 +379,7 @@ describe("SessionExecution lifecycle", () => {
|
||||
})
|
||||
|
||||
describe("SessionRestart background recovery", () => {
|
||||
it.effect("keeps shell owners idle until a user prompt delivers recovered notices exactly once", () =>
|
||||
it.effect("wakes idle shell owners and delivers recovered notices exactly once", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const store = yield* SessionStore.Service
|
||||
@@ -416,20 +416,6 @@ describe("SessionRestart background recovery", () => {
|
||||
yield* restart.resumeSuspendedSessions
|
||||
yield* Effect.forEach([parent, child], execution.awaitIdle, { discard: true })
|
||||
|
||||
expect(drained).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, parent)).toHaveLength(1)
|
||||
expect(yield* SessionInbox.list(database.db, child)).toHaveLength(1)
|
||||
expect(yield* restarted.pendingBackground).toEqual([])
|
||||
yield* restart.resumeSuspendedSessions
|
||||
expect(drained).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, parent)).toHaveLength(1)
|
||||
expect(yield* SessionInbox.list(database.db, child)).toHaveLength(1)
|
||||
|
||||
yield* seedInbox(database, parent, ["steer"])
|
||||
yield* seedInbox(database, child, ["steer"])
|
||||
yield* execution.wake(parent)
|
||||
yield* execution.wake(child)
|
||||
yield* Effect.forEach([parent, child], execution.awaitIdle, { discard: true })
|
||||
expect(drained.toSorted()).toEqual([parent, child].toSorted())
|
||||
expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toMatchObject([
|
||||
{
|
||||
@@ -523,7 +509,7 @@ describe("SessionRestart background recovery", () => {
|
||||
yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions
|
||||
yield* Context.get(context, SessionExecution.Service).awaitIdle(sessionID)
|
||||
|
||||
expect(drained).toEqual([])
|
||||
expect(drained).toEqual([sessionID])
|
||||
const inbox = yield* SessionInbox.list(database.db, sessionID)
|
||||
expect(inbox).toMatchObject([
|
||||
{
|
||||
@@ -575,6 +561,7 @@ describe("SessionRestart background recovery", () => {
|
||||
expect(yield* restarted.pendingBackground).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, sessionID)).toHaveLength(delivered ? 0 : 1)
|
||||
yield* SessionInbox.promote(database.db, bus, sessionID, "steer")
|
||||
// Recovery ends a busy period, so an idle marker follows the notification.
|
||||
const messages = (yield* sessions.messages({ sessionID })).filter((message) => message.type !== "idle")
|
||||
expect(messages).toMatchObject([
|
||||
{
|
||||
|
||||
@@ -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> = []
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -21,6 +21,8 @@ export function Spinner(props: { children?: JSX.Element; color?: RGBA; shimmer?:
|
||||
const timer = setInterval(() => setFrame((value) => (value + 1) % SPINNER_FRAMES.length), 80)
|
||||
onCleanup(() => clearInterval(timer))
|
||||
})
|
||||
// A bare spinner beside a long sibling shrinks to a fractional width that rounds to 0, so the
|
||||
// glyph paints on top of the row gap. Only a labeled spinner may shrink (to wrap its text).
|
||||
return (
|
||||
<Show
|
||||
when={config.animations ?? true}
|
||||
@@ -29,8 +31,10 @@ export function Spinner(props: { children?: JSX.Element; color?: RGBA; shimmer?:
|
||||
<Show
|
||||
when={props.shimmer}
|
||||
fallback={
|
||||
<box flexDirection="row" gap={1}>
|
||||
<spinner frames={SPINNER_FRAMES} interval={80} color={color()} />
|
||||
<box flexDirection="row" gap={1} flexShrink={props.children ? 1 : 0}>
|
||||
<box flexShrink={0}>
|
||||
<spinner frames={SPINNER_FRAMES} interval={80} color={color()} />
|
||||
</box>
|
||||
<Show when={props.children}>
|
||||
<text fg={color()}>{props.children}</text>
|
||||
</Show>
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user