diff --git a/packages/core/src/session/model-request.ts b/packages/core/src/session/model-request.ts index 8c4970433c0c..9a9292eb6beb 100644 --- a/packages/core/src/session/model-request.ts +++ b/packages/core/src/session/model-request.ts @@ -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) => + 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)), + }) : undefined return { diff --git a/packages/core/src/session/model-transport.ts b/packages/core/src/session/model-transport.ts index eef1a3d5ebdb..590cabf3c157 100644 --- a/packages/core/src/session/model-transport.ts +++ b/packages/core/src/session/model-transport.ts @@ -59,11 +59,18 @@ export interface Handshake { readonly headers: Record } +/** + * 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 + readonly send?: (frame: string) => Effect.Effect + readonly receive?: (frame: string) => Effect.Effect +} + export interface Interface { - readonly bind: ( - sessionID: SessionSchema.ID, - handshake?: (connect: Handshake) => Effect.Effect, - ) => WebSocketChannelExecutor + readonly bind: (sessionID: SessionSchema.ID, interceptor?: Interceptor) => WebSocketChannelExecutor readonly close: (sessionID: SessionSchema.ID) => Effect.Effect readonly closeAll: Effect.Effect } @@ -278,7 +285,7 @@ export const makeLayer = (connector: WebSocketConnector) => const start = Effect.fn("SessionModelTransport.start")(function* ( owner: State, input: WebSocketChannelExchange, - handshake?: (connect: Handshake) => Effect.Effect, + interceptor?: Interceptor, ) { if (owner.closed) return yield* transportError("Session WebSocket owner is closed", { @@ -288,8 +295,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) } } @@ -354,6 +361,9 @@ 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", @@ -364,7 +374,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, @@ -405,6 +415,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(() => { @@ -482,10 +493,7 @@ export const makeLayer = (connector: WebSocketConnector) => return { frames, complete, http: channel.connection.http } }) - const bind = ( - sessionID: SessionSchema.ID, - handshake?: (connect: Handshake) => Effect.Effect, - ): WebSocketChannelExecutor => ({ + const bind = (sessionID: SessionSchema.ID, interceptor?: Interceptor): WebSocketChannelExecutor => ({ execute: (exchange) => { const owner = state(sessionID) let execution: WebSocketChannelExecution | undefined @@ -495,7 +503,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 diff --git a/packages/core/test/session-model-request-hooks.test.ts b/packages/core/test/session-model-request-hooks.test.ts index ae6642ba95a7..3b7b1d45f45b 100644 --- a/packages/core/test/session-model-request-hooks.test.ts +++ b/packages/core/test/session-model-request-hooks.test.ts @@ -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[] = [] @@ -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.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 }> = [] + 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")) + frames.push(yield* interceptor.receive("created")) return { frames: Stream.empty, complete: Effect.void } }), }), @@ -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:create", + "receive:primary:created", + ]) }), ) }) diff --git a/packages/core/test/session-model-transport.test.ts b/packages/core/test/session-model-transport.test.ts index f4a445348402..6933d49957f8 100644 --- a/packages/core/test/session-model-transport.test.ts +++ b/packages/core/test/session-model-transport.test.ts @@ -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" } })) @@ -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 }> = [] + 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 = [] diff --git a/packages/plugin/src/effect/session.ts b/packages/plugin/src/effect/session.ts index 9e125a1c8850..bea12a931283 100644 --- a/packages/plugin/src/effect/session.ts +++ b/packages/plugin/src/effect/session.ts @@ -99,6 +99,31 @@ export interface SessionWebSocketHandshake { headers: Record } +/** + * 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 { @@ -120,6 +145,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 } diff --git a/packages/plugin/src/promise/session.ts b/packages/plugin/src/promise/session.ts index a31872466059..ac3aa3888a5b 100644 --- a/packages/plugin/src/promise/session.ts +++ b/packages/plugin/src/promise/session.ts @@ -99,6 +99,31 @@ export interface SessionWebSocketHandshake { headers: Record } +/** + * 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 { @@ -120,6 +145,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 } diff --git a/services/www/src/docs/content/build/plugins/effect.mdx b/services/www/src/docs/content/build/plugins/effect.mdx index c50a5039ade3..75e0f1ebabfe 100644 --- a/services/www/src/docs/content/build/plugins/effect.mdx +++ b/services/www/src/docs/content/build/plugins/effect.mdx @@ -1277,6 +1277,26 @@ 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. @@ -1317,6 +1337,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 } diff --git a/services/www/src/docs/content/build/plugins/index.mdx b/services/www/src/docs/content/build/plugins/index.mdx index 3237a05cf6f3..279cb1130109 100644 --- a/services/www/src/docs/content/build/plugins/index.mdx +++ b/services/www/src/docs/content/build/plugins/index.mdx @@ -1410,7 +1410,31 @@ await ctx.session.hook( ) ``` -This hook is experimental and its name or shape may change. +`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. #### Retry policy @@ -1458,6 +1482,8 @@ interface SessionHooks { "http.request": SessionHttpRequestHook "http.response": SessionHttpResponseHook "experimental.ws.handshake": SessionWebSocketHandshakeHook + "experimental.ws.send": SessionWebSocketSendHook + "experimental.ws.receive": SessionWebSocketReceiveHook retry: SessionRetryHook } @@ -1470,6 +1496,22 @@ interface SessionWebSocketHandshakeHook { headers: Record } +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 {