diff --git a/packages/ai/src/route/client.ts b/packages/ai/src/route/client.ts index 03f43d17a19..f4873f347ef 100644 --- a/packages/ai/src/route/client.ts +++ b/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( 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}` diff --git a/packages/ai/src/route/executor.ts b/packages/ai/src/route/executor.ts index f7a0fb465b3..9b90bf1aa4f 100644 --- a/packages/ai/src/route/executor.ts +++ b/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 } +export type HttpHandler = ( + request: HttpClientRequest.HttpClientRequest, +) => Effect.Effect +export type HttpMiddleware = ( + request: HttpClientRequest.HttpClientRequest, + handler: HttpHandler, +) => Effect.Effect + export class Service extends Context.Service()("@opencode/AI/RequestExecutor") {} const BODY_LIMIT = 16_384 @@ -261,7 +270,7 @@ const toHttpError = (redactedNames: ReadonlyArray) => (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 = 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, diff --git a/packages/ai/src/route/index.ts b/packages/ai/src/route/index.ts index 8b85cfc3f76..eb9f1759c55 100644 --- a/packages/ai/src/route/index.ts +++ b/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" diff --git a/packages/ai/src/route/transport/http.ts b/packages/ai/src/route/transport/http.ts index a1de9ec5735..223afb353c1 100644 --- a/packages/ai/src/route/transport/http.ts +++ b/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 { export interface HttpPrepared { readonly request: HttpClientRequest.HttpClientRequest readonly framing: Framing.Definition + readonly middleware?: HttpMiddleware } const applyQuery = (url: string, query: Record | undefined) => { @@ -74,21 +75,21 @@ export const httpJson = (input: HttpJsonInput): 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( diff --git a/packages/ai/src/route/transport/index.ts b/packages/ai/src/route/transport/index.ts index 16624cce563..c74e578e8bb 100644 --- a/packages/ai/src/route/transport/index.ts +++ b/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 - body: string | undefined -} - -export type HttpRequestTransform = (request: HttpRequest) => Effect.Effect - export interface Transport { readonly id: string readonly prepare: (input: TransportPrepareInput) => Effect.Effect @@ -32,8 +23,9 @@ export interface TransportPrepareInput { readonly auth: Auth.Definition readonly encodeBody: (body: Body) => string readonly headers?: (input: { readonly request: LLMRequest }) => Record - readonly transform?: HttpRequestTransform + readonly middleware?: HttpMiddleware } export * as HttpTransport from "./http" +export type { HttpHandler, HttpMiddleware } from "../executor" export { WebSocketExecutor, WebSocketTransport } from "./websocket" diff --git a/packages/ai/test/compile.test.ts b/packages/ai/test/compile.test.ts index 27971bf5c8e..0f1260873e0 100644 --- a/packages/ai/test/compile.test.ts +++ b/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({ diff --git a/packages/ai/test/executor.test.ts b/packages/ai/test/executor.test.ts index 91417c8f08f..b4e5fe62b8e 100644 --- a/packages/ai/test/executor.test.ts +++ b/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 diff --git a/packages/core/src/plugin/promise.ts b/packages/core/src/plugin/promise.ts index 651bc7c8ba5..4880b9e3c3b 100644 --- a/packages/core/src/plugin/promise.ts +++ b/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: Name, + callback: (event: SessionHooks[Name]) => Promise | void, + ): Promise + function sessionHook( + ...registration: { + [Name in keyof SessionHooks]: [ + name: Name, + callback: (event: SessionHooks[Name]) => Promise | 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( diff --git a/packages/core/src/plugin/provider/openai.ts b/packages/core/src/plugin/provider/openai.ts index cbd101eacaf..8a464a31415 100644 --- a/packages/core/src/plugin/provider/openai.ts +++ b/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)) }), ) diff --git a/packages/core/src/session/model-request.ts b/packages/core/src/session/model-request.ts index bae38339c00..8107f547cdf 100644 --- a/packages/core/src/session/model-request.ts +++ b/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[0], - ) => Effect.Effect + readonly executeTool: (input: Parameters[0]) => Effect.Effect /** 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() + 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( + (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) diff --git a/packages/core/test/plugin/promise.test.ts b/packages/core/test/plugin/promise.test.ts index f63c7f3eca8..50d9f630f04 100644 --- a/packages/core/test/plugin/promise.test.ts +++ b/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[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( + (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() + const interrupted = yield* Deferred.make() + const middlewares: Parameters[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( + (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}!` } }, - ) + }) }) }, }) diff --git a/packages/core/test/plugin/provider-openai.test.ts b/packages/core/test/plugin/provider-openai.test.ts index 9a73f5af939..239599e31db 100644 --- a/packages/core/test/plugin/provider-openai.test.ts +++ b/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(value: T | undefined): T { return value } +const http = Effect.fn(function* (providerID: Provider.ID, url: string) { + const middlewares: Parameters[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( + (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) }), ) diff --git a/packages/core/test/session-runner-recorded.test.ts b/packages/core/test/session-runner-recorded.test.ts index 6fd9ef67dad..3d6b3de38fd 100644 --- a/packages/core/test/session-runner-recorded.test.ts +++ b/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({ - 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) => + 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) => + Layer.effect( + SessionExecution.Service, + Effect.gen(function* () { + const sessionRunner = yield* SessionRunner.Service + const coordinator = yield* SessionRunCoordinator.make({ + 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) => 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]) + }), + ) +}) diff --git a/packages/plugin/src/effect/session.ts b/packages/plugin/src/effect/session.ts index e4206b6df25..d00eba103a5 100644 --- a/packages/plugin/src/effect/session.ts +++ b/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 } -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 } +export type SessionHttpHandler = (request: Request) => Effect.Effect + +export type SessionHttpMiddleware = ( + request: Request, + next: SessionHttpHandler, +) => Effect.Effect + export interface SessionHooks { readonly context: SessionContext - readonly request: SessionRequest + readonly http: SessionHttp } export type SessionDomain = Pick< diff --git a/packages/plugin/src/promise/session.ts b/packages/plugin/src/promise/session.ts index adc371b9510..40866a75cd3 100644 --- a/packages/plugin/src/promise/session.ts +++ b/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 } -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 + +export type SessionHttpMiddleware = ( + request: Request, + next: SessionHttpHandler, +) => Promise | Response + export interface SessionHooks { readonly context: SessionContext - readonly request: SessionRequest + readonly http: SessionHttp } export type SessionDomain = Pick< diff --git a/packages/www/content/docs/build/plugins.mdx b/packages/www/content/docs/build/plugins.mdx index fd9331c8805..955a7a150ff 100644 --- a/packages/www/content/docs/build/plugins.mdx +++ b/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 })