diff --git a/packages/core/src/plugin/promise.ts b/packages/core/src/plugin/promise.ts index 651bc7c8ba5..32e5342af26 100644 --- a/packages/core/src/plugin/promise.ts +++ b/packages/core/src/plugin/promise.ts @@ -194,7 +194,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 +265,35 @@ export function fromPromise(plugin: Plugin) { ), }, session: { - hook: (name, callback) => - register(host.session.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))), + hook: (name, callback) => { + if (name !== "http") + return register( + host.session.hook(name, (event) => + Effect.promise(() => Promise.resolve(Reflect.apply(callback, undefined, [event]))), + ), + ) + return register( + host.session.hook("http", (event) => { + const request = event.request + const output = { + ...event, + request: (input: Request) => + Effect.runPromiseWith(context)(request(input), { signal: input.signal }), + } + return Effect.promise(() => Promise.resolve(Reflect.apply(callback, undefined, [output]))).pipe( + Effect.tap(() => + Effect.sync(() => { + event.request = (input) => + Effect.tryPromise({ + try: (signal) => output.request(new Request(input, { signal })), + catch: (cause) => (cause instanceof Error ? cause : new Error(String(cause))), + }) + }), + ), + ) + }), + ) + }, 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..b57cdf5486f 100644 --- a/packages/core/src/plugin/provider/openai.ts +++ b/packages/core/src/plugin/provider/openai.ts @@ -225,15 +225,25 @@ export const OpenAIPlugin = define({ }) } }) - yield* ctx.session.hook("request", (evt) => + yield* ctx.session.hook("http", (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}` + const request = evt.request + evt.request = (input) => { + const url = new URL(input.url) + const headers = new Headers(input.headers) + headers.set("originator", "opencode") + headers.set("session-id", evt.sessionID) + if (url.origin !== "https://api.openai.com") return request(new Request(input, { headers })) + return request( + new Request(`${codexBaseURL}${url.pathname.replace(/^\/v1/, "")}${url.search}`, { + method: input.method, + headers, + body: input.body, + signal: input.signal, + }), + ) } - evt.headers.originator = "opencode" - evt.headers["session-id"] = evt.sessionID }), ) diff --git a/packages/core/src/session/model-request.ts b/packages/core/src/session/model-request.ts index fab9e4c0e46..9c66f79fc01 100644 --- a/packages/core/src/session/model-request.ts +++ b/packages/core/src/session/model-request.ts @@ -4,7 +4,8 @@ import { LLM, Message, SystemPart, type LLMRequest } from "@opencode-ai/ai" import type { StreamOptions } from "@opencode-ai/ai/route" 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 +49,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 +136,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 } }), @@ -204,9 +202,7 @@ export const layer = Layer.effect( }) const hookedTools = Object.entries(contextEvent.tools).flatMap(([name, tool]) => { const registered = toolsByName.get(name) - return registered - ? [{ ...registered, description: tool.description, inputSchema: tool.input }] - : [] + return registered ? [{ ...registered, description: tool.description, inputSchema: tool.input }] : [] }) const request = LLM.request({ model, @@ -220,24 +216,37 @@ 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 sent = request + const origins = new WeakMap() + const web = yield* HttpClientRequest.toWeb(request) + const event = yield* hooks.trigger("session", "http", { sessionID: session.id, agent: agent.id, model: resolved.ref, - ...request, - }) - .pipe( - Effect.tap((event) => - Effect.sync(() => { - request.url = event.url - request.headers = event.headers - request.body = event.body + request: (input) => + Effect.gen(function* () { + sent = HttpClientRequest.fromWeb(input) + if (input.body) + sent = HttpClientRequest.bodyUint8Array( + sent, + new Uint8Array(yield* Effect.promise(() => input.arrayBuffer())), + input.headers.get("content-type") ?? undefined, + ) + 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 }), - ), - Effect.asVoid, - ), + }) + const response = yield* event.request(web) + const origin = origins.get(response) ?? sent + return HttpClientResponse.fromWeb(origin, response) + }).pipe(Effect.mapError((cause) => (cause instanceof Error ? cause : new Error(String(cause))))), } if (promptCacheSnapshots) { const current = PromptCacheDiagnostics.snapshot(request) @@ -257,8 +266,7 @@ export const layer = Layer.effect( ) } const executeTool: Prepared["executeTool"] = (executeInput) => { - if (stepLimitReached) - return new Tool.Error({ message: "Tools are disabled after the maximum agent steps" }) + if (stepLimitReached) return new Tool.Error({ message: "Tools are disabled after the maximum agent steps" }) if (toolsByName.has(executeInput.call.name) && !Object.hasOwn(contextEvent.tools, executeInput.call.name)) return new Tool.Error({ message: `Tool is not available for this request: ${executeInput.call.name}` }) return tools diff --git a/packages/core/test/plugin/promise.test.ts b/packages/core/test/plugin/promise.test.ts index f63c7f3eca8..b0c01fe44f6 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" @@ -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,77 @@ 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) + yield* PluginPromise.fromPromise( + define({ + id: "promise-session-http", + setup: async (ctx) => { + await ctx.session.hook("http", (event) => { + const request = event.request + event.request = async (input) => { + const response = await request(new Request(input, { headers: { "x-hook": "promise" } })) + return new Response(`${await response.text()}-response`) + } + }) + }, + }), + ).effect(host) + const event: SessionHooks["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") }), + request: (input) => Effect.succeed(new Response(input.headers.get("x-hook") ?? "missing")), + } + + yield* hooks.trigger("session", "http", event) + const response = yield* event.request(new Request("https://provider.test")) + + expect(yield* Effect.promise(() => response.text())).toBe("promise-response") + }), + ) + + 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) => { + const request = event.request + event.request = (input) => request(input) + }) + }, + }), + ).effect(host) + const started = yield* Deferred.make() + const interrupted = yield* Deferred.make() + const event: SessionHooks["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") }), + request: () => + Deferred.succeed(started, undefined).pipe( + Effect.andThen(Effect.never), + Effect.onInterrupt(() => Deferred.succeed(interrupted, undefined)), + ), + } + + yield* hooks.trigger("session", "http", event) + const fiber = yield* event.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 +388,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..de5f87e09e2 100644 --- a/packages/core/test/plugin/provider-openai.test.ts +++ b/packages/core/test/plugin/provider-openai.test.ts @@ -29,6 +29,21 @@ function required(value: T | undefined): T { return value } +const http = Effect.fn(function* (providerID: Provider.ID, url: string) { + const event = 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") }), + request: (input) => { + const headers = new Headers(input.headers) + headers.set("x-seen-url", input.url) + return Effect.succeed(new Response(null, { headers })) + }, + }) + const response = yield* event.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 +115,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 +125,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 +175,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/plugin/src/effect/session.ts b/packages/plugin/src/effect/session.ts index e4206b6df25..c2bf4b0ce83 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,16 @@ 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 + request: (input: Request) => 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..b94121d4812 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,16 @@ 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 + request: (input: Request) => Promise } 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..f8e1c5cbc50 100644 --- a/packages/www/content/docs/build/plugins.mdx +++ b/packages/www/content/docs/build/plugins.mdx @@ -246,7 +246,8 @@ 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)` | `request`, wrapping the model's HTTP request and response | | `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 | @@ -259,7 +260,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 })