Преглед на файлове

feat(plugin): add session HTTP middleware (#40327)

Aiden Cline преди 1 седмица
родител
ревизия
aa820db18d

+ 3 - 3
packages/ai/src/route/client.ts

@@ -5,7 +5,7 @@ import { Endpoint, type EndpointPatch } from "./endpoint"
 import { RequestExecutor } from "./executor"
 import { Framing } from "./framing"
 import { HttpTransport } from "./transport"
-import type { HttpRequestTransform, Transport, TransportRuntime } from "./transport"
+import type { HttpMiddleware, Transport, TransportRuntime } from "./transport"
 import { WebSocketExecutor } from "./transport"
 import type { Protocol } from "./protocol"
 import { applyCachePolicy } from "../cache-policy"
@@ -155,7 +155,7 @@ export interface Interface {
 }
 
 export interface StreamOptions {
-  readonly transform?: HttpRequestTransform
+  readonly http?: HttpMiddleware
 }
 
 export interface StreamMethod {
@@ -307,7 +307,7 @@ function makeFromTransport<Body, Prepared, Frame, Event, State>(
           auth: routeInput.auth ?? Auth.none,
           encodeBody,
           headers: routeInput.headers,
-          transform: options?.transform,
+          middleware: options?.http,
         }),
       streamPrepared: (prepared: Prepared, request: LLMRequest, runtime: TransportRuntime) => {
         const route = `${request.model.provider}/${request.model.route.id}`

+ 22 - 5
packages/ai/src/route/executor.ts

@@ -20,9 +20,18 @@ import { classifyProviderFailure } from "../provider-error"
 export interface Interface {
   readonly execute: (
     request: HttpClientRequest.HttpClientRequest,
+    middleware?: HttpMiddleware,
   ) => Effect.Effect<HttpClientResponse.HttpClientResponse, AIError>
 }
 
+export type HttpHandler = (
+  request: HttpClientRequest.HttpClientRequest,
+) => Effect.Effect<HttpClientResponse.HttpClientResponse, Error>
+export type HttpMiddleware = (
+  request: HttpClientRequest.HttpClientRequest,
+  handler: HttpHandler,
+) => Effect.Effect<HttpClientResponse.HttpClientResponse, Error>
+
 export class Service extends Context.Service<Service, Interface>()("@opencode/AI/RequestExecutor") {}
 
 const BODY_LIMIT = 16_384
@@ -261,7 +270,7 @@ const toHttpError = (redactedNames: ReadonlyArray<string | RegExp>) => (error: u
     return transportError({ message: error.message, kind: "Timeout" })
   }
   if (!HttpClientError.isHttpClientError(error)) {
-    return transportError({ message: "HTTP transport failed" })
+    return transportError({ message: error instanceof Error ? error.message : "HTTP transport failed" })
   }
   const request = "request" in error ? error.request : undefined
   if (error.reason._tag === "TransportError") {
@@ -282,12 +291,20 @@ export const layer: Layer.Layer<Service, never, HttpClient.HttpClient> = Layer.e
   Service,
   Effect.gen(function* () {
     const http = yield* HttpClient.HttpClient
-    const executeOnce = (request: HttpClientRequest.HttpClientRequest) =>
+    const executeOnce = (request: HttpClientRequest.HttpClientRequest, middleware?: HttpMiddleware) =>
       Effect.gen(function* () {
         const redactedNames = yield* Headers.CurrentRedactedNames
-        return yield* http
-          .execute(request)
-          .pipe(Effect.mapError(toHttpError(redactedNames)), Effect.flatMap(statusError(request, redactedNames)))
+        if (!middleware)
+          return yield* http
+            .execute(request)
+            .pipe(Effect.mapError(toHttpError(redactedNames)), Effect.flatMap(statusError(request, redactedNames)))
+
+        const response = yield* middleware(request, (input) =>
+          http
+            .execute(input)
+            .pipe(Effect.mapError((cause) => (cause instanceof Error ? cause : new Error(String(cause))))),
+        ).pipe(Effect.mapError(toHttpError(redactedNames)))
+        return yield* statusError(response.request, redactedNames)(response)
       })
     return Service.of({
       execute: executeOnce,

+ 1 - 1
packages/ai/src/route/index.ts

@@ -23,4 +23,4 @@ export type { ApiKeyMode, AuthOverride, ProviderAuthOption } from "./auth-option
 export type { Definition as EndpointFn, EndpointInput } from "./endpoint"
 export type { Definition as FramingDef } from "./framing"
 export type { Protocol as ProtocolDef } from "./protocol"
-export type { HttpRequest, HttpRequestTransform, Transport as TransportDef, TransportRuntime } from "./transport"
+export type { HttpHandler, HttpMiddleware, Transport as TransportDef, TransportRuntime } from "./transport"

+ 10 - 9
packages/ai/src/route/transport/http.ts

@@ -3,7 +3,7 @@ import { Headers, HttpClientRequest } from "effect/unstable/http"
 import { Auth } from "../auth"
 import { render as renderEndpoint } from "../endpoint"
 import { Framing } from "../framing"
-import type { Transport, TransportPrepareInput } from "./index"
+import type { HttpMiddleware, Transport, TransportPrepareInput } from "./index"
 import * as ProviderShared from "../../protocols/shared"
 import { mergeJsonRecords, type LLMRequest } from "../../schema"
 
@@ -19,6 +19,7 @@ export interface JsonRequestParts<Body = unknown> {
 export interface HttpPrepared<Frame> {
   readonly request: HttpClientRequest.HttpClientRequest
   readonly framing: Framing.Definition<Frame>
+  readonly middleware?: HttpMiddleware
 }
 
 const applyQuery = (url: string, query: Record<string, string> | undefined) => {
@@ -74,21 +75,21 @@ export const httpJson = <Body, Frame>(input: HttpJsonInput<Body, Frame>): HttpJs
   prepare: (prepareInput) =>
     Effect.gen(function* () {
       const parts = yield* jsonRequestParts({ ...prepareInput })
-      const request = { url: parts.url, method: "POST", headers: { ...parts.headers }, body: parts.bodyText }
-      yield* (prepareInput.transform?.(request) ?? Effect.void)
+      const request = ProviderShared.jsonPost({
+        url: parts.url,
+        body: parts.bodyText,
+        headers: parts.headers,
+      })
       return {
-        request: ProviderShared.jsonPost({
-          url: request.url,
-          body: request.body ?? "",
-          headers: Headers.fromInput(request.headers),
-        }),
+        request,
         framing: input.framing,
+        middleware: prepareInput.middleware,
       }
     }),
   frames: (prepared, request, runtime) =>
     Stream.unwrap(
       runtime.http
-        .execute(prepared.request)
+        .execute(prepared.request, prepared.middleware)
         .pipe(
           Effect.map((response) =>
             prepared.framing.frame(

+ 3 - 11
packages/ai/src/route/transport/index.ts

@@ -1,7 +1,7 @@
 import type { Effect, Stream } from "effect"
 import { Endpoint } from "../endpoint"
 import { Auth } from "../auth"
-import type { Interface as RequestExecutorInterface } from "../executor"
+import type { HttpMiddleware, Interface as RequestExecutorInterface } from "../executor"
 import type { Interface as WebSocketExecutorInterface } from "./websocket"
 import type { AIError, LLMRequest } from "../../schema"
 
@@ -10,15 +10,6 @@ export interface TransportRuntime {
   readonly webSocket?: WebSocketExecutorInterface
 }
 
-export interface HttpRequest {
-  url: string
-  readonly method: string
-  headers: Record<string, string>
-  body: string | undefined
-}
-
-export type HttpRequestTransform = (request: HttpRequest) => Effect.Effect<void>
-
 export interface Transport<Body, Prepared, Frame> {
   readonly id: string
   readonly prepare: (input: TransportPrepareInput<Body>) => Effect.Effect<Prepared, AIError>
@@ -32,8 +23,9 @@ export interface TransportPrepareInput<Body> {
   readonly auth: Auth.Definition
   readonly encodeBody: (body: Body) => string
   readonly headers?: (input: { readonly request: LLMRequest }) => Record<string, string>
-  readonly transform?: HttpRequestTransform
+  readonly middleware?: HttpMiddleware
 }
 
 export * as HttpTransport from "./http"
+export type { HttpHandler, HttpMiddleware } from "../executor"
 export { WebSocketExecutor, WebSocketTransport } from "./websocket"

+ 90 - 8
packages/ai/test/compile.test.ts

@@ -1,6 +1,6 @@
 import { describe, expect, test } from "bun:test"
-import { Effect, Schema } from "effect"
-import { HttpClientRequest } from "effect/unstable/http"
+import { Effect, Ref, Schema } from "effect"
+import { HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
 import { LLM, mergeProviderOptions } from "../src"
 import { AnthropicMessages, OpenAIChat } from "../src/protocols"
 import { Auth, LLMClient } from "../src/route"
@@ -146,12 +146,16 @@ describe("request option precedence", () => {
         prompt: "Say hello.",
       }),
       {
-        transform: (request) =>
-          Effect.sync(() => {
-            expect(request.headers.authorization).toBe("Bearer fresh-key")
-            request.url = "https://proxy.test/v1/chat/completions"
-            request.headers["x-plugin"] = "transformed"
-            request.body = JSON.stringify({ transformed: true })
+        http: (request, handler) =>
+          Effect.gen(function* () {
+            return yield* handler(
+              request.pipe(
+                HttpClientRequest.setUrl("https://proxy.test/v1/chat/completions"),
+                HttpClientRequest.setMethod("PUT"),
+                HttpClientRequest.setHeader("x-plugin", "transformed"),
+                HttpClientRequest.bodyText(JSON.stringify({ transformed: true }), "application/custom+json"),
+              ),
+            )
           }),
       },
     ).pipe(
@@ -160,7 +164,9 @@ describe("request option precedence", () => {
           Effect.gen(function* () {
             const web = yield* HttpClientRequest.toWeb(input.request).pipe(Effect.orDie)
             expect(web.url).toBe("https://proxy.test/v1/chat/completions")
+            expect(web.method).toBe("PUT")
             expect(web.headers.get("x-plugin")).toBe("transformed")
+            expect(web.headers.get("content-type")).toBe("application/custom+json")
             expect(decodeJson(input.text)).toEqual({ transformed: true })
             return input.respond(sseEvents(deltaChunk({}, "stop")), {
               headers: { "content-type": "text/event-stream" },
@@ -171,6 +177,82 @@ describe("request option precedence", () => {
     ),
   )
 
+  it.effect("transforms the HTTP response before protocol decoding", () =>
+    Effect.gen(function* () {
+      const response = yield* LLMClient.generate(
+        LLM.request({
+          model: OpenAIChat.route
+            .with({ endpoint: { baseURL: "https://api.openai.test/v1/" }, auth: Auth.bearer("test") })
+            .model({ id: "gpt-4o-mini" }),
+          prompt: "Say hello.",
+        }),
+        {
+          http: (request, handler) =>
+            Effect.gen(function* () {
+              const response = yield* handler(request)
+              return HttpClientResponse.fromWeb(
+                response.request,
+                new Response((yield* response.text).replace("network", "hooked"), {
+                  status: response.status,
+                  headers: response.headers,
+                }),
+              )
+            }),
+        },
+      ).pipe(
+        Effect.provide(
+          dynamicResponse((input) =>
+            Effect.succeed(
+              input.respond(sseEvents(deltaChunk({ content: "network" }, "stop")), {
+                headers: { "content-type": "text/event-stream" },
+              }),
+            ),
+          ),
+        ),
+      )
+
+      expect(response.text).toBe("hooked")
+    }),
+  )
+
+  it.effect("can inspect an error response and retry the native request", () =>
+    Effect.gen(function* () {
+      const attempts = yield* Ref.make(0)
+      const response = yield* LLMClient.generate(
+        LLM.request({
+          model: OpenAIChat.route
+            .with({ endpoint: { baseURL: "https://api.openai.test/v1/" }, auth: Auth.bearer("stale") })
+            .model({ id: "gpt-4o-mini" }),
+          prompt: "Say hello.",
+        }),
+        {
+          http: (request, handler) =>
+            Effect.gen(function* () {
+              const response = yield* handler(request)
+              expect(response.status).toBe(401)
+              return yield* handler(HttpClientRequest.setHeader(request, "authorization", "Bearer refreshed"))
+            }),
+        },
+      ).pipe(
+        Effect.provide(
+          dynamicResponse((input) =>
+            Effect.gen(function* () {
+              yield* Ref.update(attempts, (value) => value + 1)
+              if (input.request.headers.authorization !== "Bearer refreshed")
+                return input.respond("unauthorized", { status: 401 })
+              return input.respond(sseEvents(deltaChunk({ content: "retried" }, "stop")), {
+                headers: { "content-type": "text/event-stream" },
+              })
+            }),
+          ),
+        ),
+      )
+
+      expect(response.text).toBe("retried")
+      expect(yield* Ref.get(attempts)).toBe(2)
+    }),
+  )
+
   it.effect("applies raw body overlays after protocol lowering", () =>
     LLMClient.generate(
       LLM.request({

+ 12 - 0
packages/ai/test/executor.test.ts

@@ -67,6 +67,18 @@ const expectAIError = (error: unknown) => {
 const errorHttp = (error: AIError) => ("http" in error.reason ? error.reason.http : undefined)
 
 describe("RequestExecutor", () => {
+  it.effect("preserves middleware error messages", () =>
+    Effect.gen(function* () {
+      const executor = yield* RequestExecutor.Service
+      const error = yield* executor
+        .execute(request, () => Effect.fail(new Error("plugin rejected request")))
+        .pipe(Effect.flip)
+
+      expectAIError(error)
+      expect(error.reason.message).toBe("plugin rejected request")
+    }).pipe(Effect.provide(responsesLayer([]))),
+  )
+
   it.effect("classifies context overflow responses", () =>
     Effect.gen(function* () {
       const executor = yield* RequestExecutor.Service

+ 61 - 3
packages/core/src/plugin/promise.ts

@@ -2,6 +2,7 @@ export * as PluginPromise from "./promise"
 
 import { define } from "@opencode-ai/plugin/effect/plugin"
 import type { Context, Plugin } from "@opencode-ai/plugin/promise/plugin"
+import type { SessionHooks, SessionHttp, SessionHttpMiddleware } from "@opencode-ai/plugin/promise/session"
 import type { Info } from "@opencode-ai/plugin/promise/tool"
 import { Agent } from "@opencode-ai/schema/agent"
 import { Integration } from "@opencode-ai/schema/integration"
@@ -57,6 +58,62 @@ export function fromPromise(plugin: Plugin) {
               }),
             )
 
+        function sessionHook<Name extends keyof SessionHooks>(
+          name: Name,
+          callback: (event: SessionHooks[Name]) => Promise<void> | void,
+        ): Promise<Registration>
+        function sessionHook(
+          ...registration: {
+            [Name in keyof SessionHooks]: [
+              name: Name,
+              callback: (event: SessionHooks[Name]) => Promise<void> | void,
+            ]
+          }[keyof SessionHooks]
+        ) {
+          if (registration[0] !== "http")
+            return register(
+              host.session.hook(registration[0], (event) =>
+                Effect.promise(() => Promise.resolve(registration[1](event))),
+              ),
+            )
+          return register(
+            host.session.hook("http", (event) => {
+              const middlewares: SessionHttpMiddleware[] = []
+              const output: SessionHttp = {
+                ...event,
+                use: (item) => {
+                  middlewares.push(item)
+                },
+              }
+              return Effect.promise(() => Promise.resolve(registration[1](output))).pipe(
+                Effect.flatMap(() =>
+                  Effect.forEach(
+                    middlewares,
+                    (item) =>
+                      event.use((input, next) =>
+                        Effect.tryPromise({
+                          try: (signal) => {
+                            const inputSignal = AbortSignal.any([signal, input.signal])
+                            return Promise.resolve(
+                              item(new Request(input, { signal: inputSignal }), (request) => {
+                                const requestSignal = AbortSignal.any([signal, request.signal])
+                                return Effect.runPromiseWith(
+                                  context,
+                                )(next(new Request(request, { signal: requestSignal })), { signal: requestSignal })
+                              }),
+                            )
+                          },
+                          catch: (cause) => (cause instanceof Error ? cause : new Error(String(cause))),
+                        }),
+                      ),
+                    { discard: true },
+                  ),
+                ),
+              )
+            }),
+          )
+        }
+
         const context2: Context = {
           app: host.app,
           options: host.options,
@@ -194,7 +251,9 @@ export function fromPromise(plugin: Plugin) {
                               ),
                             ),
                           refresh:
-                            refresh === undefined ? undefined : (credential) => Effect.promise(() => refresh(credential)),
+                            refresh === undefined
+                              ? undefined
+                              : (credential) => Effect.promise(() => refresh(credential)),
                         })
                       },
                       remove: draft.method.remove,
@@ -263,8 +322,7 @@ export function fromPromise(plugin: Plugin) {
               ),
           },
           session: {
-            hook: (name, callback) =>
-              register(host.session.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
+            hook: sessionHook,
             create: (input) =>
               run(
                 host.session.create(

+ 8 - 9
packages/core/src/plugin/provider/openai.ts

@@ -225,15 +225,14 @@ export const OpenAIPlugin = define({
         })
       }
     })
-    yield* ctx.session.hook("request", (evt) =>
-      Effect.sync(() => {
-        if (!chatgpt || evt.model.providerID !== Provider.ID.openai) return
-        const url = new URL(evt.url)
-        if (url.origin === "https://api.openai.com") {
-          evt.url = `${codexBaseURL}${url.pathname.replace(/^\/v1/, "")}${url.search}`
-        }
-        evt.headers.originator = "opencode"
-        evt.headers["session-id"] = evt.sessionID
+    yield* ctx.session.hook("http", (evt) =>
+      evt.use((request, next) => {
+        if (!chatgpt || evt.model.providerID !== Provider.ID.openai) return next(request)
+        const url = new URL(request.url)
+        request.headers.set("originator", "opencode")
+        request.headers.set("session-id", evt.sessionID)
+        if (url.origin !== "https://api.openai.com") return next(request)
+        return next(new Request(`${codexBaseURL}${url.pathname.replace(/^\/v1/, "")}${url.search}`, request))
       }),
     )
 

+ 41 - 19
packages/core/src/session/model-request.ts

@@ -2,9 +2,11 @@ export * as SessionModelRequest from "./model-request"
 
 import { LLM, Message, SystemPart, type LLMRequest } from "@opencode-ai/ai"
 import type { StreamOptions } from "@opencode-ai/ai/route"
+import type { SessionHttpHandler, SessionHttpMiddleware } from "@opencode-ai/plugin/effect/session"
 import type { Content } from "@opencode-ai/schema/tool"
 import { SessionError } from "@opencode-ai/schema/session-error"
-import { Cause, Config, Context, Effect, Layer, Result } from "effect"
+import { Cause, Config, Context, Effect, Layer, Result, Stream } from "effect"
+import { HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
 import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
 import { App } from "../app"
 import { Model } from "../model"
@@ -48,9 +50,7 @@ interface Prepared {
    * One request-scoped execution operation. Unknown, hook-removed, and
    * step-limit-violating calls fail individually through the same seam.
    */
-  readonly executeTool: (
-    input: Parameters<Tool.Snapshot["execute"]>[0],
-  ) => Effect.Effect<Tool.Result, ExecuteError>
+  readonly executeTool: (input: Parameters<Tool.Snapshot["execute"]>[0]) => Effect.Effect<Tool.Result, ExecuteError>
   /** True when this request is the final Step; violating calls are rejected and no continuation follows. */
   readonly stepLimitReached: boolean
 }
@@ -137,8 +137,7 @@ export const boundImages = (messages: LLMRequest["messages"]) => {
           result: {
             ...part.result,
             value: part.result.value.map((item: Content) => {
-              if (item.type !== "file" || !isImage(item.mime) || imageBytes - removed <= IMAGE_BYTES_TARGET)
-                return item
+              if (item.type !== "file" || !isImage(item.mime) || imageBytes - removed <= IMAGE_BYTES_TARGET) return item
               removed += Buffer.byteLength(item.uri)
               return { type: "text" as const, text: IMAGE_REMOVED }
             }),
@@ -229,24 +228,47 @@ export const layer = Layer.effect(
         toolChoice: stepLimitReached ? "none" : undefined,
       })
       const options: StreamOptions = {
-        transform: (request) =>
-          hooks
-            .trigger("session", "request", {
+        http: (request, handler) =>
+          Effect.gen(function* () {
+            let latest = request
+            const origins = new WeakMap<Response, HttpClientRequest.HttpClientRequest>()
+            const middlewares: SessionHttpMiddleware[] = []
+            const web = yield* HttpClientRequest.toWeb(request)
+            yield* hooks.trigger("session", "http", {
               sessionID: session.id,
               agent: agent.id,
               model: resolved.ref,
-              ...request,
-            })
-            .pipe(
-              Effect.tap((event) =>
+              use: (item) =>
                 Effect.sync(() => {
-                  request.url = event.url
-                  request.headers = event.headers
-                  request.body = event.body
+                  middlewares.push(item)
                 }),
-              ),
-              Effect.asVoid,
-            ),
+            })
+            const send = (input: Request) =>
+              Effect.gen(function* () {
+                let sent = HttpClientRequest.fromWeb(input)
+                if (input.body)
+                  sent = HttpClientRequest.bodyUint8Array(
+                    sent,
+                    new Uint8Array(yield* Effect.promise(() => input.clone().arrayBuffer())),
+                    input.headers.get("content-type") ?? undefined,
+                  )
+                latest = sent
+                const response = yield* handler(sent)
+                const body = [204, 205, 304].includes(response.status)
+                  ? null
+                  : yield* Stream.toReadableStreamEffect(response.stream)
+                const output = new Response(body, { status: response.status, headers: response.headers })
+                origins.set(output, sent)
+                return output
+              })
+            const dispatch = middlewares.reduce<SessionHttpHandler>(
+              (next, item) => (input: Request) => item(input, next),
+              send,
+            )
+            const response = yield* dispatch(web)
+            const origin = origins.get(response) ?? latest
+            return HttpClientResponse.fromWeb(origin, response)
+          }).pipe(Effect.mapError((cause) => (cause instanceof Error ? cause : new Error(String(cause))))),
       }
       if (promptCacheSnapshots) {
         const current = PromptCacheDiagnostics.snapshot(request)

+ 114 - 15
packages/core/test/plugin/promise.test.ts

@@ -1,6 +1,6 @@
 import { describe, expect } from "bun:test"
 import { Message, SystemPart } from "@opencode-ai/ai"
-import { DateTime, Effect, Schema } from "effect"
+import { DateTime, Deferred, Effect, Fiber, Schema } from "effect"
 import { Agent } from "@opencode-ai/core/agent"
 import { Catalog } from "@opencode-ai/core/catalog"
 import { Model } from "@opencode-ai/core/model"
@@ -15,7 +15,7 @@ import { SessionPending } from "@opencode-ai/core/session/pending"
 import { Tool } from "@opencode-ai/core/tool"
 import { Provider } from "@opencode-ai/core/provider"
 import { define } from "@opencode-ai/plugin/promise/plugin"
-import type { SessionHooks } from "@opencode-ai/plugin/effect/session"
+import type { SessionHooks, SessionHttpHandler } from "@opencode-ai/plugin/effect/session"
 import { testEffect } from "../lib/effect"
 import { PluginTestLayer } from "./fixture"
 import { host as testHost } from "./host"
@@ -148,7 +148,9 @@ describe("fromPromise", () => {
             expect((await ctx.agent.get({ agentID: Agent.ID.make("reviewer") })).data).toMatchObject({
               description: "Reviews code",
             })
-            await expect(ctx.agent.get({ agentID: Agent.ID.make("missing") })).rejects.toThrow("Agent not found: missing")
+            await expect(ctx.agent.get({ agentID: Agent.ID.make("missing") })).rejects.toThrow(
+              "Agent not found: missing",
+            )
             const models = (await ctx.catalog.model.list()).data
             expect(models.find((model) => model.providerID === "test" && model.id === "alias")).toMatchObject({
               modelID: "gpt-5",
@@ -221,6 +223,105 @@ describe("fromPromise", () => {
     }),
   )
 
+  it.effect("adapts promise session HTTP hooks", () =>
+    Effect.gen(function* () {
+      const plugin = yield* Plugin.Service
+      const hooks = yield* PluginHooks.Service
+      const host = yield* PluginHost.make(plugin)
+      const bodies: string[] = []
+      yield* PluginPromise.fromPromise(
+        define({
+          id: "promise-session-http",
+          setup: async (ctx) => {
+            await ctx.session.hook("http", (event) => {
+              event.use(async (request, next) => {
+                request.headers.set("x-hook", "promise")
+                await next(request)
+                const response = await next(request)
+                return new Response(`${await response.text()}-response`)
+              })
+            })
+            await ctx.session.hook("http", (event) => {
+              event.use(async (request, next) => {
+                const response = await next(request)
+                return new Response(`${await response.text()}-outer`)
+              })
+            })
+          },
+        }),
+      ).effect(host)
+      const middlewares: Parameters<PluginHooks.Domains["session"]["http"]["use"]>[0][] = []
+      const event: PluginHooks.Domains["session"]["http"] = {
+        sessionID: Session.ID.make("ses_promise_session_http"),
+        agent: Agent.ID.make("build"),
+        model: Model.Ref.make({ providerID: Provider.ID.make("test"), id: Model.ID.make("model") }),
+        use: (item) =>
+          Effect.sync(() => {
+            middlewares.push(item)
+          }),
+      }
+
+      yield* hooks.trigger("session", "http", event)
+      const request = middlewares.reduce<SessionHttpHandler>(
+        (next, item) => (input: Request) => item(input, next),
+        (input: Request) =>
+          Effect.promise(() => input.text()).pipe(
+            Effect.tap((body) => Effect.sync(() => bodies.push(body))),
+            Effect.as(new Response(input.headers.get("x-hook") ?? "missing")),
+          ),
+      )
+      const response = yield* request(new Request("https://provider.test", { method: "POST", body: "payload" }))
+
+      expect(bodies).toEqual(["payload", "payload"])
+      expect(yield* Effect.promise(() => response.text())).toBe("promise-response-outer")
+    }),
+  )
+
+  it.effect("interrupts the Effect request through a promise session HTTP hook", () =>
+    Effect.gen(function* () {
+      const plugin = yield* Plugin.Service
+      const hooks = yield* PluginHooks.Service
+      const host = yield* PluginHost.make(plugin)
+      yield* PluginPromise.fromPromise(
+        define({
+          id: "promise-session-http-interrupt",
+          setup: async (ctx) => {
+            await ctx.session.hook("http", (event) => {
+              event.use((request, next) => next(request))
+            })
+          },
+        }),
+      ).effect(host)
+      const started = yield* Deferred.make<void>()
+      const interrupted = yield* Deferred.make<void>()
+      const middlewares: Parameters<PluginHooks.Domains["session"]["http"]["use"]>[0][] = []
+      const event: PluginHooks.Domains["session"]["http"] = {
+        sessionID: Session.ID.make("ses_promise_session_http_interrupt"),
+        agent: Agent.ID.make("build"),
+        model: Model.Ref.make({ providerID: Provider.ID.make("test"), id: Model.ID.make("model") }),
+        use: (item) =>
+          Effect.sync(() => {
+            middlewares.push(item)
+          }),
+      }
+
+      yield* hooks.trigger("session", "http", event)
+      const request = middlewares.reduce<SessionHttpHandler>(
+        (next, item) => (input: Request) => item(input, next),
+        () =>
+          Deferred.succeed(started, undefined).pipe(
+            Effect.andThen(Effect.never),
+            Effect.onInterrupt(() => Deferred.succeed(interrupted, undefined)),
+          ),
+      )
+      const fiber = yield* request(new Request("https://provider.test")).pipe(Effect.forkChild)
+      yield* Deferred.await(started)
+      yield* Fiber.interrupt(fiber)
+
+      expect(yield* Deferred.isDone(interrupted)).toBeTrue()
+    }),
+  )
+
   it.effect("disposes a hook registration on request", () =>
     Effect.gen(function* () {
       const agents = yield* Agent.Service
@@ -315,19 +416,17 @@ describe("fromPromise", () => {
         id: "promise-tool",
         setup: async (ctx) => {
           await ctx.tool.transform((tools) => {
-            tools.add(
-              {
-                name: "hello",
-                options: { codemode: false },
-                description: "Hello",
-                input: Schema.Struct({ name: Schema.String }),
-                output: Schema.String,
-                execute: async ({ name }, context) => {
-                  await context.progress({ phase: "greeting" })
-                  return { output: `Hello, ${name}!` }
-                },
+            tools.add({
+              name: "hello",
+              options: { codemode: false },
+              description: "Hello",
+              input: Schema.Struct({ name: Schema.String }),
+              output: Schema.String,
+              execute: async ({ name }, context) => {
+                await context.progress({ phase: "greeting" })
+                return { output: `Hello, ${name}!` }
               },
-            )
+            })
           })
         },
       })

+ 30 - 38
packages/core/test/plugin/provider-openai.test.ts

@@ -12,6 +12,7 @@ import { PluginHost } from "@opencode-ai/core/plugin/host"
 import { PluginHooks } from "@opencode-ai/core/plugin/hooks"
 import { OpenAIPlugin } from "@opencode-ai/core/plugin/provider/openai"
 import { Provider } from "@opencode-ai/core/provider"
+import type { SessionHttpHandler } from "@opencode-ai/plugin/effect/session"
 import { testEffect } from "../lib/effect"
 import { PluginTestLayer } from "./fixture"
 
@@ -29,6 +30,29 @@ function required<T>(value: T | undefined): T {
   return value
 }
 
+const http = Effect.fn(function* (providerID: Provider.ID, url: string) {
+  const middlewares: Parameters<PluginHooks.Domains["session"]["http"]["use"]>[0][] = []
+  yield* (yield* PluginHooks.Service).trigger("session", "http", {
+    sessionID: Session.ID.make("ses_test"),
+    agent: Agent.ID.make("build"),
+    model: Model.Ref.make({ providerID, id: Model.ID.make("gpt-5.5") }),
+    use: (item) =>
+      Effect.sync(() => {
+        middlewares.push(item)
+      }),
+  })
+  const request = middlewares.reduce<SessionHttpHandler>(
+    (next, item) => (input: Request) => item(input, next),
+    (input: Request) => {
+      const headers = new Headers(input.headers)
+      headers.set("x-seen-url", input.url)
+      return Effect.succeed(new Response(null, { headers }))
+    },
+  )
+  const response = yield* request(new Request(url, { method: "POST", body: "{}" }))
+  return { url: response.headers.get("x-seen-url"), headers: Object.fromEntries(response.headers.entries()) }
+})
+
 describe("OpenAIPlugin", () => {
   it.effect("registers browser and headless ChatGPT OAuth methods", () =>
     Effect.gen(function* () {
@@ -100,33 +124,9 @@ describe("OpenAIPlugin", () => {
       })
       yield* addPlugin()
 
-      const request = yield* (yield* PluginHooks.Service).trigger("session", "request", {
-        sessionID: Session.ID.make("ses_test"),
-        agent: Agent.ID.make("build"),
-        model: Model.Ref.make({ providerID: Provider.ID.openai, id: Model.ID.make("gpt-5.5") }),
-        url: "https://api.openai.com/v1/responses",
-        method: "POST",
-        headers: {},
-        body: "{}",
-      })
-      const custom = yield* (yield* PluginHooks.Service).trigger("session", "request", {
-        sessionID: Session.ID.make("ses_test"),
-        agent: Agent.ID.make("build"),
-        model: Model.Ref.make({ providerID: Provider.ID.make("custom-openai"), id: Model.ID.make("gpt-5.5") }),
-        url: "https://custom.example/v1/responses",
-        method: "POST",
-        headers: {},
-        body: "{}",
-      })
-      const proxy = yield* (yield* PluginHooks.Service).trigger("session", "request", {
-        sessionID: Session.ID.make("ses_test"),
-        agent: Agent.ID.make("build"),
-        model: Model.Ref.make({ providerID: Provider.ID.openai, id: Model.ID.make("gpt-5.5") }),
-        url: "https://proxy.example/v1/responses?region=us",
-        method: "POST",
-        headers: {},
-        body: "{}",
-      })
+      const request = yield* http(Provider.ID.openai, "https://api.openai.com/v1/responses")
+      const custom = yield* http(Provider.ID.make("custom-openai"), "https://custom.example/v1/responses")
+      const proxy = yield* http(Provider.ID.openai, "https://proxy.example/v1/responses?region=us")
 
       const provider = required(yield* catalog.provider.get(Provider.ID.openai))
       expect(provider.package).toBe("@opencode-ai/ai/providers/openai")
@@ -134,7 +134,7 @@ describe("OpenAIPlugin", () => {
       expect(provider.headers).toMatchObject({ "chatgpt-account-id": "acct_123" })
       expect(request.url).toBe("https://chatgpt.com/backend-api/codex/responses")
       expect(request.headers).toMatchObject({ originator: "opencode", "session-id": "ses_test" })
-      expect(custom.headers).toEqual({})
+      expect(custom.headers).not.toHaveProperty("originator")
       expect(proxy.url).toBe("https://proxy.example/v1/responses?region=us")
       expect(proxy.headers).toMatchObject({ originator: "opencode", "session-id": "ses_test" })
       const eligible = required(yield* catalog.model.get(Provider.ID.openai, Model.ID.make("gpt-5.5")))
@@ -184,21 +184,13 @@ describe("OpenAIPlugin", () => {
       })
       yield* addPlugin()
 
-      const request = yield* (yield* PluginHooks.Service).trigger("session", "request", {
-        sessionID: Session.ID.make("ses_test"),
-        agent: Agent.ID.make("build"),
-        model: Model.Ref.make({ providerID: Provider.ID.openai, id: Model.ID.make("gpt-5.5") }),
-        url: "https://api.openai.com/v1/responses",
-        method: "POST",
-        headers: {},
-        body: "{}",
-      })
+      const request = yield* http(Provider.ID.openai, "https://api.openai.com/v1/responses")
 
       const model = required(yield* catalog.model.get(Provider.ID.openai, Model.ID.make("gpt-5.5")))
       expect(model.package).toBe("@opencode-ai/ai/providers/openai")
       expect(model.enabled).toBe(true)
       expect(model.limit).toEqual({ context: 1_050_000, input: 922_000, output: 128_000 })
-      expect(request.headers).toEqual({})
+      expect(request.headers).not.toHaveProperty("originator")
       expect(required(yield* catalog.model.get(Provider.ID.openai, Model.ID.make("gpt-4.1"))).enabled).toBe(true)
     }),
   )

+ 125 - 35
packages/core/test/session-runner-recorded.test.ts

@@ -41,6 +41,7 @@ import { SystemPromptPlugin } from "@opencode-ai/core/plugin/system-prompt"
 import { describe, expect } from "bun:test"
 import { eq } from "drizzle-orm"
 import { Effect, Layer, Stream } from "effect"
+import { HttpClient, HttpClientResponse } from "effect/unstable/http"
 import path from "node:path"
 import { testEffect } from "./lib/effect"
 import { agentHost, catalogHost, host } from "./plugin/host"
@@ -104,37 +105,39 @@ const promptCatalog = Layer.mock(Catalog.Service, {
     small: () => Effect.succeed(undefined),
   },
 })
-const runnerLayer = AppNodeBuilder.build(SessionRunnerLLM.node, [
-  [Snapshot.node, Snapshot.noopLayer],
-  [LayerNodePlatform.llmClient, client],
-  [SessionRunnerModel.node, models],
-  [InstructionBuiltIns.node, systemContext],
-  [InstructionDiscovery.node, instructionContext],
-  [Location.node, Location.boundNode({ directory: AbsolutePath.make("/project") })],
-  [SkillInstructions.node, skillInstructions],
-  [ReferenceInstructions.node, referenceInstructions],
-  [McpInstructions.node, mcpInstructions],
-  [Config.node, config],
-  [Permission.node, permission],
-  [PluginSupervisor.node, pluginSupervisor],
-])
-const execution = Layer.effect(
-  SessionExecution.Service,
-  Effect.gen(function* () {
-    const sessionRunner = yield* SessionRunner.Service
-    const coordinator = yield* SessionRunCoordinator.make<Session.ID, SessionRunner.RunError>({
-      drain: (sessionID, force) => sessionRunner.drain({ sessionID, force }),
-    })
-    return SessionExecution.Service.of({
-      active: coordinator.active,
-      resume: coordinator.run,
-      wake: coordinator.wake,
-      interrupt: coordinator.interrupt,
-      awaitIdle: coordinator.awaitIdle,
-    })
-  }),
-).pipe(Layer.provide(runnerLayer))
-const it = testEffect(
+const runnerLayer = (llmClient: Layer.Layer<typeof LLMClient.Service>) =>
+  AppNodeBuilder.build(SessionRunnerLLM.node, [
+    [Snapshot.node, Snapshot.noopLayer],
+    [LayerNodePlatform.llmClient, llmClient],
+    [SessionRunnerModel.node, models],
+    [InstructionBuiltIns.node, systemContext],
+    [InstructionDiscovery.node, instructionContext],
+    [Location.node, Location.boundNode({ directory: AbsolutePath.make("/project") })],
+    [SkillInstructions.node, skillInstructions],
+    [ReferenceInstructions.node, referenceInstructions],
+    [McpInstructions.node, mcpInstructions],
+    [Config.node, config],
+    [Permission.node, permission],
+    [PluginSupervisor.node, pluginSupervisor],
+  ])
+const execution = (llmClient: Layer.Layer<typeof LLMClient.Service>) =>
+  Layer.effect(
+    SessionExecution.Service,
+    Effect.gen(function* () {
+      const sessionRunner = yield* SessionRunner.Service
+      const coordinator = yield* SessionRunCoordinator.make<Session.ID, SessionRunner.RunError>({
+        drain: (sessionID, force) => sessionRunner.drain({ sessionID, force }),
+      })
+      return SessionExecution.Service.of({
+        active: coordinator.active,
+        resume: coordinator.run,
+        wake: coordinator.wake,
+        interrupt: coordinator.interrupt,
+        awaitIdle: coordinator.awaitIdle,
+      })
+    }),
+  ).pipe(Layer.provide(runnerLayer(llmClient)))
+const testLayer = (llmClient: Layer.Layer<typeof LLMClient.Service>) =>
   AppNodeBuilder.build(
     LayerNode.group([
       Database.node,
@@ -156,7 +159,7 @@ const it = testEffect(
       Session.node,
     ]),
     [
-      [LayerNodePlatform.llmClient, client],
+      [LayerNodePlatform.llmClient, llmClient],
       [Permission.node, permission],
       [Catalog.node, promptCatalog],
       [SessionRunnerModel.node, models],
@@ -168,10 +171,10 @@ const it = testEffect(
       [Config.node, config],
       [Snapshot.node, Snapshot.noopLayer],
       [PluginSupervisor.node, pluginSupervisor],
-      [SessionExecution.node, execution],
+      [SessionExecution.node, execution(llmClient)],
     ],
-  ),
-)
+  )
+const it = testEffect(testLayer(client))
 const sessionID = Session.ID.make("ses_runner_recorded")
 
 describe("SessionRunnerLLM recorded", () => {
@@ -247,3 +250,90 @@ describe("SessionRunnerLLM recorded", () => {
     }),
   )
 })
+
+describe("SessionModelRequest HTTP bridge", () => {
+  const bodies: Uint8Array[] = []
+  const methods: string[] = []
+  const response = [
+    'data: {"id":"chatcmpl_test","object":"chat.completion.chunk","created":0,"model":"gpt-4o-mini","choices":[{"index":0,"delta":{"role":"assistant","content":"Hello!"},"finish_reason":null}]}',
+    'data: {"id":"chatcmpl_test","object":"chat.completion.chunk","created":0,"model":"gpt-4o-mini","choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}',
+    "data: [DONE]",
+    "",
+  ].join("\n\n")
+  const transport = Layer.succeed(
+    HttpClient.HttpClient,
+    HttpClient.make((request) =>
+      Effect.sync(() => {
+        if (request.body._tag !== "Uint8Array") throw new Error(`Unexpected request body: ${request.body._tag}`)
+        methods.push(request.method)
+        bodies.push(request.body.body.slice())
+        return HttpClientResponse.fromWeb(
+          request,
+          new Response(response, { headers: { "content-type": "text/event-stream" } }),
+        )
+      }),
+    ),
+  )
+  const retryIt = testEffect(
+    testLayer(LLMClient.layer.pipe(Layer.provide(RequestExecutor.layer.pipe(Layer.provide(transport))))),
+  )
+
+  retryIt.effect("lets an Effect plugin send the same POST Request twice", () =>
+    Effect.gen(function* () {
+      bodies.length = 0
+      methods.length = 0
+      const agents = yield* Agent.Service
+      const catalog = yield* Catalog.Service
+      const hooks = yield* PluginHooks.Service
+      yield* agents.transform((draft) =>
+        draft.update(Agent.ID.make("build"), (agent) => {
+          agent.mode = "primary"
+          agent.permissions.push({ action: "execute", resource: "*", effect: "deny" })
+        }),
+      )
+      const pluginHost = host({
+        agent: agentHost(agents),
+        catalog: catalogHost(catalog),
+        session: { hook: (name, callback) => hooks.register("session", name, callback) },
+      })
+      yield* pluginHost.session.hook("http", (event) =>
+        event.use((request, next) =>
+          Effect.gen(function* () {
+            yield* next(request).pipe(Effect.flatMap((response) => Effect.promise(() => response.text())))
+            return yield* next(request)
+          }),
+        ),
+      )
+      yield* Effect.forEach(SystemPromptPlugin.Plugins, (plugin) => plugin.effect(pluginHost), { discard: true })
+      const { db } = yield* Database.Service
+      yield* db
+        .insert(ProjectTable)
+        .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
+        .onConflictDoNothing()
+        .run()
+        .pipe(Effect.orDie)
+      const retrySessionID = Session.ID.make("ses_model_request_http_retry")
+      yield* db
+        .insert(SessionTable)
+        .values({
+          id: retrySessionID,
+          project_id: Project.ID.global,
+          slug: "test",
+          directory: "/project",
+          title: "test",
+          version: "test",
+        })
+        .run()
+        .pipe(Effect.orDie)
+      const session = yield* Session.Service
+      yield* session.prompt({ sessionID: retrySessionID, text: "Say hello.", resume: false })
+
+      yield* session.resume(retrySessionID)
+
+      expect(methods).toEqual(["POST", "POST"])
+      expect(bodies).toHaveLength(2)
+      expect(bodies[0]?.byteLength).toBeGreaterThan(0)
+      expect(bodies[1]).toEqual(bodies[0])
+    }),
+  )
+})

+ 11 - 4
packages/plugin/src/effect/session.ts

@@ -1,10 +1,9 @@
 import type { SessionApi } from "@opencode-ai/client/effect/api"
 import type { Message, SystemPart } from "@opencode-ai/ai"
-import type { HttpRequest } from "@opencode-ai/ai/route"
 import type { Agent } from "@opencode-ai/schema/agent"
 import type { Model } from "@opencode-ai/schema/model"
 import type { Session } from "@opencode-ai/schema/session"
-import type { JsonSchema } from "effect"
+import type { Effect, JsonSchema } from "effect"
 import type { Hooks } from "./registration.js"
 
 export interface SessionContext {
@@ -16,15 +15,23 @@ export interface SessionContext {
   tools: Record<string, { description: string; input: JsonSchema.JsonSchema }>
 }
 
-export interface SessionRequest extends HttpRequest {
+export interface SessionHttp {
   readonly sessionID: Session.ID
   readonly agent: Agent.ID
   readonly model: Model.Ref
+  readonly use: (middleware: SessionHttpMiddleware) => Effect.Effect<void>
 }
 
+export type SessionHttpHandler = (request: Request) => Effect.Effect<Response, Error>
+
+export type SessionHttpMiddleware = (
+  request: Request,
+  next: SessionHttpHandler,
+) => Effect.Effect<Response, Error>
+
 export interface SessionHooks {
   readonly context: SessionContext
-  readonly request: SessionRequest
+  readonly http: SessionHttp
 }
 
 export type SessionDomain = Pick<

+ 10 - 3
packages/plugin/src/promise/session.ts

@@ -1,6 +1,5 @@
 import type { SessionApi } from "@opencode-ai/client/promise/api"
 import type { Message, SystemPart } from "@opencode-ai/ai"
-import type { HttpRequest } from "@opencode-ai/ai/route"
 import type { Agent } from "@opencode-ai/schema/agent"
 import type { Model } from "@opencode-ai/schema/model"
 import type { Session } from "@opencode-ai/schema/session"
@@ -16,15 +15,23 @@ export interface SessionContext {
   tools: Record<string, { description: string; input: JsonSchema.JsonSchema }>
 }
 
-export interface SessionRequest extends HttpRequest {
+export interface SessionHttp {
   readonly sessionID: Session.ID
   readonly agent: Agent.ID
   readonly model: Model.Ref
+  readonly use: (middleware: SessionHttpMiddleware) => void
 }
 
+export type SessionHttpHandler = (request: Request) => Promise<Response>
+
+export type SessionHttpMiddleware = (
+  request: Request,
+  next: SessionHttpHandler,
+) => Promise<Response> | Response
+
 export interface SessionHooks {
   readonly context: SessionContext
-  readonly request: SessionRequest
+  readonly http: SessionHttp
 }
 
 export type SessionDomain = Pick<

+ 18 - 5
packages/www/content/docs/build/plugins.mdx

@@ -239,16 +239,29 @@ without restarting OpenCode.
 
 ### Runtime hooks
 
-Runtime hooks intercept live operations. Their event objects expose specific
-mutable fields:
+Runtime hooks intercept live operations:
 
 | Hook                                        | Mutable fields                                                                 |
 | ------------------------------------------- | ------------------------------------------------------------------------------ |
 | `ctx.aisdk.hook("sdk", callback)`           | `sdk`, after inspecting `model`, `package`, and `options`                      |
 | `ctx.aisdk.hook("language", callback)`      | `language`, after inspecting `model`, `sdk`, and `options`                     |
-| `ctx.session.hook("request", callback)`     | `system`, `messages`, and the `tools` record immediately before model dispatch |
+| `ctx.session.hook("context", callback)`     | `system`, `messages`, and the `tools` record immediately before model dispatch |
+| `ctx.session.hook("http", callback)`       | `use`, registering request and response handling                               |
 | `ctx.tool.hook("execute.before", callback)` | `input`, before the selected tool executes                                     |
-| `ctx.tool.hook("execute.after", callback)`  | Terminal `result` on success or `error` on failure                            |
+| `ctx.tool.hook("execute.after", callback)`  | Terminal `result` on success or `error` on failure                             |
+
+HTTP hooks can modify requests, inspect responses, retry, or return a
+response without calling the provider. It applies to native models; AI SDK
+models do not currently pass through this hook.
+
+```ts
+await ctx.session.hook("http", (event) => {
+  event.use((request, next) => {
+    request.headers.set("x-session-id", event.sessionID)
+    return next(request)
+  })
+})
+```
 
 For example, remove a tool from selected model requests and normalize another
 tool's input:
@@ -259,7 +272,7 @@ import { Plugin } from "@opencode-ai/plugin"
 export default Plugin.define({
   id: "acme.guards",
   setup: async (ctx) => {
-    await ctx.session.hook("request", (event) => {
+    await ctx.session.hook("context", (event) => {
       delete event.tools.write
     })