Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 18 additions & 9 deletions packages/core/src/session/model-request.ts
Original file line number Diff line number Diff line change
Expand Up @@ -320,15 +320,24 @@ 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, (connect) =>
hooks
.trigger("session", "experimental.ws.handshake", {
...scope,
url: connect.url,
headers: connect.headers,
})
.pipe(Effect.map((event) => ({ url: event.url, headers: event.headers }))),
)
? 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, mode) =>
hooks
.trigger("session", "experimental.ws.send", { ...scope, mode, frame })
.pipe(Effect.map((event) => event.frame)),
receive: (frame) =>
hooks
.trigger("session", "experimental.ws.receive", { ...scope, frame })
.pipe(Effect.map((event) => event.frame)),
})
: undefined

return {
Expand Down
37 changes: 24 additions & 13 deletions packages/core/src/session/model-transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ export * as SessionModelTransport from "./model-transport.js"

import {
WebSocketTransport,
type ChannelCreate,
type ChannelObservation,
type ChannelCheckpoint,
type WebSocketChannelExchange,
Expand Down Expand Up @@ -59,11 +60,18 @@ 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, mode: ChannelCreate["mode"]) => Effect.Effect<string>
readonly receive?: (frame: string) => Effect.Effect<string>
}

export interface Interface {
readonly bind: (
sessionID: SessionSchema.ID,
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
) => WebSocketChannelExecutor
readonly bind: (sessionID: SessionSchema.ID, interceptor?: Interceptor) => WebSocketChannelExecutor
readonly close: (sessionID: SessionSchema.ID) => Effect.Effect<void>
readonly closeAll: Effect.Effect<void>
}
Expand Down Expand Up @@ -278,7 +286,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
const start = Effect.fn("SessionModelTransport.start")(function* (
owner: State,
input: WebSocketChannelExchange,
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
interceptor?: Interceptor,
) {
if (owner.closed)
return yield* transportError("Session WebSocket owner is closed", {
Expand All @@ -288,8 +296,8 @@ export const makeLayer = (connector: WebSocketConnector) =>
delivery: "not-sent",
})
if (owner.httpFallback) return fallback(input)
const selected = handshake
? yield* handshake({ url: input.connect.url, headers: { ...input.connect.headers } })
const selected = interceptor?.handshake
? yield* interceptor.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) } }
Expand Down Expand Up @@ -354,6 +362,11 @@ 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, create.mode)
.pipe(Effect.onInterrupt(() => closeChannel(owner, channel)))
: create.message
yield* Effect.logDebug("session websocket sending", {
sessionTransport: "websocket",
phase: "send",
Expand All @@ -364,7 +377,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
delivery: "send-attempted",
}
channel.active = active
const sent = yield* channel.connection.sendText(create.message).pipe(
const sent = yield* channel.connection.sendText(message).pipe(
Effect.withSpan("SessionModelTransport.send"),
Effect.onInterrupt(() => closeChannel(owner, channel)),
Effect.result,
Expand Down Expand Up @@ -405,6 +418,7 @@ 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(() => {
Expand Down Expand Up @@ -482,10 +496,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
return { frames, complete, http: channel.connection.http }
})

const bind = (
sessionID: SessionSchema.ID,
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
): WebSocketChannelExecutor => ({
const bind = (sessionID: SessionSchema.ID, interceptor?: Interceptor): WebSocketChannelExecutor => ({
execute: (exchange) => {
const owner = state(sessionID)
let execution: WebSocketChannelExecution | undefined
Expand All @@ -495,7 +506,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, handshake)),
Effect.andThen(start(owner, exchange, interceptor)),
Effect.tap((started) =>
Effect.sync(() => {
execution = started
Expand Down
33 changes: 28 additions & 5 deletions packages/core/test/session-model-request-hooks.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 handshake hook", () =>
it.effect("offers the WebSocket executor alongside HTTP hooks and routes the WebSocket hooks", () =>
Effect.gen(function* () {
const hooks = yield* PluginHooks.Service
const seen: string[] = []
Expand All @@ -92,13 +92,31 @@ 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.mode}:${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, handshake) => ({
bind: (_sessionID, interceptor) => ({
execute: () =>
Effect.gen(function* () {
if (!handshake) throw new Error("Expected a handshake interceptor")
bound.push(yield* handshake({ url: "wss://example.test/v1/responses", headers: { "api-key": "k" } }))
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", "incremental"))
frames.push(yield* interceptor.receive("created"))
return { frames: Stream.empty, complete: Effect.void }
}),
}),
Expand Down Expand Up @@ -127,7 +145,12 @@ 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(seen).toEqual(["handshake:primary:wss://example.test/v1/responses"])
expect(frames).toEqual(["create+plugin", "CREATED"])
expect(seen).toEqual([
"handshake:primary:wss://example.test/v1/responses",
"send:primary:incremental:create",
"receive:primary:created",
])
}),
)
})
43 changes: 37 additions & 6 deletions packages/core/test/session-model-transport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -178,12 +178,13 @@ describe("SessionModelTransport", () => {
fixture.connector,
Effect.gen(function* () {
const transport = yield* SessionModelTransport.Service
const executor = transport.bind(session, (connect) =>
Effect.succeed({
url: connect.url,
headers: { ...connect.headers, authorization: `Bearer ${tokens.shift()}` },
}),
)
const executor = transport.bind(session, {
handshake: (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" } }))
Expand All @@ -196,6 +197,36 @@ 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; mode?: string }> = []
await run(
fixture.connector,
Effect.gen(function* () {
const transport = yield* SessionModelTransport.Service
const executor = transport.bind(session, {
send: (frame, mode) => {
seen.push({ tap: "send", frame, mode })
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", mode: "full" },
{ tap: "receive", frame: "completed:first:rewritten" },
])
}),
)
})

test("does not carry a checkpoint across physical connection rotation", async () => {
const fixture = automatic()
const checkpoints: Array<unknown> = []
Expand Down
29 changes: 29 additions & 0 deletions packages/plugin/src/effect/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,33 @@ 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
/** `"incremental"` frames carry only what changed since the provider's last checkpoint. */
readonly mode: "full" | "incremental"
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 {
Expand All @@ -120,6 +147,8 @@ 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
}

Expand Down
29 changes: 29 additions & 0 deletions packages/plugin/src/promise/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,33 @@ 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
/** `"incremental"` frames carry only what changed since the provider's last checkpoint. */
readonly mode: "full" | "incremental"
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 {
Expand All @@ -120,6 +147,8 @@ 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
}

Expand Down
24 changes: 24 additions & 0 deletions services/www/src/docs/content/build/plugins/effect.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -1277,6 +1277,27 @@ effect: (ctx) =>
}),
```

`experimental.ws.send` and `experimental.ws.receive` expose the frames themselves: `send` runs after the provider
driver builds an outbound frame (`mode` is `"full"` or `"incremental"`), `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.
Expand Down Expand Up @@ -1317,6 +1338,9 @@ 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
}

Expand Down
Loading
Loading