Compare commits

...

6 Commits

Author SHA1 Message Date
Aiden Cline 2b249974fa refactor(ai): keep HTTP middleware Effect-native 2026-08-03 14:50:04 -05:00
Aiden Cline 1eec3e640a fix(ai): expose raw HTTP responses to hooks 2026-08-02 11:44:53 -05:00
Aiden Cline 4dc5ba66d8 feat(plugin): wrap native session HTTP 2026-08-01 22:02:56 -05:00
Aiden Cline a21598901d refactor(plugin): rename session HTTP hook 2026-08-01 21:54:49 -05:00
Aiden Cline 3a7f507135 Revert "feat(plugin): wrap session HTTP requests"
This reverts commit b1f86ee72b.
2026-08-01 21:54:28 -05:00
Aiden Cline b1f86ee72b feat(plugin): wrap session HTTP requests 2026-08-01 11:44:08 -05:00
14 changed files with 329 additions and 132 deletions
+7 -4
View File
@@ -5,7 +5,7 @@ import { Endpoint, type EndpointPatch } from "./endpoint"
import { RequestExecutor } from "./executor" import { RequestExecutor } from "./executor"
import { Framing } from "./framing" import { Framing } from "./framing"
import { HttpTransport } from "./transport" import { HttpTransport } from "./transport"
import type { HttpRequestTransform, Transport, TransportRuntime } from "./transport" import type { HttpMiddleware, Transport, TransportRuntime } from "./transport"
import { WebSocketExecutor } from "./transport" import { WebSocketExecutor } from "./transport"
import type { Protocol } from "./protocol" import type { Protocol } from "./protocol"
import { applyCachePolicy } from "../cache-policy" import { applyCachePolicy } from "../cache-policy"
@@ -96,7 +96,10 @@ export interface RoutePatch<Body, Prepared> extends RouteDefaultsInput {
type RouteMappedModelInput = RouteModelInput | RouteRoutedModelInput type RouteMappedModelInput = RouteModelInput | RouteRoutedModelInput
const makeRouteModel = <Options extends ProviderOptions = ProviderOptions>(route: AnyRoute, mapped: RouteMappedModelInput) => { const makeRouteModel = <Options extends ProviderOptions = ProviderOptions>(
route: AnyRoute,
mapped: RouteMappedModelInput,
) => {
const provider = route.provider ?? ("provider" in mapped ? mapped.provider : undefined) const provider = route.provider ?? ("provider" in mapped ? mapped.provider : undefined)
if (!provider) throw new Error(`Route.model(${route.id}) requires a provider`) if (!provider) throw new Error(`Route.model(${route.id}) requires a provider`)
if (!endpointBaseURL(route.endpoint)) if (!endpointBaseURL(route.endpoint))
@@ -150,7 +153,7 @@ export interface Interface {
} }
export interface StreamOptions { export interface StreamOptions {
readonly transform?: HttpRequestTransform readonly http?: HttpMiddleware
} }
export interface StreamMethod { export interface StreamMethod {
@@ -302,7 +305,7 @@ function makeFromTransport<Body, Prepared, Frame, Event, State>(
auth: routeInput.auth ?? Auth.none, auth: routeInput.auth ?? Auth.none,
encodeBody, encodeBody,
headers: routeInput.headers, headers: routeInput.headers,
transform: options?.transform, middleware: options?.http,
}), }),
streamPrepared: (prepared: Prepared, request: LLMRequest, runtime: TransportRuntime) => { streamPrepared: (prepared: Prepared, request: LLMRequest, runtime: TransportRuntime) => {
const route = `${request.model.provider}/${request.model.route.id}` const route = `${request.model.provider}/${request.model.route.id}`
+21 -4
View File
@@ -20,9 +20,18 @@ import { classifyProviderFailure } from "../provider-error"
export interface Interface { export interface Interface {
readonly execute: ( readonly execute: (
request: HttpClientRequest.HttpClientRequest, request: HttpClientRequest.HttpClientRequest,
middleware?: HttpMiddleware,
) => Effect.Effect<HttpClientResponse.HttpClientResponse, LLMError> ) => Effect.Effect<HttpClientResponse.HttpClientResponse, LLMError>
} }
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/LLM/RequestExecutor") {} export class Service extends Context.Service<Service, Interface>()("@opencode/LLM/RequestExecutor") {}
const BODY_LIMIT = 16_384 const BODY_LIMIT = 16_384
@@ -282,12 +291,20 @@ export const layer: Layer.Layer<Service, never, HttpClient.HttpClient> = Layer.e
Service, Service,
Effect.gen(function* () { Effect.gen(function* () {
const http = yield* HttpClient.HttpClient const http = yield* HttpClient.HttpClient
const executeOnce = (request: HttpClientRequest.HttpClientRequest) => const executeOnce = (request: HttpClientRequest.HttpClientRequest, middleware?: HttpMiddleware) =>
Effect.gen(function* () { Effect.gen(function* () {
const redactedNames = yield* Headers.CurrentRedactedNames const redactedNames = yield* Headers.CurrentRedactedNames
return yield* http if (!middleware)
.execute(request) return yield* http
.pipe(Effect.mapError(toHttpError(redactedNames)), Effect.flatMap(statusError(request, redactedNames))) .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({ return Service.of({
execute: executeOnce, execute: executeOnce,
+1 -1
View File
@@ -23,4 +23,4 @@ export type { ApiKeyMode, AuthOverride, ProviderAuthOption } from "./auth-option
export type { Definition as EndpointFn, EndpointInput } from "./endpoint" export type { Definition as EndpointFn, EndpointInput } from "./endpoint"
export type { Definition as FramingDef } from "./framing" export type { Definition as FramingDef } from "./framing"
export type { Protocol as ProtocolDef } from "./protocol" 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
View File
@@ -3,7 +3,7 @@ import { Headers, HttpClientRequest } from "effect/unstable/http"
import { Auth } from "../auth" import { Auth } from "../auth"
import { render as renderEndpoint } from "../endpoint" import { render as renderEndpoint } from "../endpoint"
import { Framing } from "../framing" import { Framing } from "../framing"
import type { Transport, TransportPrepareInput } from "./index" import type { HttpMiddleware, Transport, TransportPrepareInput } from "./index"
import * as ProviderShared from "../../protocols/shared" import * as ProviderShared from "../../protocols/shared"
import { mergeJsonRecords, type LLMRequest } from "../../schema" import { mergeJsonRecords, type LLMRequest } from "../../schema"
@@ -19,6 +19,7 @@ export interface JsonRequestParts<Body = unknown> {
export interface HttpPrepared<Frame> { export interface HttpPrepared<Frame> {
readonly request: HttpClientRequest.HttpClientRequest readonly request: HttpClientRequest.HttpClientRequest
readonly framing: Framing.Definition<Frame> readonly framing: Framing.Definition<Frame>
readonly middleware?: HttpMiddleware
} }
const applyQuery = (url: string, query: Record<string, string> | undefined) => { 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) => prepare: (prepareInput) =>
Effect.gen(function* () { Effect.gen(function* () {
const parts = yield* jsonRequestParts({ ...prepareInput }) const parts = yield* jsonRequestParts({ ...prepareInput })
const request = { url: parts.url, method: "POST", headers: { ...parts.headers }, body: parts.bodyText } const request = ProviderShared.jsonPost({
yield* (prepareInput.transform?.(request) ?? Effect.void) url: parts.url,
body: parts.bodyText,
headers: parts.headers,
})
return { return {
request: ProviderShared.jsonPost({ request,
url: request.url,
body: request.body ?? "",
headers: Headers.fromInput(request.headers),
}),
framing: input.framing, framing: input.framing,
middleware: prepareInput.middleware,
} }
}), }),
frames: (prepared, request, runtime) => frames: (prepared, request, runtime) =>
Stream.unwrap( Stream.unwrap(
runtime.http runtime.http
.execute(prepared.request) .execute(prepared.request, prepared.middleware)
.pipe( .pipe(
Effect.map((response) => Effect.map((response) =>
prepared.framing.frame( prepared.framing.frame(
+3 -11
View File
@@ -1,7 +1,7 @@
import type { Effect, Stream } from "effect" import type { Effect, Stream } from "effect"
import { Endpoint } from "../endpoint" import { Endpoint } from "../endpoint"
import { Auth } from "../auth" 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 { Interface as WebSocketExecutorInterface } from "./websocket"
import type { LLMError, LLMRequest } from "../../schema" import type { LLMError, LLMRequest } from "../../schema"
@@ -10,15 +10,6 @@ export interface TransportRuntime {
readonly webSocket?: WebSocketExecutorInterface 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> { export interface Transport<Body, Prepared, Frame> {
readonly id: string readonly id: string
readonly prepare: (input: TransportPrepareInput<Body>) => Effect.Effect<Prepared, LLMError> readonly prepare: (input: TransportPrepareInput<Body>) => Effect.Effect<Prepared, LLMError>
@@ -36,8 +27,9 @@ export interface TransportPrepareInput<Body> {
readonly auth: Auth.Definition readonly auth: Auth.Definition
readonly encodeBody: (body: Body) => string readonly encodeBody: (body: Body) => string
readonly headers?: (input: { readonly request: LLMRequest }) => Record<string, string> readonly headers?: (input: { readonly request: LLMRequest }) => Record<string, string>
readonly transform?: HttpRequestTransform readonly middleware?: HttpMiddleware
} }
export * as HttpTransport from "./http" export * as HttpTransport from "./http"
export type { HttpHandler, HttpMiddleware } from "../executor"
export { WebSocketExecutor, WebSocketTransport } from "./websocket" export { WebSocketExecutor, WebSocketTransport } from "./websocket"
+90 -8
View File
@@ -1,6 +1,6 @@
import { describe, expect, test } from "bun:test" import { describe, expect, test } from "bun:test"
import { Effect, Schema } from "effect" import { Effect, Ref, Schema } from "effect"
import { HttpClientRequest } from "effect/unstable/http" import { HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
import { LLM, mergeProviderOptions } from "../src" import { LLM, mergeProviderOptions } from "../src"
import { AnthropicMessages, OpenAIChat } from "../src/protocols" import { AnthropicMessages, OpenAIChat } from "../src/protocols"
import { Auth, LLMClient } from "../src/route" import { Auth, LLMClient } from "../src/route"
@@ -146,12 +146,16 @@ describe("request option precedence", () => {
prompt: "Say hello.", prompt: "Say hello.",
}), }),
{ {
transform: (request) => http: (request, handler) =>
Effect.sync(() => { Effect.gen(function* () {
expect(request.headers.authorization).toBe("Bearer fresh-key") return yield* handler(
request.url = "https://proxy.test/v1/chat/completions" request.pipe(
request.headers["x-plugin"] = "transformed" HttpClientRequest.setUrl("https://proxy.test/v1/chat/completions"),
request.body = JSON.stringify({ transformed: true }) HttpClientRequest.setMethod("PUT"),
HttpClientRequest.setHeader("x-plugin", "transformed"),
HttpClientRequest.bodyText(JSON.stringify({ transformed: true }), "application/custom+json"),
),
)
}), }),
}, },
).pipe( ).pipe(
@@ -160,7 +164,9 @@ describe("request option precedence", () => {
Effect.gen(function* () { Effect.gen(function* () {
const web = yield* HttpClientRequest.toWeb(input.request).pipe(Effect.orDie) const web = yield* HttpClientRequest.toWeb(input.request).pipe(Effect.orDie)
expect(web.url).toBe("https://proxy.test/v1/chat/completions") 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("x-plugin")).toBe("transformed")
expect(web.headers.get("content-type")).toBe("application/custom+json")
expect(decodeJson(input.text)).toEqual({ transformed: true }) expect(decodeJson(input.text)).toEqual({ transformed: true })
return input.respond(sseEvents(deltaChunk({}, "stop")), { return input.respond(sseEvents(deltaChunk({}, "stop")), {
headers: { "content-type": "text/event-stream" }, 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", () => it.effect("applies raw body overlays after protocol lowering", () =>
LLMClient.generate( LLMClient.generate(
LLM.request({ LLM.request({
+32 -3
View File
@@ -194,7 +194,9 @@ export function fromPromise(plugin: Plugin) {
), ),
), ),
refresh: refresh:
refresh === undefined ? undefined : (credential) => Effect.promise(() => refresh(credential)), refresh === undefined
? undefined
: (credential) => Effect.promise(() => refresh(credential)),
}) })
}, },
remove: draft.method.remove, remove: draft.method.remove,
@@ -263,8 +265,35 @@ export function fromPromise(plugin: Plugin) {
), ),
}, },
session: { session: {
hook: (name, callback) => hook: (name, callback) => {
register(host.session.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))), 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) => create: (input) =>
run( run(
host.session.create( host.session.create(
+16 -6
View File
@@ -225,15 +225,25 @@ export const OpenAIPlugin = define({
}) })
} }
}) })
yield* ctx.session.hook("request", (evt) => yield* ctx.session.hook("http", (evt) =>
Effect.sync(() => { Effect.sync(() => {
if (!chatgpt || evt.model.providerID !== Provider.ID.openai) return if (!chatgpt || evt.model.providerID !== Provider.ID.openai) return
const url = new URL(evt.url) const request = evt.request
if (url.origin === "https://api.openai.com") { evt.request = (input) => {
evt.url = `${codexBaseURL}${url.pathname.replace(/^\/v1/, "")}${url.search}` 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
}), }),
) )
+33 -25
View File
@@ -4,7 +4,8 @@ import { LLM, Message, SystemPart, type LLMRequest } from "@opencode-ai/ai"
import type { StreamOptions } from "@opencode-ai/ai/route" import type { StreamOptions } from "@opencode-ai/ai/route"
import type { Content } from "@opencode-ai/schema/tool" import type { Content } from "@opencode-ai/schema/tool"
import { SessionError } from "@opencode-ai/schema/session-error" 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 { makeLocationNode } from "@opencode-ai/util/effect/app-node"
import { App } from "../app" import { App } from "../app"
import { Model } from "../model" import { Model } from "../model"
@@ -48,9 +49,7 @@ interface Prepared {
* One request-scoped execution operation. Unknown, hook-removed, and * One request-scoped execution operation. Unknown, hook-removed, and
* step-limit-violating calls fail individually through the same seam. * step-limit-violating calls fail individually through the same seam.
*/ */
readonly executeTool: ( readonly executeTool: (input: Parameters<Tool.Snapshot["execute"]>[0]) => Effect.Effect<Tool.Result, ExecuteError>
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. */ /** True when this request is the final Step; violating calls are rejected and no continuation follows. */
readonly stepLimitReached: boolean readonly stepLimitReached: boolean
} }
@@ -137,8 +136,7 @@ export const boundImages = (messages: LLMRequest["messages"]) => {
result: { result: {
...part.result, ...part.result,
value: part.result.value.map((item: Content) => { value: part.result.value.map((item: Content) => {
if (item.type !== "file" || !isImage(item.mime) || imageBytes - removed <= IMAGE_BYTES_TARGET) if (item.type !== "file" || !isImage(item.mime) || imageBytes - removed <= IMAGE_BYTES_TARGET) return item
return item
removed += Buffer.byteLength(item.uri) removed += Buffer.byteLength(item.uri)
return { type: "text" as const, text: IMAGE_REMOVED } 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 hookedTools = Object.entries(contextEvent.tools).flatMap(([name, tool]) => {
const registered = toolsByName.get(name) const registered = toolsByName.get(name)
return registered return registered ? [{ ...registered, description: tool.description, inputSchema: tool.input }] : []
? [{ ...registered, description: tool.description, inputSchema: tool.input }]
: []
}) })
const request = LLM.request({ const request = LLM.request({
model, model,
@@ -220,24 +216,37 @@ export const layer = Layer.effect(
toolChoice: stepLimitReached ? "none" : undefined, toolChoice: stepLimitReached ? "none" : undefined,
}) })
const options: StreamOptions = { const options: StreamOptions = {
transform: (request) => http: (request, handler) =>
hooks Effect.gen(function* () {
.trigger("session", "request", { let sent = request
const origins = new WeakMap<Response, HttpClientRequest.HttpClientRequest>()
const web = yield* HttpClientRequest.toWeb(request)
const event = yield* hooks.trigger("session", "http", {
sessionID: session.id, sessionID: session.id,
agent: agent.id, agent: agent.id,
model: resolved.ref, model: resolved.ref,
...request, request: (input) =>
}) Effect.gen(function* () {
.pipe( sent = HttpClientRequest.fromWeb(input)
Effect.tap((event) => if (input.body)
Effect.sync(() => { sent = HttpClientRequest.bodyUint8Array(
request.url = event.url sent,
request.headers = event.headers new Uint8Array(yield* Effect.promise(() => input.arrayBuffer())),
request.body = event.body 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) { if (promptCacheSnapshots) {
const current = PromptCacheDiagnostics.snapshot(request) const current = PromptCacheDiagnostics.snapshot(request)
@@ -257,8 +266,7 @@ export const layer = Layer.effect(
) )
} }
const executeTool: Prepared["executeTool"] = (executeInput) => { const executeTool: Prepared["executeTool"] = (executeInput) => {
if (stepLimitReached) if (stepLimitReached) return new Tool.Error({ message: "Tools are disabled after the maximum agent steps" })
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)) 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 new Tool.Error({ message: `Tool is not available for this request: ${executeInput.call.name}` })
return tools return tools
+85 -14
View File
@@ -1,6 +1,6 @@
import { describe, expect } from "bun:test" import { describe, expect } from "bun:test"
import { Message, SystemPart } from "@opencode-ai/ai" 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 { Agent } from "@opencode-ai/core/agent"
import { Catalog } from "@opencode-ai/core/catalog" import { Catalog } from "@opencode-ai/core/catalog"
import { Model } from "@opencode-ai/core/model" 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({ expect((await ctx.agent.get({ agentID: Agent.ID.make("reviewer") })).data).toMatchObject({
description: "Reviews code", 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 const models = (await ctx.catalog.model.list()).data
expect(models.find((model) => model.providerID === "test" && model.id === "alias")).toMatchObject({ expect(models.find((model) => model.providerID === "test" && model.id === "alias")).toMatchObject({
modelID: "gpt-5", 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<void>()
const interrupted = yield* Deferred.make<void>()
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", () => it.effect("disposes a hook registration on request", () =>
Effect.gen(function* () { Effect.gen(function* () {
const agents = yield* Agent.Service const agents = yield* Agent.Service
@@ -315,19 +388,17 @@ describe("fromPromise", () => {
id: "promise-tool", id: "promise-tool",
setup: async (ctx) => { setup: async (ctx) => {
await ctx.tool.transform((tools) => { await ctx.tool.transform((tools) => {
tools.add( tools.add({
{ name: "hello",
name: "hello", options: { codemode: false },
options: { codemode: false }, description: "Hello",
description: "Hello", input: Schema.Struct({ name: Schema.String }),
input: Schema.Struct({ name: Schema.String }), output: Schema.String,
output: Schema.String, execute: async ({ name }, context) => {
execute: async ({ name }, context) => { await context.progress({ phase: "greeting" })
await context.progress({ phase: "greeting" }) return { output: `Hello, ${name}!` }
return { output: `Hello, ${name}!` }
},
}, },
) })
}) })
}, },
}) })
@@ -29,6 +29,21 @@ function required<T>(value: T | undefined): T {
return value 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", () => { describe("OpenAIPlugin", () => {
it.effect("registers browser and headless ChatGPT OAuth methods", () => it.effect("registers browser and headless ChatGPT OAuth methods", () =>
Effect.gen(function* () { Effect.gen(function* () {
@@ -100,33 +115,9 @@ describe("OpenAIPlugin", () => {
}) })
yield* addPlugin() yield* addPlugin()
const request = yield* (yield* PluginHooks.Service).trigger("session", "request", { const request = yield* http(Provider.ID.openai, "https://api.openai.com/v1/responses")
sessionID: Session.ID.make("ses_test"), const custom = yield* http(Provider.ID.make("custom-openai"), "https://custom.example/v1/responses")
agent: Agent.ID.make("build"), const proxy = yield* http(Provider.ID.openai, "https://proxy.example/v1/responses?region=us")
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 provider = required(yield* catalog.provider.get(Provider.ID.openai)) const provider = required(yield* catalog.provider.get(Provider.ID.openai))
expect(provider.package).toBe("@opencode-ai/ai/providers/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(provider.headers).toMatchObject({ "chatgpt-account-id": "acct_123" })
expect(request.url).toBe("https://chatgpt.com/backend-api/codex/responses") expect(request.url).toBe("https://chatgpt.com/backend-api/codex/responses")
expect(request.headers).toMatchObject({ originator: "opencode", "session-id": "ses_test" }) 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.url).toBe("https://proxy.example/v1/responses?region=us")
expect(proxy.headers).toMatchObject({ originator: "opencode", "session-id": "ses_test" }) 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"))) const eligible = required(yield* catalog.model.get(Provider.ID.openai, Model.ID.make("gpt-5.5")))
@@ -184,21 +175,13 @@ describe("OpenAIPlugin", () => {
}) })
yield* addPlugin() yield* addPlugin()
const request = yield* (yield* PluginHooks.Service).trigger("session", "request", { const request = yield* http(Provider.ID.openai, "https://api.openai.com/v1/responses")
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 model = required(yield* catalog.model.get(Provider.ID.openai, Model.ID.make("gpt-5.5"))) 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.package).toBe("@opencode-ai/ai/providers/openai")
expect(model.enabled).toBe(true) expect(model.enabled).toBe(true)
expect(model.limit).toEqual({ context: 1_050_000, input: 922_000, output: 128_000 }) 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) expect(required(yield* catalog.model.get(Provider.ID.openai, Model.ID.make("gpt-4.1"))).enabled).toBe(true)
}), }),
) )
+4 -4
View File
@@ -1,10 +1,9 @@
import type { SessionApi } from "@opencode-ai/client/effect/api" import type { SessionApi } from "@opencode-ai/client/effect/api"
import type { Message, SystemPart } from "@opencode-ai/ai" 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 { Agent } from "@opencode-ai/schema/agent"
import type { Model } from "@opencode-ai/schema/model" import type { Model } from "@opencode-ai/schema/model"
import type { Session } from "@opencode-ai/schema/session" 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" import type { Hooks } from "./registration.js"
export interface SessionContext { export interface SessionContext {
@@ -16,15 +15,16 @@ export interface SessionContext {
tools: Record<string, { description: string; input: JsonSchema.JsonSchema }> tools: Record<string, { description: string; input: JsonSchema.JsonSchema }>
} }
export interface SessionRequest extends HttpRequest { export interface SessionHttp {
readonly sessionID: Session.ID readonly sessionID: Session.ID
readonly agent: Agent.ID readonly agent: Agent.ID
readonly model: Model.Ref readonly model: Model.Ref
request: (input: Request) => Effect.Effect<Response, Error>
} }
export interface SessionHooks { export interface SessionHooks {
readonly context: SessionContext readonly context: SessionContext
readonly request: SessionRequest readonly http: SessionHttp
} }
export type SessionDomain = Pick< export type SessionDomain = Pick<
+3 -3
View File
@@ -1,6 +1,5 @@
import type { SessionApi } from "@opencode-ai/client/promise/api" import type { SessionApi } from "@opencode-ai/client/promise/api"
import type { Message, SystemPart } from "@opencode-ai/ai" 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 { Agent } from "@opencode-ai/schema/agent"
import type { Model } from "@opencode-ai/schema/model" import type { Model } from "@opencode-ai/schema/model"
import type { Session } from "@opencode-ai/schema/session" import type { Session } from "@opencode-ai/schema/session"
@@ -16,15 +15,16 @@ export interface SessionContext {
tools: Record<string, { description: string; input: JsonSchema.JsonSchema }> tools: Record<string, { description: string; input: JsonSchema.JsonSchema }>
} }
export interface SessionRequest extends HttpRequest { export interface SessionHttp {
readonly sessionID: Session.ID readonly sessionID: Session.ID
readonly agent: Agent.ID readonly agent: Agent.ID
readonly model: Model.Ref readonly model: Model.Ref
request: (input: Request) => Promise<Response>
} }
export interface SessionHooks { export interface SessionHooks {
readonly context: SessionContext readonly context: SessionContext
readonly request: SessionRequest readonly http: SessionHttp
} }
export type SessionDomain = Pick< export type SessionDomain = Pick<
+3 -2
View File
@@ -246,7 +246,8 @@ mutable fields:
| ------------------------------------------- | ------------------------------------------------------------------------------ | | ------------------------------------------- | ------------------------------------------------------------------------------ |
| `ctx.aisdk.hook("sdk", callback)` | `sdk`, after inspecting `model`, `package`, and `options` | | `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.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.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 |
@@ -259,7 +260,7 @@ import { Plugin } from "@opencode-ai/plugin"
export default Plugin.define({ export default Plugin.define({
id: "acme.guards", id: "acme.guards",
setup: async (ctx) => { setup: async (ctx) => {
await ctx.session.hook("request", (event) => { await ctx.session.hook("context", (event) => {
delete event.tools.write delete event.tools.write
}) })