From 9156fe3a98a1fd67163db285f158a646846b51f1 Mon Sep 17 00:00:00 2001 From: Muhammad Mugni Hadi Date: Tue, 26 May 2026 18:12:40 +0700 Subject: [PATCH 1/3] feat(opencode): add experimental OpenAI responses websocket --- packages/llm/src/route/transport/websocket.ts | 22 ++- packages/opencode/src/effect/runtime-flags.ts | 1 + packages/opencode/src/plugin/codex.ts | 147 +++++++++--------- packages/opencode/src/session/llm.ts | 1 + .../src/session/llm/native-request.ts | 15 +- .../src/session/llm/native-runtime.ts | 52 ++++++- .../test/effect/runtime-flags.test.ts | 15 ++ .../opencode/test/session/llm-native.test.ts | 77 +++++++++ 8 files changed, 240 insertions(+), 90 deletions(-) diff --git a/packages/llm/src/route/transport/websocket.ts b/packages/llm/src/route/transport/websocket.ts index 310121420c4f..c00ab10a6091 100644 --- a/packages/llm/src/route/transport/websocket.ts +++ b/packages/llm/src/route/transport/websocket.ts @@ -123,15 +123,21 @@ const webSocketUrl = (value: string) => }) export const open = (input: WebSocketRequest) => - Effect.try({ - try: () => - new (globalThis.WebSocket as unknown as WebSocketConstructorWithHeaders)(input.url, { headers: input.headers }), - catch: (error) => - transportError("open", error instanceof Error ? error.message : "Failed to construct WebSocket", { - url: input.url, - kind: "open", + Effect.logInfo("llm websocket open").pipe( + Effect.annotateLogs({ "llm.websocket.url": input.url }), + Effect.andThen( + Effect.try({ + try: () => + new (globalThis.WebSocket as unknown as WebSocketConstructorWithHeaders)(input.url, { headers: input.headers }), + catch: (error) => + transportError("open", error instanceof Error ? error.message : "Failed to construct WebSocket", { + url: input.url, + kind: "open", + }), }), - }).pipe(Effect.flatMap((ws) => fromWebSocket(ws, input))) + ), + Effect.flatMap((ws) => fromWebSocket(ws, input)), + ) export const layer: Layer.Layer = Layer.succeed(Service, Service.of({ open })) diff --git a/packages/opencode/src/effect/runtime-flags.ts b/packages/opencode/src/effect/runtime-flags.ts index b12a5d5707c4..72ad01384ee0 100644 --- a/packages/opencode/src/effect/runtime-flags.ts +++ b/packages/opencode/src/effect/runtime-flags.ts @@ -54,6 +54,7 @@ export class Service extends ConfigService.Service()("@opencode/Runtime outputTokenMax: positiveInteger("OPENCODE_EXPERIMENTAL_OUTPUT_TOKEN_MAX"), bashDefaultTimeoutMs: positiveInteger("OPENCODE_EXPERIMENTAL_BASH_DEFAULT_TIMEOUT_MS"), experimentalNativeLlm: bool("OPENCODE_EXPERIMENTAL_NATIVE_LLM"), + experimentalOpenAIWebSocket: bool("OPENCODE_EXPERIMENTAL_OPENAI_WEBSOCKET"), client: Config.string("OPENCODE_CLIENT").pipe(Config.withDefault("cli")), }) {} diff --git a/packages/opencode/src/plugin/codex.ts b/packages/opencode/src/plugin/codex.ts index df4b4d0d5c05..a88b6fc6e98d 100644 --- a/packages/opencode/src/plugin/codex.ts +++ b/packages/opencode/src/plugin/codex.ts @@ -10,7 +10,7 @@ const log = Log.create({ service: "plugin.codex" }) const CLIENT_ID = "app_EMoamEEZ73f0CkXaXp7hrann" const ISSUER = "https://auth.openai.com" -const CODEX_API_ENDPOINT = "https://chatgpt.com/backend-api/codex/responses" +export const CODEX_API_ENDPOINT = "https://chatgpt.com/backend-api/codex/responses" const OAUTH_PORT = 1455 const OAUTH_POLLING_SAFETY_MARGIN_MS = 3000 const ALLOWED_MODELS = new Set([ @@ -110,6 +110,28 @@ function buildAuthorizeUrl(redirectUri: string, pkce: PkceCodes, state: string): return `${ISSUER}/oauth/authorize?${params.toString()}` } + +function codexBaseURL(codexApiEndpoint: string): string { + const url = new URL(codexApiEndpoint) + if (url.pathname.endsWith("/responses")) url.pathname = url.pathname.slice(0, -"/responses".length) + return url.toString().replace(/\/$/, "") +} + +function headersInit(init: HeadersInit | undefined): Headers { + const headers = new Headers() + if (!init) return headers + if (init instanceof Headers) { + init.forEach((value, key) => headers.set(key, value)) + return headers + } + if (Array.isArray(init)) { + for (const [key, value] of init) if (value !== undefined) headers.set(key, String(value)) + return headers + } + for (const [key, value] of Object.entries(init)) if (value !== undefined) headers.set(key, String(value)) + return headers +} + interface TokenResponse { id_token: string access_token: string @@ -419,85 +441,58 @@ export async function CodexAuthPlugin(input: PluginInput, options: CodexAuthPlug }> | undefined - return { - apiKey: OAUTH_DUMMY_KEY, - async fetch(requestInput: RequestInfo | URL, init?: RequestInit) { - // Remove dummy API key authorization header - if (init?.headers) { - if (init.headers instanceof Headers) { - init.headers.delete("authorization") - init.headers.delete("Authorization") - } else if (Array.isArray(init.headers)) { - init.headers = init.headers.filter(([key]) => key.toLowerCase() !== "authorization") - } else { - delete init.headers["authorization"] - delete init.headers["Authorization"] - } - } - - const currentAuth = await getAuth() - if (currentAuth.type !== "oauth") return fetch(requestInput, init) - - // Cast to include accountId field - const authWithAccount = currentAuth as typeof currentAuth & { accountId?: string } - - // Check if token needs refresh - if (!currentAuth.access || currentAuth.expires < Date.now()) { - if (!refreshPromise) { - log.info("refreshing codex access token") - refreshPromise = refreshAccessToken(currentAuth.refresh, issuer) - .then(async (tokens) => { - const accountId = extractAccountId(tokens) || authWithAccount.accountId - await input.client.auth.set({ - path: { id: "openai" }, - body: { - type: "oauth", - refresh: tokens.refresh_token, - access: tokens.access_token, - expires: Date.now() + (tokens.expires_in ?? 3600) * 1000, - ...(accountId && { accountId }), - }, - }) - return { + const codexAuthHeaders = async (init?: HeadersInit) => { + const currentAuth = await getAuth() + if (currentAuth.type !== "oauth") return headersInit(init) + + const authWithAccount = currentAuth as typeof currentAuth & { accountId?: string } + if (!currentAuth.access || currentAuth.expires < Date.now()) { + if (!refreshPromise) { + log.info("refreshing codex access token") + refreshPromise = refreshAccessToken(currentAuth.refresh, issuer) + .then(async (tokens) => { + const accountId = extractAccountId(tokens) || authWithAccount.accountId + await input.client.auth.set({ + path: { id: "openai" }, + body: { + type: "oauth", + refresh: tokens.refresh_token, access: tokens.access_token, - accountId, - } - }) - .finally(() => { - refreshPromise = undefined + expires: Date.now() + (tokens.expires_in ?? 3600) * 1000, + ...(accountId && { accountId }), + }, }) - } - - const refreshed = await refreshPromise - currentAuth.access = refreshed.access - authWithAccount.accountId = refreshed.accountId - } - - // Build headers - const headers = new Headers() - if (init?.headers) { - if (init.headers instanceof Headers) { - init.headers.forEach((value, key) => headers.set(key, value)) - } else if (Array.isArray(init.headers)) { - for (const [key, value] of init.headers) { - if (value !== undefined) headers.set(key, String(value)) - } - } else { - for (const [key, value] of Object.entries(init.headers)) { - if (value !== undefined) headers.set(key, String(value)) - } - } + return { access: tokens.access_token, accountId } + }) + .finally(() => { + refreshPromise = undefined + }) } - // Set authorization header with access token - headers.set("authorization", `Bearer ${currentAuth.access}`) - - // Set ChatGPT-Account-Id header for organization subscriptions - if (authWithAccount.accountId) { - headers.set("ChatGPT-Account-Id", authWithAccount.accountId) - } + const refreshed = await refreshPromise + currentAuth.access = refreshed.access + authWithAccount.accountId = refreshed.accountId + } + + const headers = headersInit(init) + headers.delete("authorization") + headers.delete("Authorization") + headers.set("authorization", `Bearer ${currentAuth.access}`) + if (authWithAccount.accountId) headers.set("ChatGPT-Account-Id", authWithAccount.accountId) + return headers + } - // Rewrite URL to Codex endpoint + return { + apiKey: OAUTH_DUMMY_KEY, + codexWebSocket: { + baseURL: codexBaseURL(codexApiEndpoint), + async headers(init?: HeadersInit) { + return Object.fromEntries((await codexAuthHeaders(init)).entries()) + }, + }, + async fetch(requestInput: RequestInfo | URL, init?: RequestInit) { + const currentAuth = await getAuth() + if (currentAuth.type !== "oauth") return fetch(requestInput, init) const parsed = requestInput instanceof URL ? requestInput @@ -509,7 +504,7 @@ export async function CodexAuthPlugin(input: PluginInput, options: CodexAuthPlug return fetch(url, { ...init, - headers, + headers: await codexAuthHeaders(init?.headers), }) }, } diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index ea2efc99d007..0e026f45c631 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -233,6 +233,7 @@ const live: Layer.Layer< providerOptions: prepared.params.options, headers: prepared.headers, abort: input.abort, + experimentalOpenAIWebSocket: flags.experimentalOpenAIWebSocket, }) if (native.type === "supported") { yield* Effect.logInfo("llm runtime selected").pipe( diff --git a/packages/opencode/src/session/llm/native-request.ts b/packages/opencode/src/session/llm/native-request.ts index b7f30e24c362..ddc7bc73b4a3 100644 --- a/packages/opencode/src/session/llm/native-request.ts +++ b/packages/opencode/src/session/llm/native-request.ts @@ -1,4 +1,5 @@ import type { JsonSchema, LLMRequest, ProviderMetadata } from "@opencode-ai/llm" +import type { AuthShape } from "@opencode-ai/llm/route" import { LLM, Message, SystemPart, ToolCallPart, ToolDefinition, ToolResultPart } from "@opencode-ai/llm" import { AmazonBedrock, @@ -32,6 +33,8 @@ export type RequestInput = { readonly maxOutputTokens?: number readonly providerOptions?: LLMRequest["providerOptions"] readonly headers?: Record + readonly routeAuth?: AuthShape + readonly useOpenAIWebSocket?: boolean } const providerMetadata = (value: unknown): ProviderMetadata | undefined => { @@ -153,8 +156,7 @@ const requireBaseURL = (model: Provider.Model, url: string | undefined) => { export const model = (input: Provider.Model | RequestInput, headers?: Record) => { const model = "model" in input ? input.model : input const url = baseURL(input) - const options = { - ...("model" in input && input.apiKey ? { apiKey: input.apiKey } : {}), + const baseOptions = { ...(url ? { baseURL: url } : {}), headers: Object.keys({ ...model.headers, ...headers }).length === 0 ? undefined : { ...model.headers, ...headers }, limits: { @@ -162,7 +164,14 @@ export const model = (input: Provider.Model | RequestInput, headers?: Record readonly headers: Record readonly abort: AbortSignal + readonly experimentalOpenAIWebSocket?: boolean } export function status(input: Pick): RuntimeStatus { @@ -82,7 +83,9 @@ export function stream(input: StreamInput): StreamResult { request: LLMNative.request({ model: input.model, apiKey: current.apiKey, - baseURL: current.baseURL, + baseURL: openAIWebSocketBaseURL(input) ?? current.baseURL, + routeAuth: openAIWebSocketAuth(input), + useOpenAIWebSocket: shouldUseOpenAIWebSocket(input), messages: ProviderTransform.message(input.messages, input.model, input.providerOptions ?? {}), toolChoice: input.toolChoice, temperature: input.temperature, @@ -101,6 +104,49 @@ export function stream(input: StreamInput): StreamResult { } } + +function shouldUseOpenAIWebSocket(input: Pick) { + if (!input.experimentalOpenAIWebSocket) return false + if (input.model.providerID !== "openai" || input.model.api.npm !== "@ai-sdk/openai") return false + if (input.auth?.type === "oauth" && !codexWebSocket(input.provider.options)) return false + return true +} + +function codexWebSocket(options: Record): + | { + readonly baseURL?: string + readonly headers?: (headers?: Headers.Headers) => Promise> + } + | undefined { + const value = options["codexWebSocket"] + if (!isRecord(value)) return undefined + const headers = value.headers + if (headers !== undefined && typeof headers !== "function") return undefined + return { + baseURL: typeof value.baseURL === "string" ? value.baseURL : undefined, + headers: + typeof headers === "function" + ? (headers as (headers?: Headers.Headers) => Promise>) + : undefined, + } +} + +function openAIWebSocketBaseURL(input: Pick) { + if (!shouldUseOpenAIWebSocket(input)) return undefined + return codexWebSocket(input.provider.options)?.baseURL +} + +function openAIWebSocketAuth(input: Pick) { + if (!shouldUseOpenAIWebSocket(input)) return undefined + const helper = codexWebSocket(input.provider.options) + if (!helper?.headers) return undefined + return RouteAuth.custom((authInput) => + Effect.promise(() => helper.headers?.(authInput.headers) ?? Promise.resolve({})).pipe( + Effect.map((headers) => Headers.setAll(authInput.headers, headers)), + ), + ) +} + function providerFetch(input: Pick): typeof globalThis.fetch | undefined { if (input.provider.id !== "openai" || input.auth?.type !== "oauth") return undefined const value: unknown = input.provider.options.fetch diff --git a/packages/opencode/test/effect/runtime-flags.test.ts b/packages/opencode/test/effect/runtime-flags.test.ts index b044d07f23e1..a3d234de8a79 100644 --- a/packages/opencode/test/effect/runtime-flags.test.ts +++ b/packages/opencode/test/effect/runtime-flags.test.ts @@ -63,6 +63,7 @@ describe("RuntimeFlags", () => { expect(flags.experimentalWorkspaces).toBe(true) expect(flags.experimentalIconDiscovery).toBe(true) expect(flags.experimentalNativeLlm).toBe(false) + expect(flags.experimentalOpenAIWebSocket).toBe(false) expect(flags.client).toBe("desktop") }), ) @@ -91,6 +92,18 @@ describe("RuntimeFlags", () => { }), ) + it.effect("enables OpenAI WebSocket via dedicated flag only", () => + Effect.gen(function* () { + const explicit = yield* readFlags.pipe( + Effect.provide(fromConfig({ OPENCODE_EXPERIMENTAL_OPENAI_WEBSOCKET: "true" })), + ) + const umbrella = yield* readFlags.pipe(Effect.provide(fromConfig({ OPENCODE_EXPERIMENTAL: "true" }))) + + expect(explicit.experimentalOpenAIWebSocket).toBe(true) + expect(umbrella.experimentalOpenAIWebSocket).toBe(false) + }), + ) + it.effect("layer accepts partial test overrides and fills defaults from Config definitions", () => Effect.gen(function* () { const flags = yield* readFlags.pipe( @@ -110,6 +123,7 @@ describe("RuntimeFlags", () => { expect(flags.enableExa).toBe(false) expect(flags.experimentalIconDiscovery).toBe(false) expect(flags.experimentalOxfmt).toBe(false) + expect(flags.experimentalOpenAIWebSocket).toBe(false) expect(flags.outputTokenMax).toBeUndefined() expect(flags.bashDefaultTimeoutMs).toBe(1_000) expect(flags.enableExperimentalModels).toBe(false) @@ -355,6 +369,7 @@ describe("RuntimeFlags", () => { expect(flags.enableExa).toBe(false) expect(flags.experimentalIconDiscovery).toBe(false) expect(flags.experimentalOxfmt).toBe(false) + expect(flags.experimentalOpenAIWebSocket).toBe(false) expect(flags.outputTokenMax).toBeUndefined() expect(flags.bashDefaultTimeoutMs).toBeUndefined() expect(flags.client).toBe("cli") diff --git a/packages/opencode/test/session/llm-native.test.ts b/packages/opencode/test/session/llm-native.test.ts index 076d4c9f789a..2142b463f248 100644 --- a/packages/opencode/test/session/llm-native.test.ts +++ b/packages/opencode/test/session/llm-native.test.ts @@ -372,6 +372,83 @@ describe("session.llm-native.request", () => { expect(openrouter.route.endpoint.baseURL).toBe("https://openrouter.ai/api/v1") }) + test("selects OpenAI Responses WebSocket route when requested", () => { + const request = LLMNative.request({ + model: baseModel, + apiKey: "test-openai-key", + messages: [{ role: "user", content: "hello" }], + useOpenAIWebSocket: true, + }) + + expect(request.model).toMatchObject({ + id: "gpt-5-mini", + provider: "openai", + route: { id: "openai-responses-websocket" }, + }) + }) + + test("streams Codex OAuth requests through OpenAI Responses WebSocket when enabled", async () => { + const opened: Array<{ url: string; headers: Record; message: string }> = [] + const webSocketLayer = Layer.succeed( + WebSocketExecutor.Service, + WebSocketExecutor.Service.of({ + open: (input) => + Effect.succeed({ + sendText: (message) => + Effect.sync(() => { + opened.push({ url: input.url, headers: input.headers, message }) + }), + messages: Stream.empty, + close: Effect.void, + }), + }), + ) + + await Effect.runPromise( + Effect.gen(function* () { + const llmClient = yield* LLMClient.Service + const native = LLMNativeRuntime.stream({ + model: baseModel, + provider: { + ...providerInfo, + options: { + apiKey: OAUTH_DUMMY_KEY, + fetch: async () => new Response(), + codexWebSocket: { + baseURL: "https://chatgpt.test/backend-api/codex", + async headers(headers?: Record) { + return { ...headers, authorization: "Bearer access", "ChatGPT-Account-Id": "acc-123" } + }, + }, + }, + }, + auth: { type: "oauth", refresh: "refresh", access: "access", expires: Date.now() + 60_000 }, + llmClient, + messages: [{ role: "user", content: "hello" }], + tools: {}, + headers: {}, + abort: new AbortController().signal, + experimentalOpenAIWebSocket: true, + }) + expect(native.type).toBe("supported") + if (native.type === "unsupported") throw new Error(native.reason) + yield* native.stream.pipe(Stream.runCollect) + }).pipe( + Effect.provide(LLMClient.layer.pipe(Layer.provide(Layer.mergeAll(RequestExecutor.defaultLayer, webSocketLayer)))), + ), + ) + + expect(opened).toHaveLength(1) + expect(opened[0]?.url).toBe("wss://chatgpt.test/backend-api/codex/responses") + expect(opened[0]?.headers.authorization).toBe("Bearer access") + expect(opened[0]?.headers["ChatGPT-Account-Id"] ?? opened[0]?.headers["chatgpt-account-id"]).toBe("acc-123") + expect(JSON.parse(opened[0]?.message ?? "{}")).toMatchObject({ + type: "response.create", + model: "gpt-5-mini", + input: [{ role: "user", content: [{ type: "input_text", text: "hello" }] }], + }) + }) + test("fails fast for unsupported provider packages", () => { expect(() => LLMNative.request({ From 4ba8ab8d4cbc3c9f54d63dbaac3348bdcbd1ea9f Mon Sep 17 00:00:00 2001 From: Muhammad Mugni Hadi Date: Tue, 26 May 2026 19:42:52 +0700 Subject: [PATCH 2/3] fix(llm): redact websocket open URLs --- packages/llm/src/route/executor.ts | 46 ++++-------------- packages/llm/src/route/redaction.ts | 40 ++++++++++++++++ packages/llm/src/route/transport/websocket.ts | 3 +- packages/llm/test/executor.test.ts | 47 ++++++++++++++++++- 4 files changed, 96 insertions(+), 40 deletions(-) create mode 100644 packages/llm/src/route/redaction.ts diff --git a/packages/llm/src/route/executor.ts b/packages/llm/src/route/executor.ts index 815b2c289c8a..ae4651eb97c8 100644 --- a/packages/llm/src/route/executor.ts +++ b/packages/llm/src/route/executor.ts @@ -22,6 +22,15 @@ import { TransportReason, UnknownProviderReason, } from "../schema" +import { + REDACTED, + REDACT_JSON_FIELD, + REDACT_QUERY_FIELD, + isSensitiveHeaderName, + isSensitiveQueryName, + redactHeaders, + redactUrl, +} from "./redaction" export interface Interface { readonly execute: ( @@ -35,43 +44,6 @@ const BODY_LIMIT = 16_384 const MAX_RETRIES = 2 const BASE_DELAY_MS = 500 const MAX_DELAY_MS = 10_000 -const REDACTED = "" - -// One source of truth for what counts as a sensitive name across headers, -// URL query keys, and field names embedded inside request/response bodies. -// -// `SENSITIVE_NAME` is used as both a substring matcher (for free-form header -// names like `Authorization` / `X-API-Key`) and as the body-field alternation -// list. `SHORT_QUERY_NAME` covers anchored short keys like `?key=…` / `?sig=…` -// that are too generic to redact substring-style without false positives. -const SENSITIVE_NAME_SOURCE = - "authorization|api[-_]?key|access[-_]?token|refresh[-_]?token|id[-_]?token|token|secret|credential|signature|x-amz-signature" -const SENSITIVE_NAME = new RegExp(SENSITIVE_NAME_SOURCE, "i") -const SHORT_QUERY_NAME = /^(key|sig)$/i -const SENSITIVE_BODY_FIELD = new RegExp(`(?:${SENSITIVE_NAME_SOURCE}|key)`, "i") -const REDACT_JSON_FIELD = new RegExp(`("(?:${SENSITIVE_BODY_FIELD.source})"\\s*:\\s*)"[^"]*"`, "gi") -const REDACT_QUERY_FIELD = new RegExp(`((?:${SENSITIVE_BODY_FIELD.source})=)[^&\\s"]+`, "gi") - -const isSensitiveHeaderName = (name: string) => SENSITIVE_NAME.test(name) - -const isSensitiveQueryName = (name: string) => isSensitiveHeaderName(name) || SHORT_QUERY_NAME.test(name) - -const redactHeaders = (headers: Headers.Headers, redactedNames: ReadonlyArray) => - Object.fromEntries( - Object.entries(Headers.redact(headers, [...redactedNames, SENSITIVE_NAME])).map(([name, value]) => [ - name, - String(value), - ]), - ) - -const redactUrl = (value: string) => { - if (!URL.canParse(value)) return REDACTED - const url = new URL(value) - url.searchParams.forEach((_, key) => { - if (isSensitiveQueryName(key)) url.searchParams.set(key, REDACTED) - }) - return url.toString() -} const normalizedHeaders = (headers: Headers.Headers) => Object.fromEntries(Object.entries(headers).map(([key, value]) => [key.toLowerCase(), value])) diff --git a/packages/llm/src/route/redaction.ts b/packages/llm/src/route/redaction.ts new file mode 100644 index 000000000000..0023e34f3f60 --- /dev/null +++ b/packages/llm/src/route/redaction.ts @@ -0,0 +1,40 @@ +import { Headers } from "effect/unstable/http" + +export const REDACTED = "" + +// One source of truth for what counts as a sensitive name across headers, +// URL query keys, and field names embedded inside request/response bodies. +// +// `SENSITIVE_NAME` is used as both a substring matcher (for free-form header +// names like `Authorization` / `X-API-Key`) and as the body-field alternation +// list. `SHORT_QUERY_NAME` covers anchored short keys like `?key=...` / `?sig=...` +// that are too generic to redact substring-style without false positives. +const SENSITIVE_NAME_SOURCE = + "authorization|api[-_]?key|access[-_]?token|refresh[-_]?token|id[-_]?token|token|secret|credential|signature|x-amz-signature" +const SENSITIVE_NAME = new RegExp(SENSITIVE_NAME_SOURCE, "i") +const SHORT_QUERY_NAME = /^(key|sig)$/i +const SENSITIVE_BODY_FIELD = new RegExp(`(?:${SENSITIVE_NAME_SOURCE}|key)`, "i") + +export const REDACT_JSON_FIELD = new RegExp(`("(?:${SENSITIVE_BODY_FIELD.source})"\\s*:\\s*)"[^"]*"`, "gi") +export const REDACT_QUERY_FIELD = new RegExp(`((?:${SENSITIVE_BODY_FIELD.source})=)[^&\\s"]+`, "gi") + +export const isSensitiveHeaderName = (name: string) => SENSITIVE_NAME.test(name) + +export const isSensitiveQueryName = (name: string) => isSensitiveHeaderName(name) || SHORT_QUERY_NAME.test(name) + +export const redactHeaders = (headers: Headers.Headers, redactedNames: ReadonlyArray) => + Object.fromEntries( + Object.entries(Headers.redact(headers, [...redactedNames, SENSITIVE_NAME])).map(([name, value]) => [ + name, + String(value), + ]), + ) + +export const redactUrl = (value: string) => { + if (!URL.canParse(value)) return REDACTED + const url = new URL(value) + url.searchParams.forEach((_, key) => { + if (isSensitiveQueryName(key)) url.searchParams.set(key, REDACTED) + }) + return url.toString() +} diff --git a/packages/llm/src/route/transport/websocket.ts b/packages/llm/src/route/transport/websocket.ts index c00ab10a6091..df58f25129b1 100644 --- a/packages/llm/src/route/transport/websocket.ts +++ b/packages/llm/src/route/transport/websocket.ts @@ -1,6 +1,7 @@ import { Cause, Context, Effect, Layer, Queue, Stream } from "effect" import { Headers } from "effect/unstable/http" import { LLMError, TransportReason } from "../../schema" +import { redactUrl } from "../redaction" import * as HttpTransport from "./http" import type { Transport } from "./index" @@ -124,7 +125,7 @@ const webSocketUrl = (value: string) => export const open = (input: WebSocketRequest) => Effect.logInfo("llm websocket open").pipe( - Effect.annotateLogs({ "llm.websocket.url": input.url }), + Effect.annotateLogs({ "llm.websocket.url": redactUrl(input.url) }), Effect.andThen( Effect.try({ try: () => diff --git a/packages/llm/test/executor.test.ts b/packages/llm/test/executor.test.ts index fb8a46707875..ec10b89fdc4a 100644 --- a/packages/llm/test/executor.test.ts +++ b/packages/llm/test/executor.test.ts @@ -1,9 +1,9 @@ import { describe, expect } from "bun:test" -import { Effect, Fiber, Layer, Random, Ref } from "effect" +import { Effect, Fiber, Layer, Logger, Random, Ref, References } from "effect" import * as TestClock from "effect/testing/TestClock" import { Headers, HttpClient, HttpClientRequest, HttpClientResponse } from "effect/unstable/http" import { LLM, LLMError } from "../src" -import { LLMClient, RequestExecutor } from "../src/route" +import { LLMClient, RequestExecutor, WebSocketExecutor } from "../src/route" import * as OpenAIChat from "../src/protocols/openai-chat" import { dynamicResponse } from "./lib/http" import { deltaChunk } from "./lib/openai-chunks" @@ -73,6 +73,49 @@ const expectLLMError = (error: unknown) => { const errorHttp = (error: LLMError) => ("http" in error.reason ? error.reason.http : undefined) describe("RequestExecutor", () => { + it.effect("redacts sensitive WebSocket URL query params in open logs", () => + Effect.gen(function* () { + class OpenWebSocket { + static readonly OPEN = 1 + static readonly CLOSING = 2 + static readonly CLOSED = 3 + readyState = OpenWebSocket.OPEN + addEventListener() {} + removeEventListener() {} + send() {} + close() { + this.readyState = OpenWebSocket.CLOSED + } + } + + const annotations: Array> = [] + const logger = Logger.make((options) => { + annotations.push(options.fiber.getRef(References.CurrentLogAnnotations)) + }) + + yield* Effect.acquireUseRelease( + Effect.sync(() => { + const original = globalThis.WebSocket + globalThis.WebSocket = OpenWebSocket as unknown as typeof globalThis.WebSocket + return original + }), + () => + WebSocketExecutor.open({ + url: "wss://provider.test/realtime?api_key=query-secret-123&key=short-secret&debug=1", + headers: Headers.empty, + }).pipe(Effect.flatMap((connection) => connection.close)), + (original) => + Effect.sync(() => { + globalThis.WebSocket = original + }), + ).pipe(Effect.provide(Logger.layer([logger]))) + + expect(annotations.find((item) => item["llm.websocket.url"])?.["llm.websocket.url"]).toBe( + "wss://provider.test/realtime?api_key=%3Credacted%3E&key=%3Credacted%3E&debug=1", + ) + }), + ) + it.effect("returns redacted diagnostics for retryable rate limits", () => Effect.gen(function* () { const executor = yield* RequestExecutor.Service From 6ae7af4b628f8994c7e6fb1c3abce1266d00ea59 Mon Sep 17 00:00:00 2001 From: Muhammad Mugni Hadi Date: Tue, 26 May 2026 19:44:34 +0700 Subject: [PATCH 3/3] fix(llm): preserve responses tool result narrowing --- packages/llm/src/protocols/openai-responses.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/packages/llm/src/protocols/openai-responses.ts b/packages/llm/src/protocols/openai-responses.ts index ad673a263f4b..9973cfbe45ea 100644 --- a/packages/llm/src/protocols/openai-responses.ts +++ b/packages/llm/src/protocols/openai-responses.ts @@ -320,7 +320,8 @@ const lowerToolResultOutput = Effect.fn("OpenAIResponses.lowerToolResultOutput") // Text/json/error results are encoded as a plain string for backward // compatibility with existing cassettes and provider expectations. if (part.result.type !== "content") return ProviderShared.toolResultText(part) - return yield* Effect.forEach(part.result.value, lowerToolResultContentItem) + const content: ReadonlyArray = part.result.value + return yield* Effect.forEach(content, lowerToolResultContentItem) }) const lowerMessages = Effect.fn("OpenAIResponses.lowerMessages")(function* (request: LLMRequest) {