From 089f9f07e8281d99232f92e13cc8629e43a2590a Mon Sep 17 00:00:00 2001 From: halysondev Date: Fri, 25 Sep 2026 23:11:46 -0300 Subject: [PATCH 1/6] feat(plugins): load local plugins and expose an upstream rewrite slot `ocx start` now imports `$OPENCODEX_HOME/plugins/*.{ts,js,mjs}` before the server binds. A plugin default-exports `{ name?, setup(context) }` and can register an upstream rewriter that runs synchronously at the physical send: fetchWithHeaderTimeout, fetchWithAttemptDeadline and the CodexWsSession dial. This lets a local sidecar (for example a compression proxy) sit in front of providers without editing the core. - upstream-hooks.ts imports nothing; with no plugin the send is untouched. - A throwing rewriter is disabled and the send continues unmodified. - The loader refuses files owned by other users or writable by group/others, bounds setup to 5s, removes hooks from a failed setup, and honours OCX_PLUGINS=0. - Documented in structure/ops/plugins.md and docs-site guides/local-plugins. Co-Authored-By: Claude Opus 5.5 --- docs-site/astro.config.mjs | 1 + .../src/content/docs/guides/local-plugins.md | 81 +++++++++ scripts/test-layout/layout.json | 2 + src/cli/index.ts | 1 + src/lib/upstream-retry.ts | 6 +- src/plugins/loader.ts | 156 ++++++++++++++++++ src/plugins/upstream-hooks.ts | 88 ++++++++++ src/server/responses/codex-ws-session.ts | 5 +- src/server/responses/fetch-helpers.ts | 6 +- structure/INDEX.md | 2 + structure/manifest.json | 9 + structure/ops/plugins.md | 40 +++++ tests/fixtures/test-layout-expected.json | 2 + tests/lib/plugin-loader.test.ts | 96 +++++++++++ tests/lib/plugin-upstream-hooks.test.ts | 64 +++++++ .../responses-fetch-helpers-boundary.test.ts | 2 + 16 files changed, 556 insertions(+), 5 deletions(-) create mode 100644 docs-site/src/content/docs/guides/local-plugins.md create mode 100644 src/plugins/loader.ts create mode 100644 src/plugins/upstream-hooks.ts create mode 100644 structure/ops/plugins.md create mode 100644 tests/lib/plugin-loader.test.ts create mode 100644 tests/lib/plugin-upstream-hooks.test.ts diff --git a/docs-site/astro.config.mjs b/docs-site/astro.config.mjs index fc4162f8879..628a59942d5 100644 --- a/docs-site/astro.config.mjs +++ b/docs-site/astro.config.mjs @@ -110,6 +110,7 @@ export default defineConfig({ { label: "Integrations", translations: { fr: "Intégrations", ko: "연동", "zh-CN": "集成", "zh-TW": "整合", ru: "Интеграции", ja: "連携", tr: "Entegrasyonlar" }, slug: "guides/integrations" }, { label: "MiniMax clients", translations: { fr: "Clients MiniMax", ko: "MiniMax 클라이언트", "zh-CN": "MiniMax 客户端", "zh-TW": "MiniMax 客戶端", ru: "Клиенты MiniMax", ja: "MiniMax クライアント", tr: "MiniMax İstemcileri" }, slug: "guides/minimax" }, { label: "Sidecars: Web Search & Vision", translations: { fr: "Services auxiliaires : recherche web et vision", ko: "사이드카: 웹 검색 & 비전", "zh-CN": "边车:网络搜索与视觉", "zh-TW": "邊車:網路搜尋與視覺", ru: "Сайдкары: веб-поиск и зрение", ja: "サイドカー: ウェブ検索 & ビジョン", tr: "Sidecar'lar: Web Arama ve Görme" }, slug: "guides/sidecars" }, + { label: "Local Plugins", translations: { fr: "Plugins locaux", ko: "로컬 플러그인", "zh-CN": "本地插件", "zh-TW": "本機外掛", ru: "Локальные плагины", ja: "ローカルプラグイン", tr: "Yerel Eklentiler" }, slug: "guides/local-plugins" }, { label: "Image Bridge", translations: { fr: "Pont d’images", ko: "이미지 브릿지", "zh-CN": "图像桥接", "zh-TW": "圖像橋接", ru: "Image Bridge", ja: "画像ブリッジ", tr: "Image Bridge" }, slug: "guides/image-bridge" }, { label: "Video Bridge", translations: { fr: "Pont vidéo", ko: "비디오 브릿지", "zh-CN": "视频桥接", "zh-TW": "影片橋接", ru: "Video Bridge", ja: "動画ブリッジ", tr: "Video Bridge" }, slug: "guides/video-bridge" }, { label: "Web Dashboard", translations: { fr: "Tableau de bord web", ko: "웹 대시보드", "zh-CN": "网页控制台", "zh-TW": "網頁儀表板", ru: "Веб-дашборд", ja: "ウェブダッシュボード", tr: "Web Kontrol Paneli" }, slug: "guides/web-dashboard" }, diff --git a/docs-site/src/content/docs/guides/local-plugins.md b/docs-site/src/content/docs/guides/local-plugins.md new file mode 100644 index 00000000000..64969434dbf --- /dev/null +++ b/docs-site/src/content/docs/guides/local-plugins.md @@ -0,0 +1,81 @@ +--- +title: Local Plugins +description: Load your own code into the proxy at startup to rewrite provider sends, for example to put a local compression proxy in front of providers. +--- + +A local plugin is a TypeScript or JavaScript file that `ocx start` loads before the proxy begins +serving. It can see every provider send just before it leaves the process and redirect it or add +headers — enough to place a local sidecar (a compression proxy, a recorder) in front of providers +without changing opencodex itself. + +Plugins are local to one install. opencodex does not download, update or sign them. + +## Where plugins live + +Put plugin files in `plugins/` inside the opencodex home (`~/.opencodex/plugins/`, or +`$OPENCODEX_HOME/plugins/` when that variable is set): + +```text +~/.opencodex/plugins/ + my-sidecar.ts +``` + +- Files ending in `.ts`, `.js` or `.mjs` are loaded in name order. +- Names starting with `.` or `_`, and `*.d.ts`, are ignored — rename a plugin to `_my-sidecar.ts` + to switch it off. +- The directory is optional. Without it nothing is loaded. +- A plugin runs inside the proxy with your credentials, so opencodex refuses a file owned by another + user or writable by group or others. Fix it with `chmod go-w ~/.opencodex/plugins/*`. + +Restart the proxy after adding, changing or removing a plugin (`ocx service restart`, or stop and +start `ocx start`). Each loaded plugin prints a `Plugin loaded: ` line at startup; a skipped +plugin prints the reason. + +To start once without plugins, set `OCX_PLUGINS=0`. + +## Writing a plugin + +A plugin default-exports an object with an optional `name` and a `setup` function. `setup` receives a +context and has five seconds to finish. + +```ts +interface UpstreamTarget { + url: string; // absolute upstream URL; assign a new one to redirect + headers: Headers; // outbound headers, including credentials — never log them + readonly transport: "http" | "websocket"; +} + +export default { + name: "my-sidecar", + setup(ctx: { + log(message: string): void; + registerUpstreamRewriter(rewrite: (target: UpstreamTarget) => void): void; + onShutdown(teardown: () => void): void; + }) { + ctx.registerUpstreamRewriter(target => { + const upstream = new URL(target.url); + if (!upstream.pathname.endsWith("/chat/completions")) return; + target.url = `http://127.0.0.1:9000${upstream.pathname}`; + target.headers.set("x-original-origin", upstream.origin); + }); + }, +}; +``` + +The context also carries `name`, `configDir` (the opencodex home) and `pluginDir`. + +Plugins cannot import opencodex modules — in the packaged binary they are not on disk. Declare the +small interfaces you need locally, as above. + +## How rewrites behave + +- The rewriter runs synchronously on every provider send over HTTP and on the Codex WebSocket + connection. Keep it fast; do network checks (health probes) in the background and read a cached + result in the rewriter. +- It runs after opencodex has chosen the provider, account and route, so it does not change routing, + account selection, retries or request logs. +- A redirected send goes to the host you chose. That host sees the request exactly as the provider + would, credentials included. +- If a rewriter throws, opencodex disables it for the rest of the process and sends the request + unmodified. If `setup` throws or times out, the plugin is skipped and anything it registered is + removed; other plugins and the proxy start normally. diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 47991be1119..87e1a6ac8b0 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -170,6 +170,8 @@ "explicit": { "openai-chat-serialized-tool-call-scaling.test.ts": "adapters/openai", "coding-agent-json-lines-scaling.test.ts": "providers", + "plugin-loader.test.ts": "lib", + "plugin-upstream-hooks.test.ts": "lib", "usage-snapshot-digest-reuse.test.ts": "usage", "release-desktop-scripts.test.ts": "ci-workflows", "installed-gate-drivers.test.ts": "ci-workflows", diff --git a/src/cli/index.ts b/src/cli/index.ts index 827dbd6701b..11a8c9a2523 100755 --- a/src/cli/index.ts +++ b/src/cli/index.ts @@ -529,6 +529,7 @@ async function handleStart(options: { block?: boolean } = {}) { const readinessGate = createReadinessGate(); const localAttestationSecret = createLocalAttestationSecret(); const config = loadConfig(); + await (await import("../plugins/loader")).loadAndReportOcxPlugins(); let server: ReturnType; for (let attempt = 0; ; attempt++) { try { diff --git a/src/lib/upstream-retry.ts b/src/lib/upstream-retry.ts index 2b2b546388c..7594b0b8458 100644 --- a/src/lib/upstream-retry.ts +++ b/src/lib/upstream-retry.ts @@ -17,6 +17,7 @@ */ import { clearableDeadline } from "./abort"; import { redactSecretString } from "./redact"; +import { rewriteUpstream } from "../plugins/upstream-hooks"; /** * Responses the origin may already be executing. RFC 9110 §9.2.2 forbids an intermediary @@ -472,10 +473,11 @@ export async function fetchWithAttemptDeadline( if (preferIdentityEncoding && !headers.has("accept-encoding")) { headers.set("accept-encoding", "identity"); } + const target = rewriteUpstream(url, headers, "http"); try { - return await executor(url, { + return await executor(target.url, { ...init, - headers, + headers: target.headers, redirect: "manual", signal: attemptTimeout.signal, }); diff --git a/src/plugins/loader.ts b/src/plugins/loader.ts new file mode 100644 index 00000000000..2083f8a08ee --- /dev/null +++ b/src/plugins/loader.ts @@ -0,0 +1,156 @@ +/** + * Local plugin loader. + * + * `ocx start` imports every `*.ts`, `*.js` and `*.mjs` file in `$OPENCODEX_HOME/plugins/` + * before the server binds, so a plugin's hooks are in place for the first request. A missing + * directory means no plugins and no work. `OCX_PLUGINS=0` disables loading for one start. + * + * A plugin is a module whose default export is `{ name?, setup(ctx) }`. It runs in the proxy + * process with the operator's credentials, so the loader accepts only files owned by the + * current user that no other user can write — the same trust boundary as `config.json`. + * Plugins cannot import ocx internals (a compiled binary keeps them inside `$bunfs`); they + * receive everything they may use through `OcxPluginContext`. + * + * Every failure is contained: a plugin that throws, times out or has the wrong shape is + * reported and skipped, and the remaining plugins and the proxy start normally. + */ + +import { readdirSync, statSync } from "node:fs"; +import { basename, join } from "node:path"; +import { pathToFileURL } from "node:url"; +import { getConfigDir } from "../config/paths"; +import { registerOptionalShutdownHook } from "../lib/optional-shutdown-hooks"; +import { registerUpstreamRewriter, type UpstreamRewriter } from "./upstream-hooks"; + +export type { UpstreamRewriter, UpstreamTarget, UpstreamTransport } from "./upstream-hooks"; + +export interface OcxPluginContext { + /** The plugin's own name, as reported in logs. */ + readonly name: string; + /** `$OPENCODEX_HOME`, for plugins that keep state next to the proxy's. */ + readonly configDir: string; + /** `$OPENCODEX_HOME/plugins`. */ + readonly pluginDir: string; + log(message: string): void; + /** See `src/plugins/upstream-hooks.ts`. Called synchronously on every provider send. */ + registerUpstreamRewriter(rewrite: UpstreamRewriter): void; + /** Runs once when the proxy shuts down. Must not throw or block. */ + onShutdown(teardown: () => void): void; +} + +export interface OcxPlugin { + name?: string; + setup(context: OcxPluginContext): void | Promise; +} + +export interface PluginLoadResult { + file: string; + name: string; + loaded: boolean; + error?: string; +} + +const PLUGIN_EXTENSIONS = [".ts", ".js", ".mjs"]; +const SETUP_TIMEOUT_MS = 5_000; + +export function pluginDirectory(): string { + return join(getConfigDir(), "plugins"); +} + +function listPluginFiles(dir: string): string[] { + let entries: string[]; + try { + entries = readdirSync(dir); + } catch { + return []; + } + return entries + .filter(entry => !entry.startsWith(".") && !entry.startsWith("_") && !entry.endsWith(".d.ts")) + .filter(entry => PLUGIN_EXTENSIONS.some(extension => entry.endsWith(extension))) + .sort() + .map(entry => join(dir, entry)); +} + +/** Null when the file is safe to execute, otherwise the reason it is refused. */ +export function pluginFileTrustError(file: string): string | null { + let stats: ReturnType; + try { + stats = statSync(file); + } catch (error) { + return error instanceof Error ? error.message : String(error); + } + if (!stats.isFile()) return "not a regular file"; + if (process.platform === "win32") return null; + const uid = process.getuid?.(); + if (uid !== undefined && stats.uid !== uid) return "owned by another user"; + if ((stats.mode & 0o022) !== 0) return "writable by group or others (chmod go-w)"; + return null; +} + +function isPlugin(value: unknown): value is OcxPlugin { + return typeof value === "object" && value !== null && typeof (value as OcxPlugin).setup === "function"; +} + +async function withTimeout(work: Promise, ms: number, label: string): Promise { + let timer: ReturnType | undefined; + const timeout = new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(`${label} did not finish within ${ms}ms`)), ms); + }); + try { + return await Promise.race([work, timeout]); + } finally { + clearTimeout(timer); + } +} + +export async function loadOcxPlugins(dir = pluginDirectory()): Promise { + if (process.env["OCX_PLUGINS"] === "0") return []; + const results: PluginLoadResult[] = []; + for (const file of listPluginFiles(dir)) { + const fallbackName = basename(file).replace(/\.(ts|js|mjs)$/, ""); + const refused = pluginFileTrustError(file); + if (refused) { + results.push({ file, name: fallbackName, loaded: false, error: `refused: ${refused}` }); + continue; + } + const unregister: Array<() => void> = []; + try { + const module = await import(pathToFileURL(file).href) as { default?: unknown; plugin?: unknown }; + const plugin = module.default ?? module.plugin; + if (!isPlugin(plugin)) throw new Error("default export must be { name?, setup(context) }"); + const name = typeof plugin.name === "string" && plugin.name.trim() ? plugin.name.trim() : fallbackName; + const context: OcxPluginContext = { + name, + configDir: getConfigDir(), + pluginDir: dir, + log: message => console.log(`[plugin:${name}] ${message}`), + registerUpstreamRewriter: rewrite => { unregister.push(registerUpstreamRewriter(name, rewrite)); }, + onShutdown: teardown => { unregister.push(registerOptionalShutdownHook(`plugin:${name}`, teardown)); }, + }; + await withTimeout(Promise.resolve(plugin.setup(context)), SETUP_TIMEOUT_MS, `plugin "${name}" setup`); + results.push({ file, name, loaded: true }); + } catch (error) { + // A half-initialised plugin must not leave hooks behind. + for (const undo of unregister) undo(); + results.push({ + file, + name: fallbackName, + loaded: false, + error: error instanceof Error ? error.message : String(error), + }); + } + } + return results; +} + +/** `ocx start` entry: load and print one line per plugin. Never throws. */ +export async function loadAndReportOcxPlugins(): Promise { + try { + for (const result of await loadOcxPlugins()) { + if (result.loaded) console.log(`🔌 Plugin loaded: ${result.name}`); + else console.error(`⚠️ Plugin ${result.name} skipped: ${result.error}`); + } + } catch (error) { + console.error(`⚠️ Plugin loading failed: ${error instanceof Error ? error.message : String(error)}`); + } +} diff --git a/src/plugins/upstream-hooks.ts b/src/plugins/upstream-hooks.ts new file mode 100644 index 00000000000..a2b759ff5f5 --- /dev/null +++ b/src/plugins/upstream-hooks.ts @@ -0,0 +1,88 @@ +/** + * Upstream rewrite slot for local plugins. + * + * A plugin loaded by `src/plugins/loader.ts` may register a rewriter that sees every + * provider send immediately before it leaves the process: HTTP through + * `fetchWithHeaderTimeout` / `fetchWithAttemptDeadline`, and the Codex WebSocket dial in + * `CodexWsSession`. The rewriter may replace the URL and add or change headers — enough to + * put a local sidecar (a compression proxy, a recorder) in front of the provider without + * the core knowing it exists. + * + * This module imports nothing, so the request path pays one array-length check when no + * plugin is installed. A rewriter that throws is disabled for the rest of the process and + * the send continues unmodified: a broken plugin must never take the proxy down with it. + */ + +export type UpstreamTransport = "http" | "websocket"; + +export interface UpstreamTarget { + /** Absolute upstream URL. A rewriter may assign a new one. */ + url: string; + /** Mutable outbound headers. Credentials are present; a rewriter must not log them. */ + headers: Headers; + readonly transport: UpstreamTransport; +} + +export type UpstreamRewriter = (target: UpstreamTarget) => void; + +interface Registration { + readonly name: string; + readonly rewrite: UpstreamRewriter; + disabled: boolean; +} + +const registrations: Registration[] = []; + +export function registerUpstreamRewriter(name: string, rewrite: UpstreamRewriter): () => void { + const registration: Registration = { name, rewrite, disabled: false }; + registrations.push(registration); + return () => { + const index = registrations.indexOf(registration); + if (index >= 0) registrations.splice(index, 1); + }; +} + +export function hasUpstreamRewriters(): boolean { + return registrations.length > 0; +} + +/** + * Run every active rewriter over one send. Returns the input untouched (same objects) when + * no rewriter is registered, so the common path allocates nothing. + */ +export function rewriteUpstream( + url: string, + headers: H, + transport: UpstreamTransport, +): { url: string; headers: H | Headers } { + if (registrations.length === 0) return { url, headers }; + const target: UpstreamTarget = { url, headers: new Headers(headers), transport }; + for (const registration of registrations) { + if (registration.disabled) continue; + try { + registration.rewrite(target); + } catch (error) { + registration.disabled = true; + const reason = error instanceof Error ? error.message : String(error); + console.error(`[opencodex] plugin "${registration.name}" upstream rewriter disabled after an error: ${reason}`); + } + } + return { url: target.url, headers: target.headers }; +} + +/** Plain-record variant for callers that hold headers as `Record` (WebSocket dial). */ +export function rewriteUpstreamRecord( + url: string, + headers: Record, + transport: UpstreamTransport, +): { url: string; headers: Record } { + if (registrations.length === 0) return { url, headers }; + const result = rewriteUpstream(url, headers, transport); + const record: Record = {}; + new Headers(result.headers).forEach((value, key) => { record[key] = value; }); + return { url: result.url, headers: record }; +} + +export function resetUpstreamRewritersForTests(): void { + registrations.length = 0; +} diff --git a/src/server/responses/codex-ws-session.ts b/src/server/responses/codex-ws-session.ts index 32716a52979..3a8f4a84a53 100644 --- a/src/server/responses/codex-ws-session.ts +++ b/src/server/responses/codex-ws-session.ts @@ -1,3 +1,5 @@ +import { rewriteUpstreamRecord } from "../../plugins/upstream-hooks"; + export const MAX_CODEX_WS_SESSION_EXCHANGES = 32; /** Owns one physical socket; request listeners belong to the exchange, not this object. */ @@ -11,7 +13,8 @@ export class CodexWsSession { constructor(url: string, headers: Record, readonly retainable = false, private readonly changed: () => void = () => {}, proxy?: string) { - this.socket = new WebSocket(url, { headers, ...(proxy ? { proxy } : {}) } as unknown as string[]); + const target = rewriteUpstreamRecord(url, headers, "websocket"); + this.socket = new WebSocket(target.url, { headers: target.headers, ...(proxy ? { proxy } : {}) } as unknown as string[]); this.socket.addEventListener("open", this.onOpen); this.socket.addEventListener("message", this.onIdleMessage); this.socket.addEventListener("close", this.onClose); diff --git a/src/server/responses/fetch-helpers.ts b/src/server/responses/fetch-helpers.ts index 1c27951019a..789cd597ac8 100644 --- a/src/server/responses/fetch-helpers.ts +++ b/src/server/responses/fetch-helpers.ts @@ -12,6 +12,7 @@ import { waitForProviderRequestSlot } from "../../providers/request-pacing"; import { withUpstreamHttpVersion } from "../../lib/upstream-http-version"; import type { CodexWsQuotaObserver } from "./codex-ws-metadata"; import { configuredOutboundFetch } from "../../lib/proxy-env"; +import { rewriteUpstream } from "../../plugins/upstream-hooks"; import { describeProviderEgressForLog, markEgressTransparentExecutor, @@ -382,10 +383,11 @@ export async function fetchWithHeaderTimeout( if (preferIdentityEncoding && !headers.has("accept-encoding")) { headers.set("accept-encoding", "identity"); } + const target = rewriteUpstream(url, headers, "http"); try { - return await fetchExecutor(url, { + return await fetchExecutor(target.url, { ...init, - headers, + headers: target.headers, // Never replay provider credentials or request bodies to a redirect destination. // Preserve the 3xx for the owner's existing response/health policy (#914, #1471). redirect: "manual", diff --git a/structure/INDEX.md b/structure/INDEX.md index 4f653bf3fb0..d1a258a6810 100644 --- a/structure/INDEX.md +++ b/structure/INDEX.md @@ -89,6 +89,7 @@ Background service, docs, release, and design discipline. | --- | --- | | [`desktop-shell.md`](desktop-shell.md) | Tauri desktop shell, proxy attachment and sidecar lifecycle, tray controls, bootstrap navigation, and desktop companion presence. | | [`ops/service-and-sidecars.md`](ops/service-and-sidecars.md) | Service install/repair, platform launchers, tray, and sidecar processes. | +| [`ops/plugins.md`](ops/plugins.md) | Plugin loading from OPENCODEX_HOME/plugins and the upstream rewrite slot plugins attach to. | | [`ops/docs-and-release.md`](ops/docs-and-release.md) | Docs site, workflow map, branch policy, release flow, and cross-platform CI. | | [`design-methodology.md`](design-methodology.md) | Stage ordering for new GUI, CLI, and user-facing surfaces. | | [`ops/cross-platform-ci.md`](ops/cross-platform-ci.md) | Test lanes, platform coverage, aggregate gating, and release CI proof. | @@ -134,6 +135,7 @@ A source area can be described by more than one doc, because these docs are orga | `src/lib/` | [`overview.md`](overview.md)
[`runtime.md`](runtime.md)
[`transports/byte-accounting.md`](transports/byte-accounting.md)
[`transports/responses-wire-shapes.md`](transports/responses-wire-shapes.md)
[`transports/responses-failover.md`](transports/responses-failover.md)
[`transports/responses-spend.md`](transports/responses-spend.md)
[`transports/inventory.md`](transports/inventory.md)
[`gui-and-management-api.md`](gui-and-management-api.md)
[`dashboard-and-usage.md`](dashboard-and-usage.md)
[`clients/integrations.md`](clients/integrations.md)
[`ops/docs-and-release.md`](ops/docs-and-release.md) | | `src/link/` | [`remote-link.md`](remote-link.md) | | `src/oauth/` | [`runtime.md`](runtime.md)
[`transports/inventory.md`](transports/inventory.md)
[`providers-and-adapters.md`](providers-and-adapters.md)
[`providers/xai-grok.md`](providers/xai-grok.md) | +| `src/plugins/` | [`ops/plugins.md`](ops/plugins.md) | | `src/protocols/` | [`data-planes/protocol-paths.md`](data-planes/protocol-paths.md) | | `src/providers/` | [`runtime.md`](runtime.md)
[`subagents.md`](subagents.md)
[`transports/inventory.md`](transports/inventory.md)
[`providers-and-adapters.md`](providers-and-adapters.md)
[`providers/xai-grok.md`](providers/xai-grok.md) | | `src/quota/` | [`dashboard-and-usage.md`](dashboard-and-usage.md) | diff --git a/structure/manifest.json b/structure/manifest.json index 734e7d71aa8..5201bb9bf71 100644 --- a/structure/manifest.json +++ b/structure/manifest.json @@ -456,6 +456,15 @@ "src/update/" ] }, + { + "path": "ops/plugins.md", + "tier": 6, + "title": "Local Plugins", + "scope": "Plugin loading from OPENCODEX_HOME/plugins and the upstream rewrite slot plugins attach to.", + "documents": [ + "src/plugins/" + ] + }, { "path": "ops/docs-and-release.md", "tier": 6, diff --git a/structure/ops/plugins.md b/structure/ops/plugins.md new file mode 100644 index 00000000000..2d0be5bb25c --- /dev/null +++ b/structure/ops/plugins.md @@ -0,0 +1,40 @@ +# Local Plugins + +Local plugins let an operator put code in front of provider sends without editing the core. +They are local extensions of one install, not a distribution channel: nothing fetches, updates +or signs them. + +## Loading + +- `ocx start` calls `loadAndReportOcxPlugins()` from `src/plugins/loader.ts` after the config is + loaded and before `startServer`, so every hook is registered before the listener binds. + `startServer` itself stays synchronous; plugin loading is awaited in the CLI, never inside it. +- The loader imports `*.ts`, `*.js` and `*.mjs` from `$OPENCODEX_HOME/plugins/`, sorted by name. + Names starting with `.` or `_` and `*.d.ts` are ignored. A missing directory loads nothing. +- `OCX_PLUGINS=0` disables loading for that process. +- A plugin runs with the operator's credentials, so the loader refuses a file that is not a regular + file, is owned by another user, or is writable by group or others (POSIX). This is the same trust + boundary as `config.json`. +- A plugin module default-exports `{ name?, setup(context) }`. `setup` has five seconds. A plugin + that throws, times out or has the wrong shape is reported and skipped, and every hook it + registered during the failed setup is removed. The other plugins and the proxy start normally. +- Plugins cannot import ocx modules: in a compiled binary they live inside `$bunfs`. Everything a + plugin may use arrives through `OcxPluginContext` (`name`, `configDir`, `pluginDir`, `log`, + `registerUpstreamRewriter`, `onShutdown`). `onShutdown` registers through + `src/lib/optional-shutdown-hooks.ts`. + +## Upstream rewrite slot + +`src/plugins/upstream-hooks.ts` is the only core-owned seam plugins attach to. It imports nothing, +so the request path depends on it without depending on the loader. + +- It runs synchronously at the physical send: HTTP in `fetchWithHeaderTimeout` + (`src/server/responses/fetch-helpers.ts`) and `fetchWithAttemptDeadline` + (`src/lib/upstream-retry.ts`), and the Codex WebSocket dial in `CodexWsSession` + (`src/server/responses/codex-ws-session.ts`). The target carries the URL, mutable headers and the + transport (`http` or `websocket`). +- With no rewriter registered, the send is returned untouched and nothing is allocated. +- A rewriter that throws is disabled for the rest of the process and the send continues unmodified. +- Rewrites happen after the request is built and after egress/pacing decisions, so they do not + change routing, account selection, retry budgets or logging identity. A rewriter that moves a send + to another host owns that host's behaviour; the core does not re-validate it. diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 990cf779de9..7bc71ac8cbd 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1,6 +1,8 @@ { "openai-chat-serialized-tool-call-scaling.test.ts": "adapters/openai", "coding-agent-json-lines-scaling.test.ts": "providers", + "plugin-loader.test.ts": "lib", + "plugin-upstream-hooks.test.ts": "lib", "usage-snapshot-digest-reuse.test.ts": "usage", "release-desktop-scripts.test.ts": "ci-workflows", "installed-gate-drivers.test.ts": "ci-workflows", diff --git a/tests/lib/plugin-loader.test.ts b/tests/lib/plugin-loader.test.ts new file mode 100644 index 00000000000..311237846a2 --- /dev/null +++ b/tests/lib/plugin-loader.test.ts @@ -0,0 +1,96 @@ +import { afterEach, beforeEach, expect, test } from "bun:test"; +import { chmodSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { loadOcxPlugins, pluginFileTrustError } from "../../src/plugins/loader"; +import { + hasUpstreamRewriters, + resetUpstreamRewritersForTests, + rewriteUpstream, +} from "../../src/plugins/upstream-hooks"; + +let dir: string; + +beforeEach(() => { + dir = mkdtempSync(join(tmpdir(), "ocx-plugins-")); + delete process.env["OCX_PLUGINS"]; +}); + +afterEach(() => { + resetUpstreamRewritersForTests(); + delete process.env["OCX_PLUGINS"]; + rmSync(dir, { recursive: true, force: true }); +}); + +function writePlugin(file: string, source: string, mode = 0o600): string { + const path = join(dir, file); + writeFileSync(path, source); + chmodSync(path, mode); + return path; +} + +const REDIRECT_PLUGIN = ` +export default { + name: "redirect", + setup(ctx) { + ctx.registerUpstreamRewriter(target => { target.url = "http://127.0.0.1:8787" + new URL(target.url).pathname; }); + }, +}; +`; + +test("a missing plugin directory loads nothing", async () => { + expect(await loadOcxPlugins(join(dir, "absent"))).toEqual([]); + expect(hasUpstreamRewriters()).toBe(false); +}); + +test("a valid plugin registers its upstream rewriter", async () => { + writePlugin("redirect.ts", REDIRECT_PLUGIN); + const results = await loadOcxPlugins(dir); + expect(results.map(result => [result.name, result.loaded])).toEqual([["redirect", true]]); + expect(rewriteUpstream("https://api.example.com/v1/responses", undefined, "http").url) + .toBe("http://127.0.0.1:8787/v1/responses"); +}); + +test("OCX_PLUGINS=0 skips loading", async () => { + writePlugin("redirect.ts", REDIRECT_PLUGIN); + process.env["OCX_PLUGINS"] = "0"; + expect(await loadOcxPlugins(dir)).toEqual([]); + expect(hasUpstreamRewriters()).toBe(false); +}); + +test.skipIf(process.platform === "win32")("a group- or world-writable plugin is refused", async () => { + const path = writePlugin("redirect.ts", REDIRECT_PLUGIN, 0o664); + expect(pluginFileTrustError(path)).toContain("writable by group or others"); + const [result] = await loadOcxPlugins(dir); + expect(result?.loaded).toBe(false); + expect(hasUpstreamRewriters()).toBe(false); +}); + +test("a wrong export shape or a throwing setup is skipped and leaves no hooks behind", async () => { + writePlugin("a-shape.ts", "export default { name: 'shape' };"); + writePlugin("b-throws.ts", ` +export default { + setup(ctx) { + ctx.registerUpstreamRewriter(target => { target.url = "http://leaked/"; }); + throw new Error("setup failed"); + }, +}; +`); + writePlugin("c-ok.ts", REDIRECT_PLUGIN); + const results = await loadOcxPlugins(dir); + expect(results.map(result => [result.name, result.loaded])).toEqual([ + ["a-shape", false], + ["b-throws", false], + ["redirect", true], + ]); + expect(results[1]?.error).toBe("setup failed"); + expect(rewriteUpstream("https://api.example.com/v1/x", undefined, "http").url).toBe("http://127.0.0.1:8787/v1/x"); +}); + +test("hidden, underscore-prefixed and declaration files are ignored", async () => { + writePlugin(".hidden.ts", REDIRECT_PLUGIN); + writePlugin("_draft.ts", REDIRECT_PLUGIN); + writePlugin("types.d.ts", "export {};"); + writePlugin("notes.md", "# not a plugin"); + expect(await loadOcxPlugins(dir)).toEqual([]); +}); diff --git a/tests/lib/plugin-upstream-hooks.test.ts b/tests/lib/plugin-upstream-hooks.test.ts new file mode 100644 index 00000000000..8806ff32a65 --- /dev/null +++ b/tests/lib/plugin-upstream-hooks.test.ts @@ -0,0 +1,64 @@ +import { afterEach, expect, test } from "bun:test"; +import { + hasUpstreamRewriters, + registerUpstreamRewriter, + resetUpstreamRewritersForTests, + rewriteUpstream, + rewriteUpstreamRecord, +} from "../../src/plugins/upstream-hooks"; + +afterEach(() => resetUpstreamRewritersForTests()); + +test("with no plugin registered the send is returned untouched and unallocated", () => { + const headers = new Headers({ authorization: "Bearer x" }); + const result = rewriteUpstream("https://api.example.com/v1/responses", headers, "http"); + expect(hasUpstreamRewriters()).toBe(false); + expect(result.url).toBe("https://api.example.com/v1/responses"); + expect(result.headers).toBe(headers); +}); + +test("a rewriter can redirect the URL and add headers while keeping credentials", () => { + registerUpstreamRewriter("sidecar", target => { + const original = new URL(target.url); + target.url = `http://127.0.0.1:8787${original.pathname}`; + target.headers.set("x-sidecar-upstream", original.origin); + }); + const result = rewriteUpstream("https://api.example.com/v1/responses", { authorization: "Bearer x" }, "http"); + const headers = new Headers(result.headers); + expect(result.url).toBe("http://127.0.0.1:8787/v1/responses"); + expect(headers.get("x-sidecar-upstream")).toBe("https://api.example.com"); + expect(headers.get("authorization")).toBe("Bearer x"); +}); + +test("rewriters see the transport and run in registration order", () => { + const seen: string[] = []; + registerUpstreamRewriter("first", target => { seen.push(`first:${target.transport}`); target.url += "?a"; }); + registerUpstreamRewriter("second", target => { seen.push(`second:${target.transport}`); target.url += "&b"; }); + const result = rewriteUpstreamRecord("wss://chatgpt.com/backend-api/codex/responses", { "x-k": "v" }, "websocket"); + expect(seen).toEqual(["first:websocket", "second:websocket"]); + expect(result.url).toBe("wss://chatgpt.com/backend-api/codex/responses?a&b"); + expect(result.headers["x-k"]).toBe("v"); +}); + +test("a throwing rewriter is disabled and never breaks the send", () => { + let calls = 0; + registerUpstreamRewriter("broken", () => { calls += 1; throw new Error("boom"); }); + const originalError = console.error; + console.error = () => {}; + try { + for (let i = 0; i < 3; i += 1) { + expect(rewriteUpstream("https://api.example.com/v1/messages", undefined, "http").url) + .toBe("https://api.example.com/v1/messages"); + } + } finally { + console.error = originalError; + } + expect(calls).toBe(1); +}); + +test("unregistering removes the rewriter", () => { + const off = registerUpstreamRewriter("temp", target => { target.url = "http://changed/"; }); + off(); + expect(hasUpstreamRewriters()).toBe(false); + expect(rewriteUpstream("https://a.example/x", undefined, "http").url).toBe("https://a.example/x"); +}); diff --git a/tests/responses/responses-fetch-helpers-boundary.test.ts b/tests/responses/responses-fetch-helpers-boundary.test.ts index e2f63efeacd..455c3d0d671 100644 --- a/tests/responses/responses-fetch-helpers-boundary.test.ts +++ b/tests/responses/responses-fetch-helpers-boundary.test.ts @@ -49,6 +49,8 @@ describe("Responses fetch-helper import boundary", () => { "../../lib/proxy-env", "../../lib/redact", "../../lib/upstream-http-version", + // Import-free plugin rewrite slot (src/plugins/upstream-hooks.ts). + "../../plugins/upstream-hooks", "../../providers/request-pacing", "./ws-upstream", ]); From ca210f95888ea15801dda3f6cf2b4f0b7c5888ef Mon Sep 17 00:00:00 2001 From: halysondev Date: Fri, 25 Sep 2026 23:37:48 -0300 Subject: [PATCH 2/6] fix(plugins): rewrite at the physical send and harden plugin setup Addresses the CodeRabbit review on #5896. - Move the HTTP rewrite from fetchWithHeaderTimeout/fetchWithAttemptDeadline into sendWithConnectionPolicy. Rewriting before the transport choice hid the chatgpt.com origin from the Codex WebSocket selection and pushed Codex turns onto HTTP; egress now also follows the rewritten destination. Nested passes rewrite once via an init mark, like EGRESS_DECIDED. - WebSocket dials redirected to loopback drop the caller's proxy, which was chosen for the original destination and cannot reach local loopback. - Undo a rewriter's edits to the target when it throws. - Report plugin directory read failures other than ENOENT. - Key shutdown hooks by plugin file, not display name. - Close a plugin's context when setup fails or times out, so a setup that resumes later cannot register; state that the deadline only bounds setup that yields. - Guide example keeps the upstream query string. Co-Authored-By: Claude Opus 5.5 --- .../src/content/docs/guides/local-plugins.md | 21 ++++--- src/lib/upstream-retry.ts | 6 +- src/plugins/loader.ts | 50 ++++++++++++--- src/plugins/upstream-hooks.ts | 41 +++++++++++-- src/server/responses/codex-ws-session.ts | 6 +- src/server/responses/fetch-helpers.ts | 21 +++++-- structure/ops/plugins.md | 35 +++++++---- tests/lib/plugin-loader.test.ts | 56 +++++++++++++++++ tests/lib/plugin-upstream-hooks.test.ts | 61 +++++++++++++++++++ 9 files changed, 251 insertions(+), 46 deletions(-) diff --git a/docs-site/src/content/docs/guides/local-plugins.md b/docs-site/src/content/docs/guides/local-plugins.md index 64969434dbf..de90a9b8078 100644 --- a/docs-site/src/content/docs/guides/local-plugins.md +++ b/docs-site/src/content/docs/guides/local-plugins.md @@ -36,7 +36,8 @@ To start once without plugins, set `OCX_PLUGINS=0`. ## Writing a plugin A plugin default-exports an object with an optional `name` and a `setup` function. `setup` receives a -context and has five seconds to finish. +context. An asynchronous `setup` has five seconds to finish; plugins run in the proxy's own thread, so +a `setup` that blocks synchronously cannot be interrupted and delays startup until it returns. ```ts interface UpstreamTarget { @@ -55,7 +56,7 @@ export default { ctx.registerUpstreamRewriter(target => { const upstream = new URL(target.url); if (!upstream.pathname.endsWith("/chat/completions")) return; - target.url = `http://127.0.0.1:9000${upstream.pathname}`; + target.url = `http://127.0.0.1:9000${upstream.pathname}${upstream.search}`; target.headers.set("x-original-origin", upstream.origin); }); }, @@ -70,12 +71,18 @@ small interfaces you need locally, as above. ## How rewrites behave - The rewriter runs synchronously on every provider send over HTTP and on the Codex WebSocket - connection. Keep it fast; do network checks (health probes) in the background and read a cached - result in the rewriter. + connection, after opencodex has picked the transport. Keep it fast; do network checks (health + probes) in the background and read a cached result in the rewriter. +- Egress settings (proxy, `noProxy`) are applied to the rewritten HTTP destination. A Codex + WebSocket redirected to a loopback address connects directly instead of through a configured + proxy, which could not reach this machine's loopback. - It runs after opencodex has chosen the provider, account and route, so it does not change routing, account selection, retries or request logs. - A redirected send goes to the host you chose. That host sees the request exactly as the provider would, credentials included. -- If a rewriter throws, opencodex disables it for the rest of the process and sends the request - unmodified. If `setup` throws or times out, the plugin is skipped and anything it registered is - removed; other plugins and the proxy start normally. +- If a rewriter throws, opencodex undoes its changes to that send, disables it for the rest of the + process and sends the request unmodified. If `setup` throws or times out, the plugin is skipped, + anything it registered is removed, and later registration attempts from it are ignored; other + plugins and the proxy start normally. +- A plugin directory that exists but cannot be read (for example, wrong permissions) is reported at + startup rather than treated as empty. diff --git a/src/lib/upstream-retry.ts b/src/lib/upstream-retry.ts index 7594b0b8458..2b2b546388c 100644 --- a/src/lib/upstream-retry.ts +++ b/src/lib/upstream-retry.ts @@ -17,7 +17,6 @@ */ import { clearableDeadline } from "./abort"; import { redactSecretString } from "./redact"; -import { rewriteUpstream } from "../plugins/upstream-hooks"; /** * Responses the origin may already be executing. RFC 9110 §9.2.2 forbids an intermediary @@ -473,11 +472,10 @@ export async function fetchWithAttemptDeadline( if (preferIdentityEncoding && !headers.has("accept-encoding")) { headers.set("accept-encoding", "identity"); } - const target = rewriteUpstream(url, headers, "http"); try { - return await executor(target.url, { + return await executor(url, { ...init, - headers: target.headers, + headers, redirect: "manual", signal: attemptTimeout.signal, }); diff --git a/src/plugins/loader.ts b/src/plugins/loader.ts index 2083f8a08ee..8e650311381 100644 --- a/src/plugins/loader.ts +++ b/src/plugins/loader.ts @@ -11,8 +11,10 @@ * Plugins cannot import ocx internals (a compiled binary keeps them inside `$bunfs`); they * receive everything they may use through `OcxPluginContext`. * - * Every failure is contained: a plugin that throws, times out or has the wrong shape is - * reported and skipped, and the remaining plugins and the proxy start normally. + * Failures are contained: a plugin that throws, times out or has the wrong shape is reported + * and skipped, its context stops accepting registrations, and the remaining plugins and the + * proxy start normally. The setup deadline bounds setup that yields to the event loop; plugins + * run in the proxy's own thread, so synchronous work that never yields cannot be interrupted. */ import { readdirSync, statSync } from "node:fs"; @@ -57,12 +59,14 @@ export function pluginDirectory(): string { return join(getConfigDir(), "plugins"); } +/** A missing directory is "no plugins"; any other read failure propagates to be reported. */ function listPluginFiles(dir: string): string[] { let entries: string[]; try { entries = readdirSync(dir); - } catch { - return []; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return []; + throw error; } return entries .filter(entry => !entry.startsWith(".") && !entry.startsWith("_") && !entry.endsWith(".d.ts")) @@ -103,10 +107,24 @@ async function withTimeout(work: Promise, ms: number, label: string): Prom } } -export async function loadOcxPlugins(dir = pluginDirectory()): Promise { +export interface LoadOcxPluginsOptions { + /** Deadline for a setup that yields; see the module comment. */ + setupTimeoutMs?: number; +} + +export async function loadOcxPlugins( + dir = pluginDirectory(), + options: LoadOcxPluginsOptions = {}, +): Promise { if (process.env["OCX_PLUGINS"] === "0") return []; + let files: string[]; + try { + files = listPluginFiles(dir); + } catch (error) { + return [{ file: dir, name: "plugins directory", loaded: false, error: error instanceof Error ? error.message : String(error) }]; + } const results: PluginLoadResult[] = []; - for (const file of listPluginFiles(dir)) { + for (const file of files) { const fallbackName = basename(file).replace(/\.(ts|js|mjs)$/, ""); const refused = pluginFileTrustError(file); if (refused) { @@ -114,6 +132,15 @@ export async function loadOcxPlugins(dir = pluginDirectory()): Promise void> = []; + // Closed when setup fails or times out: a setup that resumes later must not register. + let active = true; + const whileActive = (name: string, register: () => () => void): void => { + if (!active) { + console.error(`[plugin:${name}] registration after a failed setup was ignored`); + return; + } + unregister.push(register()); + }; try { const module = await import(pathToFileURL(file).href) as { default?: unknown; plugin?: unknown }; const plugin = module.default ?? module.plugin; @@ -124,13 +151,16 @@ export async function loadOcxPlugins(dir = pluginDirectory()): Promise console.log(`[plugin:${name}] ${message}`), - registerUpstreamRewriter: rewrite => { unregister.push(registerUpstreamRewriter(name, rewrite)); }, - onShutdown: teardown => { unregister.push(registerOptionalShutdownHook(`plugin:${name}`, teardown)); }, + registerUpstreamRewriter: rewrite => whileActive(name, () => registerUpstreamRewriter(name, rewrite)), + // Keyed by file, not name: two plugins may share a display name. + onShutdown: teardown => whileActive(name, () => registerOptionalShutdownHook(`plugin:${file}`, teardown)), }; - await withTimeout(Promise.resolve(plugin.setup(context)), SETUP_TIMEOUT_MS, `plugin "${name}" setup`); + const timeoutMs = options.setupTimeoutMs ?? SETUP_TIMEOUT_MS; + await withTimeout(Promise.resolve(plugin.setup(context)), timeoutMs, `plugin "${name}" setup`); results.push({ file, name, loaded: true }); } catch (error) { - // A half-initialised plugin must not leave hooks behind. + // A half-initialised plugin must not leave hooks behind, now or later. + active = false; for (const undo of unregister) undo(); results.push({ file, diff --git a/src/plugins/upstream-hooks.ts b/src/plugins/upstream-hooks.ts index a2b759ff5f5..a12b068e990 100644 --- a/src/plugins/upstream-hooks.ts +++ b/src/plugins/upstream-hooks.ts @@ -2,11 +2,10 @@ * Upstream rewrite slot for local plugins. * * A plugin loaded by `src/plugins/loader.ts` may register a rewriter that sees every - * provider send immediately before it leaves the process: HTTP through - * `fetchWithHeaderTimeout` / `fetchWithAttemptDeadline`, and the Codex WebSocket dial in - * `CodexWsSession`. The rewriter may replace the URL and add or change headers — enough to - * put a local sidecar (a compression proxy, a recorder) in front of the provider without - * the core knowing it exists. + * provider send at the physical boundary, after the transport was chosen: HTTP in + * `sendWithConnectionPolicy` and the Codex WebSocket dial in `CodexWsSession`. The rewriter + * may replace the URL and add or change headers — enough to put a local sidecar (a + * compression proxy, a recorder) in front of the provider without the core knowing it exists. * * This module imports nothing, so the request path pays one array-length check when no * plugin is installed. A rewriter that throws is disabled for the rest of the process and @@ -59,9 +58,15 @@ export function rewriteUpstream( const target: UpstreamTarget = { url, headers: new Headers(headers), transport }; for (const registration of registrations) { if (registration.disabled) continue; + // A rewriter that edits the target and then throws must not leave a half-rewritten send + // for the next rewriter or the network. + const urlBefore = target.url; + const headersBefore = new Headers(target.headers); try { registration.rewrite(target); } catch (error) { + target.url = urlBefore; + target.headers = headersBefore; registration.disabled = true; const reason = error instanceof Error ? error.message : String(error); console.error(`[opencodex] plugin "${registration.name}" upstream rewriter disabled after an error: ${reason}`); @@ -83,6 +88,32 @@ export function rewriteUpstreamRecord( return { url: result.url, headers: record }; } +function isLoopbackUrl(raw: string): boolean { + let host: string; + try { + host = new URL(raw).hostname.toLowerCase().replace(/^\[|\]$/g, ""); + } catch { + return false; + } + return host === "localhost" || host.endsWith(".localhost") || host === "::1" || /^127\.\d{1,3}\.\d{1,3}\.\d{1,3}$/.test(host); +} + +/** + * WebSocket dial variant. The caller chose `proxy` for the original destination; a proxy + * elsewhere on the network cannot reach this machine's loopback, so a rewrite onto a loopback + * sidecar dials directly. Any other rewrite keeps the caller's proxy decision. + */ +export function rewriteWebSocketDial( + url: string, + headers: Record, + proxy: string | undefined, +): { url: string; headers: Record; proxy: string | undefined } { + if (registrations.length === 0) return { url, headers, proxy }; + const target = rewriteUpstreamRecord(url, headers, "websocket"); + const redirectedToLoopback = target.url !== url && isLoopbackUrl(target.url); + return { ...target, proxy: redirectedToLoopback ? undefined : proxy }; +} + export function resetUpstreamRewritersForTests(): void { registrations.length = 0; } diff --git a/src/server/responses/codex-ws-session.ts b/src/server/responses/codex-ws-session.ts index 3a8f4a84a53..097a28802bf 100644 --- a/src/server/responses/codex-ws-session.ts +++ b/src/server/responses/codex-ws-session.ts @@ -1,4 +1,4 @@ -import { rewriteUpstreamRecord } from "../../plugins/upstream-hooks"; +import { rewriteWebSocketDial } from "../../plugins/upstream-hooks"; export const MAX_CODEX_WS_SESSION_EXCHANGES = 32; @@ -13,8 +13,8 @@ export class CodexWsSession { constructor(url: string, headers: Record, readonly retainable = false, private readonly changed: () => void = () => {}, proxy?: string) { - const target = rewriteUpstreamRecord(url, headers, "websocket"); - this.socket = new WebSocket(target.url, { headers: target.headers, ...(proxy ? { proxy } : {}) } as unknown as string[]); + const dial = rewriteWebSocketDial(url, headers, proxy); + this.socket = new WebSocket(dial.url, { headers: dial.headers, ...(dial.proxy ? { proxy: dial.proxy } : {}) } as unknown as string[]); this.socket.addEventListener("open", this.onOpen); this.socket.addEventListener("message", this.onIdleMessage); this.socket.addEventListener("close", this.onClose); diff --git a/src/server/responses/fetch-helpers.ts b/src/server/responses/fetch-helpers.ts index 789cd597ac8..1f69baf45d5 100644 --- a/src/server/responses/fetch-helpers.ts +++ b/src/server/responses/fetch-helpers.ts @@ -35,6 +35,7 @@ const EGRESS_DOWNGRADE_NOTICE_LIMIT = 64; * `dispatchOverride` performs, and an unknown symbol on a `RequestInit` is inert at the wire. */ const EGRESS_DECIDED = Symbol.for("opencodex.provider-egress.decided"); +const UPSTREAM_REWRITTEN = Symbol.for("opencodex.plugins.upstream-rewritten"); /** * Announce once, per provider, that an explicit egress route moved this provider off the @@ -144,11 +145,21 @@ export type ProviderFetch = typeof globalThis.fetch & PaceAwareFetch; */ export function sendWithConnectionPolicy( physicalFetch: typeof globalThis.fetch, - input: Parameters[0], + rawInput: Parameters[0], init?: RequestInit, egress?: ProviderEgressBinding, ): Promise { - const headers = new Headers(init?.headers ?? (input instanceof Request ? input.headers : undefined)); + let input = rawInput; + let headers = new Headers(init?.headers ?? (input instanceof Request ? input.headers : undefined)); + // Plugin rewrites (src/plugins/upstream-hooks.ts) run here, after the caller chose between + // the Codex WebSocket and HTTP, and before the connection and egress decisions below so + // those follow the rewritten destination. Nested passes rewrite once, like the egress mark. + const rewriteDone = (init as Record | undefined)?.[UPSTREAM_REWRITTEN] === true; + if (!rewriteDone && (typeof input === "string" || input instanceof URL)) { + const target = rewriteUpstream(String(input), headers, "http"); + input = target.url; + headers = target.headers as Headers; + } const fresh = wantsFreshConnection(input); if (fresh) { headers.set("Connection", "close"); @@ -171,6 +182,7 @@ export function sendWithConnectionPolicy( ...(fresh ? { keepalive: false } : {}), ...egressInit, ...(decide ? { [EGRESS_DECIDED]: true } : {}), + ...{ [UPSTREAM_REWRITTEN]: true }, }); } @@ -383,11 +395,10 @@ export async function fetchWithHeaderTimeout( if (preferIdentityEncoding && !headers.has("accept-encoding")) { headers.set("accept-encoding", "identity"); } - const target = rewriteUpstream(url, headers, "http"); try { - return await fetchExecutor(target.url, { + return await fetchExecutor(url, { ...init, - headers: target.headers, + headers, // Never replay provider credentials or request bodies to a redirect destination. // Preserve the 3xx for the owner's existing response/health policy (#914, #1471). redirect: "manual", diff --git a/structure/ops/plugins.md b/structure/ops/plugins.md index 2d0be5bb25c..f4654bb56c8 100644 --- a/structure/ops/plugins.md +++ b/structure/ops/plugins.md @@ -15,26 +15,37 @@ or signs them. - A plugin runs with the operator's credentials, so the loader refuses a file that is not a regular file, is owned by another user, or is writable by group or others (POSIX). This is the same trust boundary as `config.json`. -- A plugin module default-exports `{ name?, setup(context) }`. `setup` has five seconds. A plugin - that throws, times out or has the wrong shape is reported and skipped, and every hook it - registered during the failed setup is removed. The other plugins and the proxy start normally. +- A missing plugin directory means no plugins. Any other read failure (`EACCES`, `ENOTDIR`) is + reported as a skipped `plugins directory` entry. +- A plugin module default-exports `{ name?, setup(context) }`. An asynchronous `setup` has five + seconds; plugins share the proxy thread, so a setup that blocks synchronously cannot be + interrupted. A plugin that throws, times out or has the wrong shape is reported and skipped: + its context is closed, every hook it registered is removed, and a setup that resumes after the + deadline cannot register again. The other plugins and the proxy start normally. - Plugins cannot import ocx modules: in a compiled binary they live inside `$bunfs`. Everything a plugin may use arrives through `OcxPluginContext` (`name`, `configDir`, `pluginDir`, `log`, `registerUpstreamRewriter`, `onShutdown`). `onShutdown` registers through - `src/lib/optional-shutdown-hooks.ts`. + `src/lib/optional-shutdown-hooks.ts` under a per-file key, so plugins sharing a display name keep + separate teardowns. ## Upstream rewrite slot `src/plugins/upstream-hooks.ts` is the only core-owned seam plugins attach to. It imports nothing, so the request path depends on it without depending on the loader. -- It runs synchronously at the physical send: HTTP in `fetchWithHeaderTimeout` - (`src/server/responses/fetch-helpers.ts`) and `fetchWithAttemptDeadline` - (`src/lib/upstream-retry.ts`), and the Codex WebSocket dial in `CodexWsSession` - (`src/server/responses/codex-ws-session.ts`). The target carries the URL, mutable headers and the - transport (`http` or `websocket`). +- It runs synchronously at the physical send, after the transport was chosen: HTTP in + `sendWithConnectionPolicy` (`src/server/responses/fetch-helpers.ts`), and the Codex WebSocket + dial in `CodexWsSession` (`src/server/responses/codex-ws-session.ts`). Rewriting any earlier would + hide the ChatGPT origin from the WebSocket selection and push Codex turns onto HTTP. The target + carries the URL, mutable headers and the transport (`http` or `websocket`). +- `sendWithConnectionPolicy` can run twice for one send (an override handing back to the supplied + executor). The outer pass rewrites and marks the init; the inner pass does not rewrite again. +- HTTP connection and egress decisions in `sendWithConnectionPolicy` follow the rewritten + destination. The WebSocket proxy is chosen by the caller for the original destination, so + `rewriteWebSocketDial` drops it when the rewrite targets loopback. - With no rewriter registered, the send is returned untouched and nothing is allocated. -- A rewriter that throws is disabled for the rest of the process and the send continues unmodified. -- Rewrites happen after the request is built and after egress/pacing decisions, so they do not - change routing, account selection, retry budgets or logging identity. A rewriter that moves a send +- A rewriter that throws has its edits to that send undone, is disabled for the rest of the + process, and the send continues unmodified. +- Rewrites happen after the request is built, routed and paced, so they do not change routing, + account selection, retry budgets or logging identity. A rewriter that moves a send to another host owns that host's behaviour; the core does not re-validate it. diff --git a/tests/lib/plugin-loader.test.ts b/tests/lib/plugin-loader.test.ts index 311237846a2..c7b4278d5b3 100644 --- a/tests/lib/plugin-loader.test.ts +++ b/tests/lib/plugin-loader.test.ts @@ -2,6 +2,7 @@ import { afterEach, beforeEach, expect, test } from "bun:test"; import { chmodSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; +import { resetOptionalShutdownHooksForTests, runOptionalShutdownHooks } from "../../src/lib/optional-shutdown-hooks"; import { loadOcxPlugins, pluginFileTrustError } from "../../src/plugins/loader"; import { hasUpstreamRewriters, @@ -87,6 +88,61 @@ export default { expect(rewriteUpstream("https://api.example.com/v1/x", undefined, "http").url).toBe("http://127.0.0.1:8787/v1/x"); }); +test("a plugin path that cannot be read is reported, not treated as empty", async () => { + const notADirectory = writePlugin("file-not-dir", "x"); + const results = await loadOcxPlugins(notADirectory); + expect(results).toHaveLength(1); + expect(results[0]?.loaded).toBe(false); + expect(results[0]?.name).toBe("plugins directory"); + expect(results[0]?.error).toContain("ENOTDIR"); +}); + +test("two plugins with the same name keep separate shutdown teardowns", async () => { + const ran: string[] = []; + (globalThis as Record)["__ocxTeardownLog"] = ran; + const source = (tag: string) => ` +export default { + name: "same", + setup(ctx) { ctx.onShutdown(() => { globalThis.__ocxTeardownLog.push("${tag}"); }); }, +}; +`; + writePlugin("a.ts", source("a")); + writePlugin("b.ts", source("b")); + resetOptionalShutdownHooksForTests(); + try { + const results = await loadOcxPlugins(dir); + expect(results.map(result => result.loaded)).toEqual([true, true]); + runOptionalShutdownHooks(); + expect(ran.sort()).toEqual(["a", "b"]); + } finally { + resetOptionalShutdownHooksForTests(); + delete (globalThis as Record)["__ocxTeardownLog"]; + } +}); + +test("a setup that resumes after its deadline cannot leave registrations behind", async () => { + writePlugin("slow.ts", ` +export default { + name: "slow", + async setup(ctx) { + await new Promise(resolve => setTimeout(resolve, 60)); + ctx.registerUpstreamRewriter(target => { target.url = "http://leaked/"; }); + }, +}; +`); + const originalError = console.error; + console.error = () => {}; + try { + const results = await loadOcxPlugins(dir, { setupTimeoutMs: 20 }); + expect(results[0]?.loaded).toBe(false); + expect(results[0]?.error).toContain("did not finish within 20ms"); + await new Promise(resolve => setTimeout(resolve, 120)); + } finally { + console.error = originalError; + } + expect(hasUpstreamRewriters()).toBe(false); +}); + test("hidden, underscore-prefixed and declaration files are ignored", async () => { writePlugin(".hidden.ts", REDIRECT_PLUGIN); writePlugin("_draft.ts", REDIRECT_PLUGIN); diff --git a/tests/lib/plugin-upstream-hooks.test.ts b/tests/lib/plugin-upstream-hooks.test.ts index 8806ff32a65..f3f78579100 100644 --- a/tests/lib/plugin-upstream-hooks.test.ts +++ b/tests/lib/plugin-upstream-hooks.test.ts @@ -5,7 +5,9 @@ import { resetUpstreamRewritersForTests, rewriteUpstream, rewriteUpstreamRecord, + rewriteWebSocketDial, } from "../../src/plugins/upstream-hooks"; +import { sendWithConnectionPolicy } from "../../src/server/responses/fetch-helpers"; afterEach(() => resetUpstreamRewritersForTests()); @@ -56,6 +58,65 @@ test("a throwing rewriter is disabled and never breaks the send", () => { expect(calls).toBe(1); }); +test("a rewriter that edits the target and then throws leaves the send unmodified", () => { + registerUpstreamRewriter("half", target => { + target.url = "http://127.0.0.1:9/partial"; + target.headers.set("x-partial", "1"); + target.headers.delete("authorization"); + throw new Error("boom"); + }); + let seenByNext: { url: string; partial: string | null; auth: string | null } | undefined; + registerUpstreamRewriter("next", target => { + seenByNext = { url: target.url, partial: target.headers.get("x-partial"), auth: target.headers.get("authorization") }; + }); + const originalError = console.error; + console.error = () => {}; + let result: ReturnType; + try { + result = rewriteUpstream("https://api.example.com/v1/responses", { authorization: "Bearer x" }, "http"); + } finally { + console.error = originalError; + } + const headers = new Headers(result.headers); + expect(seenByNext).toEqual({ url: "https://api.example.com/v1/responses", partial: null, auth: "Bearer x" }); + expect(result.url).toBe("https://api.example.com/v1/responses"); + expect(headers.get("x-partial")).toBeNull(); + expect(headers.get("authorization")).toBe("Bearer x"); +}); + +test("a WebSocket dial redirected to loopback drops the caller's proxy; other dials keep it", () => { + const proxy = "http://corp-proxy.example:3128"; + expect(rewriteWebSocketDial("wss://chatgpt.com/backend-api/codex/responses", {}, proxy).proxy).toBe(proxy); + + const off = registerUpstreamRewriter("loopback", target => { target.url = "ws://127.0.0.1:8787/backend-api/codex/responses"; }); + const local = rewriteWebSocketDial("wss://chatgpt.com/backend-api/codex/responses", { a: "1" }, proxy); + expect(local).toEqual({ url: "ws://127.0.0.1:8787/backend-api/codex/responses", headers: { a: "1" }, proxy: undefined }); + off(); + + registerUpstreamRewriter("remote", target => { target.url = "wss://relay.example.com/backend-api/codex/responses"; }); + expect(rewriteWebSocketDial("wss://chatgpt.com/backend-api/codex/responses", {}, proxy).proxy).toBe(proxy); +}); + +test("the physical HTTP send rewrites once, even through a nested override pass", async () => { + let calls = 0; + registerUpstreamRewriter("count", target => { + calls += 1; + target.url = target.url.replace("https://api.example.com", "http://127.0.0.1:8787"); + target.headers.set("x-hop", String(calls)); + }); + const seen: Array<{ url: string; hop: string | null }> = []; + const physical = (async (input: Parameters[0], init?: RequestInit) => { + seen.push({ url: String(input), hop: new Headers(init?.headers).get("x-hop") }); + return new Response("ok"); + }) as typeof fetch; + // An override that hands the send back to the supplied executor passes through twice. + const inner = ((input: Parameters[0], init?: RequestInit) => + sendWithConnectionPolicy(physical, input, init)) as typeof fetch; + await sendWithConnectionPolicy(inner, "https://api.example.com/v1/responses", { method: "POST" }); + expect(calls).toBe(1); + expect(seen).toEqual([{ url: "http://127.0.0.1:8787/v1/responses", hop: "1" }]); +}); + test("unregistering removes the rewriter", () => { const off = registerUpstreamRewriter("temp", target => { target.url = "http://changed/"; }); off(); From 0dafeba570830a08ab421d80ea9a0ac228371903 Mon Sep 17 00:00:00 2001 From: halysondev Date: Fri, 25 Sep 2026 23:59:39 -0300 Subject: [PATCH 3/6] fix(plugins): direct loopback HTTP, per-exchange WS rewrite, stricter trust Addresses the maintainer review and the second CodeRabbit pass on #5896. - HTTP sends rewritten onto loopback now dial directly (`proxy: false`, egress marked decided), matching the Codex WebSocket; provider proxies and HTTP_PROXY cannot reach a local sidecar. - The Codex WebSocket rewrite moved from the CodexWsSession constructor to codexWsUpstreamFetch, before the pool lookup, and the dialled destination joins codexWsReuseIdentity. Every exchange, including one on a pooled socket, now passes the plugin, and a socket is never reused for another destination. CodexWsSession is back to its dev shape. - Request inputs are rewritten too; the rewritten mark is only set after a rewrite pass actually ran. - The loader uses lstat: symbolic links are refused, and the plugins directory itself must be owned by the user and not group/other-writable. - Each onShutdown call gets its own key (plugin:#), so a plugin can register several teardowns. - Docs: per-rewriter rollback, Windows skips owner/mode checks, a timed-out setup keeps running, loopback sends bypass proxies, WS rewrite per turn. Co-Authored-By: Claude Opus 5.5 --- .../src/content/docs/guides/local-plugins.md | 26 ++++++--- src/plugins/loader.ts | 40 +++++++++++--- src/plugins/upstream-hooks.ts | 9 ++-- src/server/responses/codex-ws-pool.ts | 5 +- src/server/responses/codex-ws-session.ts | 5 +- src/server/responses/fetch-helpers.ts | 23 +++++--- src/server/responses/ws-upstream.ts | 10 ++-- structure/ops/plugins.md | 36 ++++++++----- tests/lib/plugin-loader.test.ts | 53 ++++++++++++++++++- tests/lib/plugin-upstream-hooks.test.ts | 40 ++++++++++++++ 10 files changed, 196 insertions(+), 51 deletions(-) diff --git a/docs-site/src/content/docs/guides/local-plugins.md b/docs-site/src/content/docs/guides/local-plugins.md index de90a9b8078..2d8bf955b17 100644 --- a/docs-site/src/content/docs/guides/local-plugins.md +++ b/docs-site/src/content/docs/guides/local-plugins.md @@ -24,8 +24,11 @@ Put plugin files in `plugins/` inside the opencodex home (`~/.opencodex/plugins/ - Names starting with `.` or `_`, and `*.d.ts`, are ignored — rename a plugin to `_my-sidecar.ts` to switch it off. - The directory is optional. Without it nothing is loaded. -- A plugin runs inside the proxy with your credentials, so opencodex refuses a file owned by another - user or writable by group or others. Fix it with `chmod go-w ~/.opencodex/plugins/*`. +- A plugin runs inside the proxy with your credentials, so opencodex refuses a plugin file or a + `plugins/` directory that is owned by another user or writable by group or others, and refuses + symbolic links. Fix permissions with `chmod go-w ~/.opencodex/plugins ~/.opencodex/plugins/*`. +- On Windows these owner and permission checks are not performed; only regular files are loaded. + Keep the `plugins/` directory writable by your account only. Restart the proxy after adding, changing or removing a plugin (`ocx service restart`, or stop and start `ocx start`). Each loaded plugin prints a `Plugin loaded: ` line at startup; a skipped @@ -37,7 +40,9 @@ To start once without plugins, set `OCX_PLUGINS=0`. A plugin default-exports an object with an optional `name` and a `setup` function. `setup` receives a context. An asynchronous `setup` has five seconds to finish; plugins run in the proxy's own thread, so -a `setup` that blocks synchronously cannot be interrupted and delays startup until it returns. +a `setup` that blocks synchronously cannot be interrupted and delays startup until it returns. A +`setup` that times out is not stopped either: servers or timers it already started keep running, so +start long-lived resources only after the work that can fail. ```ts interface UpstreamTarget { @@ -73,15 +78,20 @@ small interfaces you need locally, as above. - The rewriter runs synchronously on every provider send over HTTP and on the Codex WebSocket connection, after opencodex has picked the transport. Keep it fast; do network checks (health probes) in the background and read a cached result in the rewriter. -- Egress settings (proxy, `noProxy`) are applied to the rewritten HTTP destination. A Codex - WebSocket redirected to a loopback address connects directly instead of through a configured - proxy, which could not reach this machine's loopback. +- A send redirected to a loopback address (`127.0.0.1`, `::1`, `localhost`) connects directly, over + HTTP and over the Codex WebSocket, ignoring provider proxies and `HTTP_PROXY`: a proxy elsewhere + cannot reach this machine's loopback. Any other destination follows the normal egress settings, + evaluated against the rewritten URL. +- The Codex WebSocket rewriter runs for every turn, before an idle pooled socket is reused, and a + socket is only reused for the same destination. A plugin that starts or stops redirecting takes + effect on the next turn. - It runs after opencodex has chosen the provider, account and route, so it does not change routing, account selection, retries or request logs. - A redirected send goes to the host you chose. That host sees the request exactly as the provider would, credentials included. -- If a rewriter throws, opencodex undoes its changes to that send, disables it for the rest of the - process and sends the request unmodified. If `setup` throws or times out, the plugin is skipped, +- If a rewriter throws, opencodex undoes that rewriter's edits to the send and disables it for the + rest of the process. Edits made by rewriters that ran before it are kept, so the send goes out as + those left it (unmodified when it is the only plugin). If `setup` throws or times out, the plugin is skipped, anything it registered is removed, and later registration attempts from it are ignored; other plugins and the proxy start normally. - A plugin directory that exists but cannot be read (for example, wrong permissions) is reported at diff --git a/src/plugins/loader.ts b/src/plugins/loader.ts index 8e650311381..b888ce2f065 100644 --- a/src/plugins/loader.ts +++ b/src/plugins/loader.ts @@ -17,7 +17,7 @@ * run in the proxy's own thread, so synchronous work that never yields cannot be interrupted. */ -import { readdirSync, statSync } from "node:fs"; +import { lstatSync, readdirSync } from "node:fs"; import { basename, join } from "node:path"; import { pathToFileURL } from "node:url"; import { getConfigDir } from "../config/paths"; @@ -75,15 +75,21 @@ function listPluginFiles(dir: string): string[] { .map(entry => join(dir, entry)); } -/** Null when the file is safe to execute, otherwise the reason it is refused. */ -export function pluginFileTrustError(file: string): string | null { - let stats: ReturnType; +/** + * Null when `path` is safe to trust, otherwise the reason it is refused. `lstat` is used so a + * symbolic link is judged as a link — and refused — rather than as the file it points to: a + * link to a file you own would otherwise pass the owner and mode checks. Windows has no + * POSIX owner or mode bits, so there only the file type is checked. + */ +function trustError(path: string, kind: "file" | "directory"): string | null { + let stats: ReturnType; try { - stats = statSync(file); + stats = lstatSync(path); } catch (error) { return error instanceof Error ? error.message : String(error); } - if (!stats.isFile()) return "not a regular file"; + if (stats.isSymbolicLink()) return "is a symbolic link"; + if (kind === "file" ? !stats.isFile() : !stats.isDirectory()) return `not a regular ${kind}`; if (process.platform === "win32") return null; const uid = process.getuid?.(); if (uid !== undefined && stats.uid !== uid) return "owned by another user"; @@ -91,6 +97,19 @@ export function pluginFileTrustError(file: string): string | null { return null; } +/** Null when the file is safe to execute, otherwise the reason it is refused. */ +export function pluginFileTrustError(file: string): string | null { + return trustError(file, "file"); +} + +/** + * A directory another user can write lets them add or swap plugin files, so its owner and + * mode are checked like each file's. + */ +export function pluginDirectoryTrustError(dir: string): string | null { + return trustError(dir, "directory"); +} + function isPlugin(value: unknown): value is OcxPlugin { return typeof value === "object" && value !== null && typeof (value as OcxPlugin).setup === "function"; } @@ -123,6 +142,9 @@ export async function loadOcxPlugins( } catch (error) { return [{ file: dir, name: "plugins directory", loaded: false, error: error instanceof Error ? error.message : String(error) }]; } + if (files.length === 0) return []; + const dirRefused = pluginDirectoryTrustError(dir); + if (dirRefused) return [{ file: dir, name: "plugins directory", loaded: false, error: `refused: ${dirRefused}` }]; const results: PluginLoadResult[] = []; for (const file of files) { const fallbackName = basename(file).replace(/\.(ts|js|mjs)$/, ""); @@ -134,6 +156,7 @@ export async function loadOcxPlugins( const unregister: Array<() => void> = []; // Closed when setup fails or times out: a setup that resumes later must not register. let active = true; + let shutdownCount = 0; const whileActive = (name: string, register: () => () => void): void => { if (!active) { console.error(`[plugin:${name}] registration after a failed setup was ignored`); @@ -152,8 +175,9 @@ export async function loadOcxPlugins( pluginDir: dir, log: message => console.log(`[plugin:${name}] ${message}`), registerUpstreamRewriter: rewrite => whileActive(name, () => registerUpstreamRewriter(name, rewrite)), - // Keyed by file, not name: two plugins may share a display name. - onShutdown: teardown => whileActive(name, () => registerOptionalShutdownHook(`plugin:${file}`, teardown)), + // Keyed by file and registration, not name: two plugins may share a display name, and + // one plugin may register several teardowns. + onShutdown: teardown => whileActive(name, () => registerOptionalShutdownHook(`plugin:${file}#${++shutdownCount}`, teardown)), }; const timeoutMs = options.setupTimeoutMs ?? SETUP_TIMEOUT_MS; await withTimeout(Promise.resolve(plugin.setup(context)), timeoutMs, `plugin "${name}" setup`); diff --git a/src/plugins/upstream-hooks.ts b/src/plugins/upstream-hooks.ts index a12b068e990..77436d1b015 100644 --- a/src/plugins/upstream-hooks.ts +++ b/src/plugins/upstream-hooks.ts @@ -88,7 +88,7 @@ export function rewriteUpstreamRecord( return { url: result.url, headers: record }; } -function isLoopbackUrl(raw: string): boolean { +export function isLoopbackUrl(raw: string): boolean { let host: string; try { host = new URL(raw).hostname.toLowerCase().replace(/^\[|\]$/g, ""); @@ -99,9 +99,10 @@ function isLoopbackUrl(raw: string): boolean { } /** - * WebSocket dial variant. The caller chose `proxy` for the original destination; a proxy - * elsewhere on the network cannot reach this machine's loopback, so a rewrite onto a loopback - * sidecar dials directly. Any other rewrite keeps the caller's proxy decision. + * WebSocket dial variant, run for every exchange before a pooled socket is chosen so the pool + * identity can include the rewritten destination. The caller chose `proxy` for the original + * destination; a proxy elsewhere on the network cannot reach this machine's loopback, so a + * rewrite onto a loopback sidecar dials directly. Any other rewrite keeps the caller's proxy. */ export function rewriteWebSocketDial( url: string, diff --git a/src/server/responses/codex-ws-pool.ts b/src/server/responses/codex-ws-pool.ts index 5d406bee4f8..cb11ea643d3 100644 --- a/src/server/responses/codex-ws-pool.ts +++ b/src/server/responses/codex-ws-pool.ts @@ -25,7 +25,7 @@ function digest(input: unknown): string { } /** Identity comes from the selected outgoing request, never a model label or caller hint. */ -export function codexWsReuseIdentity(url: string, headers: Record, frameText: string, proxy?: string): CodexWsReuseIdentity | null { +export function codexWsReuseIdentity(url: string, headers: Record, frameText: string, proxy?: string, dialUrl?: string): CodexWsReuseIdentity | null { if (url !== CODEX_RESPONSES_HTTP_URL) return null; let body: unknown; try { body = JSON.parse(frameText); } catch { return null; } @@ -52,7 +52,8 @@ export function codexWsReuseIdentity(url: string, headers: Record, readonly retainable = false, private readonly changed: () => void = () => {}, proxy?: string) { - const dial = rewriteWebSocketDial(url, headers, proxy); - this.socket = new WebSocket(dial.url, { headers: dial.headers, ...(dial.proxy ? { proxy: dial.proxy } : {}) } as unknown as string[]); + this.socket = new WebSocket(url, { headers, ...(proxy ? { proxy } : {}) } as unknown as string[]); this.socket.addEventListener("open", this.onOpen); this.socket.addEventListener("message", this.onIdleMessage); this.socket.addEventListener("close", this.onClose); diff --git a/src/server/responses/fetch-helpers.ts b/src/server/responses/fetch-helpers.ts index 1f69baf45d5..02ebb6bc032 100644 --- a/src/server/responses/fetch-helpers.ts +++ b/src/server/responses/fetch-helpers.ts @@ -12,7 +12,7 @@ import { waitForProviderRequestSlot } from "../../providers/request-pacing"; import { withUpstreamHttpVersion } from "../../lib/upstream-http-version"; import type { CodexWsQuotaObserver } from "./codex-ws-metadata"; import { configuredOutboundFetch } from "../../lib/proxy-env"; -import { rewriteUpstream } from "../../plugins/upstream-hooks"; +import { isLoopbackUrl, rewriteUpstream } from "../../plugins/upstream-hooks"; import { describeProviderEgressForLog, markEgressTransparentExecutor, @@ -154,11 +154,18 @@ export function sendWithConnectionPolicy( // Plugin rewrites (src/plugins/upstream-hooks.ts) run here, after the caller chose between // the Codex WebSocket and HTTP, and before the connection and egress decisions below so // those follow the rewritten destination. Nested passes rewrite once, like the egress mark. + // A rewrite onto this machine's loopback dials directly: a proxy chosen for the provider + // (per-provider or HTTP_PROXY) cannot reach a local sidecar. The WebSocket dial does the same. const rewriteDone = (init as Record | undefined)?.[UPSTREAM_REWRITTEN] === true; - if (!rewriteDone && (typeof input === "string" || input instanceof URL)) { - const target = rewriteUpstream(String(input), headers, "http"); - input = target.url; + let redirectedToLoopback = false; + if (!rewriteDone) { + const original = input instanceof Request ? input.url : String(input); + const target = rewriteUpstream(original, headers, "http"); headers = target.headers as Headers; + if (target.url !== original) { + redirectedToLoopback = isLoopbackUrl(target.url); + input = input instanceof Request ? new Request(target.url, input) : target.url; + } } const fresh = wantsFreshConnection(input); if (fresh) { @@ -173,15 +180,17 @@ export function sendWithConnectionPolicy( // the reselected provider and the rebuilt destination, so it decides and marks the init; the // inner pass honours that mark rather than recomputing from a stale closure. const alreadyDecided = (init as Record | undefined)?.[EGRESS_DECIDED] === true; - const decide = egress !== undefined && !alreadyDecided; - const egressInit = decide ? providerEgressSendInit(egress, physicalFetch, input) : {}; + const decide = egress !== undefined && !alreadyDecided && !redirectedToLoopback; + const egressInit = redirectedToLoopback + ? { proxy: false as const } + : decide ? providerEgressSendInit(egress, physicalFetch, input) : {}; return physicalFetch(input, { ...init, headers, redirect: "manual", ...(fresh ? { keepalive: false } : {}), ...egressInit, - ...(decide ? { [EGRESS_DECIDED]: true } : {}), + ...(decide || redirectedToLoopback ? { [EGRESS_DECIDED]: true } : {}), ...{ [UPSTREAM_REWRITTEN]: true }, }); } diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index 7e7af0791e4..f3d7ab4b2a6 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -23,6 +23,7 @@ import { codexWsExchange } from "./codex-ws-exchange"; import { CodexWsSession } from "./codex-ws-session"; import { codexWsPool, codexWsReuseIdentity } from "./codex-ws-pool"; import { codexWsCreateFrameExceedsLimit } from "./codex-ws-wire"; +import { rewriteWebSocketDial } from "../../plugins/upstream-hooks"; export { CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, MAX_CODEX_WS_QUEUE_BYTES, MAX_CODEX_WS_CREATE_FRAME_BYTES, CODEX_WS_CREATE_FRAME_LIMIT_BYTES, codexWsCreateFrameExceedsLimit, isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "./codex-ws-wire"; @@ -176,9 +177,12 @@ export function codexWsUpstreamFetch( try { // Steering keeps a private physical connection across successor responses; it // must never enter the idle-socket pool or move to a different credential. - const identity = control ? null : codexWsReuseIdentity(url, headers, frameText, proxy); - session = (identity ? codexWsPool.acquire(identity, wsUrl, headers, proxy) : null) - ?? new CodexWsSession(wsUrl, headers, false, undefined, proxy); + // Plugin rewrite runs per exchange, before the pool lookup, and the dialled destination is + // part of the reuse identity: a socket opened to one destination is never reused for another. + const dial = rewriteWebSocketDial(wsUrl, headers, proxy); + const identity = control ? null : codexWsReuseIdentity(url, headers, frameText, dial.proxy, dial.url); + session = (identity ? codexWsPool.acquire(identity, dial.url, dial.headers, dial.proxy) : null) + ?? new CodexWsSession(dial.url, dial.headers, false, undefined, dial.proxy); if (!session.busy && !session.reserve()) { session.dispose(); return sseFallback(url, init); diff --git a/structure/ops/plugins.md b/structure/ops/plugins.md index f4654bb56c8..756cbbeb9bf 100644 --- a/structure/ops/plugins.md +++ b/structure/ops/plugins.md @@ -12,40 +12,48 @@ or signs them. - The loader imports `*.ts`, `*.js` and `*.mjs` from `$OPENCODEX_HOME/plugins/`, sorted by name. Names starting with `.` or `_` and `*.d.ts` are ignored. A missing directory loads nothing. - `OCX_PLUGINS=0` disables loading for that process. -- A plugin runs with the operator's credentials, so the loader refuses a file that is not a regular - file, is owned by another user, or is writable by group or others (POSIX). This is the same trust - boundary as `config.json`. +- A plugin runs with the operator's credentials, so the loader refuses a plugin file or plugin + directory that is a symbolic link (checked with `lstat`), is not a regular file/directory, is + owned by another user, or is writable by group or others. This is the same trust boundary as + `config.json`. Owner and mode checks are POSIX-only; on Windows only the file type is checked. - A missing plugin directory means no plugins. Any other read failure (`EACCES`, `ENOTDIR`) is reported as a skipped `plugins directory` entry. - A plugin module default-exports `{ name?, setup(context) }`. An asynchronous `setup` has five seconds; plugins share the proxy thread, so a setup that blocks synchronously cannot be interrupted. A plugin that throws, times out or has the wrong shape is reported and skipped: its context is closed, every hook it registered is removed, and a setup that resumes after the - deadline cannot register again. The other plugins and the proxy start normally. + deadline cannot register again. A timed-out setup keeps running; resources it already opened + are not closed. The other plugins and the proxy start normally. - Plugins cannot import ocx modules: in a compiled binary they live inside `$bunfs`. Everything a plugin may use arrives through `OcxPluginContext` (`name`, `configDir`, `pluginDir`, `log`, `registerUpstreamRewriter`, `onShutdown`). `onShutdown` registers through - `src/lib/optional-shutdown-hooks.ts` under a per-file key, so plugins sharing a display name keep - separate teardowns. + `src/lib/optional-shutdown-hooks.ts` under a per-file, per-registration key, so plugins sharing a + display name, and several teardowns from one plugin, all run. ## Upstream rewrite slot `src/plugins/upstream-hooks.ts` is the only core-owned seam plugins attach to. It imports nothing, so the request path depends on it without depending on the loader. -- It runs synchronously at the physical send, after the transport was chosen: HTTP in - `sendWithConnectionPolicy` (`src/server/responses/fetch-helpers.ts`), and the Codex WebSocket - dial in `CodexWsSession` (`src/server/responses/codex-ws-session.ts`). Rewriting any earlier would +- It runs synchronously after the transport was chosen: HTTP in `sendWithConnectionPolicy` + (`src/server/responses/fetch-helpers.ts`), including `Request` inputs, and the Codex WebSocket in + `codexWsUpstreamFetch` (`src/server/responses/ws-upstream.ts`) once per exchange, before the pool + lookup. The dialled destination is part of the reuse identity (`codexWsReuseIdentity` in + `src/server/responses/codex-ws-pool.ts`), so a socket is never reused for another destination. Rewriting any earlier would hide the ChatGPT origin from the WebSocket selection and push Codex turns onto HTTP. The target carries the URL, mutable headers and the transport (`http` or `websocket`). - `sendWithConnectionPolicy` can run twice for one send (an override handing back to the supplied executor). The outer pass rewrites and marks the init; the inner pass does not rewrite again. -- HTTP connection and egress decisions in `sendWithConnectionPolicy` follow the rewritten - destination. The WebSocket proxy is chosen by the caller for the original destination, so - `rewriteWebSocketDial` drops it when the rewrite targets loopback. +- A rewrite onto loopback dials directly on both transports: HTTP sends carry `proxy: false` and + mark egress as decided; `rewriteWebSocketDial` drops the proxy the caller chose for the original + destination. Other rewrites follow egress resolved against the rewritten HTTP URL. The pre-dispatch + egress refusal in `providerFetch` still validates the provider's configured route against the + original URL, so a misconfigured provider fails the same way with or without a plugin. - With no rewriter registered, the send is returned untouched and nothing is allocated. -- A rewriter that throws has its edits to that send undone, is disabled for the rest of the - process, and the send continues unmodified. +- A rewriter that throws has its own edits to that send undone and is disabled for the rest of the + process. Rollback is per rewriter: edits from rewriters that ran before it are kept, and the send + continues with them. `onShutdown` keys are unique per registration (`plugin:#`), so a + plugin may register several teardowns. - Rewrites happen after the request is built, routed and paced, so they do not change routing, account selection, retry budgets or logging identity. A rewriter that moves a send to another host owns that host's behaviour; the core does not re-validate it. diff --git a/tests/lib/plugin-loader.test.ts b/tests/lib/plugin-loader.test.ts index c7b4278d5b3..7bf78861c60 100644 --- a/tests/lib/plugin-loader.test.ts +++ b/tests/lib/plugin-loader.test.ts @@ -1,5 +1,5 @@ import { afterEach, beforeEach, expect, test } from "bun:test"; -import { chmodSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { chmodSync, mkdtempSync, rmSync, symlinkSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { resetOptionalShutdownHooksForTests, runOptionalShutdownHooks } from "../../src/lib/optional-shutdown-hooks"; @@ -67,6 +67,35 @@ test.skipIf(process.platform === "win32")("a group- or world-writable plugin is expect(hasUpstreamRewriters()).toBe(false); }); +test.skipIf(process.platform === "win32")("a symbolic link is refused even when it points to a trusted file", async () => { + const outside = mkdtempSync(join(tmpdir(), "ocx-plugin-target-")); + try { + const target = join(outside, "real.ts"); + writeFileSync(target, REDIRECT_PLUGIN); + chmodSync(target, 0o600); + symlinkSync(target, join(dir, "linked.ts")); + const [result] = await loadOcxPlugins(dir); + expect(result?.loaded).toBe(false); + expect(result?.error).toBe("refused: is a symbolic link"); + expect(hasUpstreamRewriters()).toBe(false); + } finally { + rmSync(outside, { recursive: true, force: true }); + } +}); + +test.skipIf(process.platform === "win32")("a plugin directory writable by group or others is refused", async () => { + writePlugin("redirect.ts", REDIRECT_PLUGIN); + chmodSync(dir, 0o775); + const results = await loadOcxPlugins(dir); + expect(results).toEqual([{ + file: dir, + name: "plugins directory", + loaded: false, + error: "refused: writable by group or others (chmod go-w)", + }]); + expect(hasUpstreamRewriters()).toBe(false); +}); + test("a wrong export shape or a throwing setup is skipped and leaves no hooks behind", async () => { writePlugin("a-shape.ts", "export default { name: 'shape' };"); writePlugin("b-throws.ts", ` @@ -120,6 +149,28 @@ export default { } }); +test("one plugin can register several shutdown teardowns", async () => { + const ran: string[] = []; + (globalThis as Record)["__ocxTeardownLog"] = ran; + writePlugin("multi.ts", ` +export default { + setup(ctx) { + ctx.onShutdown(() => { globalThis.__ocxTeardownLog.push("first"); }); + ctx.onShutdown(() => { globalThis.__ocxTeardownLog.push("second"); }); + }, +}; +`); + resetOptionalShutdownHooksForTests(); + try { + expect((await loadOcxPlugins(dir))[0]?.loaded).toBe(true); + runOptionalShutdownHooks(); + expect(ran.sort()).toEqual(["first", "second"]); + } finally { + resetOptionalShutdownHooksForTests(); + delete (globalThis as Record)["__ocxTeardownLog"]; + } +}); + test("a setup that resumes after its deadline cannot leave registrations behind", async () => { writePlugin("slow.ts", ` export default { diff --git a/tests/lib/plugin-upstream-hooks.test.ts b/tests/lib/plugin-upstream-hooks.test.ts index f3f78579100..fa70094beaa 100644 --- a/tests/lib/plugin-upstream-hooks.test.ts +++ b/tests/lib/plugin-upstream-hooks.test.ts @@ -8,6 +8,8 @@ import { rewriteWebSocketDial, } from "../../src/plugins/upstream-hooks"; import { sendWithConnectionPolicy } from "../../src/server/responses/fetch-helpers"; +import { codexWsReuseIdentity } from "../../src/server/responses/codex-ws-pool"; +import { CODEX_RESPONSES_HTTP_URL } from "../../src/server/responses/codex-ws-request"; afterEach(() => resetUpstreamRewritersForTests()); @@ -117,6 +119,44 @@ test("the physical HTTP send rewrites once, even through a nested override pass" expect(seen).toEqual([{ url: "http://127.0.0.1:8787/v1/responses", hop: "1" }]); }); +test("an HTTP send redirected to loopback dials directly, bypassing any proxy", async () => { + const seen: Array<{ url: string; proxy: unknown }> = []; + const physical = (async (input: Parameters[0], init?: RequestInit) => { + seen.push({ url: input instanceof Request ? input.url : String(input), proxy: (init as { proxy?: unknown }).proxy }); + return new Response("ok"); + }) as typeof fetch; + registerUpstreamRewriter("loopback", target => { target.url = target.url.replace("https://api.example.com", "http://127.0.0.1:8787"); }); + await sendWithConnectionPolicy(physical, "https://api.example.com/v1/responses", { method: "POST" }); + resetUpstreamRewritersForTests(); + registerUpstreamRewriter("remote", target => { target.url = target.url.replace("https://api.example.com", "https://relay.example.net"); }); + await sendWithConnectionPolicy(physical, "https://api.example.com/v1/responses", { method: "POST" }); + expect(seen).toEqual([ + { url: "http://127.0.0.1:8787/v1/responses", proxy: false }, + { url: "https://relay.example.net/v1/responses", proxy: undefined }, + ]); +}); + +test("a Request input is rewritten too", async () => { + let seenUrl = ""; + const physical = (async (input: Parameters[0]) => { + seenUrl = input instanceof Request ? input.url : String(input); + return new Response("ok"); + }) as typeof fetch; + registerUpstreamRewriter("loopback", target => { target.url = "http://127.0.0.1:8787/v1/messages"; }); + await sendWithConnectionPolicy(physical, new Request("https://api.example.com/v1/messages", { method: "POST", body: "{}" })); + expect(seenUrl).toBe("http://127.0.0.1:8787/v1/messages"); +}); + +test("the Codex WebSocket reuse identity changes with the dialled destination", () => { + const headers = { authorization: "Bearer t", "chatgpt-account-id": "acct", "thread-id": "th" }; + const frame = JSON.stringify({ model: "gpt-x", client_metadata: { thread_id: "th", turn_id: "tu" } }); + const direct = codexWsReuseIdentity(CODEX_RESPONSES_HTTP_URL, headers, frame, undefined, "wss://chatgpt.com/backend-api/codex/responses"); + const local = codexWsReuseIdentity(CODEX_RESPONSES_HTTP_URL, headers, frame, undefined, "ws://127.0.0.1:8787/backend-api/codex/responses"); + expect(direct).not.toBeNull(); + expect(local).not.toBeNull(); + expect(local?.key).not.toBe(direct?.key); +}); + test("unregistering removes the rewriter", () => { const off = registerUpstreamRewriter("temp", target => { target.url = "http://changed/"; }); off(); From 8515237e0e6e06dd3acf3a06cee272a33b8abe67 Mon Sep 17 00:00:00 2001 From: halysondev Date: Sat, 26 Sep 2026 00:29:08 -0300 Subject: [PATCH 4/6] fix(plugins): trust the whole plugin path and settle rewritten WS routes Addresses the third CodeRabbit pass on #5896. - Every ancestor of the resolved plugin directory up to `/` must be owned by the user or root and not group/other-writable unless sticky, so no other user can swap a checked path before import (OpenSSH StrictModes rule). Plugins are imported through the resolved directory. - The Codex WebSocket reuse identity hashes the rewritten headers, so a pooled socket is never reused with stale plugin headers. - A non-loopback WebSocket rewrite resolves its own proxy route (scheme, NO_PROXY), falling back to SSE when that route requires it, via the new planCodexWsDial helper. - Loader tests drop group write on the runner's own nested temp roots, which inherit umask 002 on user-private-group systems. Co-Authored-By: Claude Opus 5.5 --- .../src/content/docs/guides/local-plugins.md | 11 +++-- src/plugins/loader.ts | 43 +++++++++++++++++-- src/server/responses/ws-upstream.ts | 37 +++++++++++++--- structure/ops/plugins.md | 14 ++++-- tests/lib/plugin-loader.test.ts | 38 ++++++++++++++-- tests/lib/plugin-upstream-hooks.test.ts | 21 +++++++++ 6 files changed, 145 insertions(+), 19 deletions(-) diff --git a/docs-site/src/content/docs/guides/local-plugins.md b/docs-site/src/content/docs/guides/local-plugins.md index 2d8bf955b17..b802edde3ca 100644 --- a/docs-site/src/content/docs/guides/local-plugins.md +++ b/docs-site/src/content/docs/guides/local-plugins.md @@ -26,7 +26,10 @@ Put plugin files in `plugins/` inside the opencodex home (`~/.opencodex/plugins/ - The directory is optional. Without it nothing is loaded. - A plugin runs inside the proxy with your credentials, so opencodex refuses a plugin file or a `plugins/` directory that is owned by another user or writable by group or others, and refuses - symbolic links. Fix permissions with `chmod go-w ~/.opencodex/plugins ~/.opencodex/plugins/*`. + symbolic links. Every directory above `plugins/`, up to `/`, must also be owned by you or root and + not writable by group or others, unless it is sticky like `/tmp`. Fix permissions with + `chmod go-w ~/.opencodex/plugins ~/.opencodex/plugins/*`; on systems whose default umask is + `002`, check the parent directories too. - On Windows these owner and permission checks are not performed; only regular files are loaded. Keep the `plugins/` directory writable by your account only. @@ -80,10 +83,10 @@ small interfaces you need locally, as above. probes) in the background and read a cached result in the rewriter. - A send redirected to a loopback address (`127.0.0.1`, `::1`, `localhost`) connects directly, over HTTP and over the Codex WebSocket, ignoring provider proxies and `HTTP_PROXY`: a proxy elsewhere - cannot reach this machine's loopback. Any other destination follows the normal egress settings, - evaluated against the rewritten URL. + cannot reach this machine's loopback. Any other destination follows the normal egress settings + (including `NO_PROXY`), evaluated against the rewritten URL, on both transports. - The Codex WebSocket rewriter runs for every turn, before an idle pooled socket is reused, and a - socket is only reused for the same destination. A plugin that starts or stops redirecting takes + socket is only reused for the same destination and the same rewritten headers. A plugin that starts or stops redirecting takes effect on the next turn. - It runs after opencodex has chosen the provider, account and route, so it does not change routing, account selection, retries or request logs. diff --git a/src/plugins/loader.ts b/src/plugins/loader.ts index b888ce2f065..aaa4e27a4cc 100644 --- a/src/plugins/loader.ts +++ b/src/plugins/loader.ts @@ -17,8 +17,8 @@ * run in the proxy's own thread, so synchronous work that never yields cannot be interrupted. */ -import { lstatSync, readdirSync } from "node:fs"; -import { basename, join } from "node:path"; +import { lstatSync, readdirSync, realpathSync } from "node:fs"; +import { basename, dirname, join } from "node:path"; import { pathToFileURL } from "node:url"; import { getConfigDir } from "../config/paths"; import { registerOptionalShutdownHook } from "../lib/optional-shutdown-hooks"; @@ -110,6 +110,34 @@ export function pluginDirectoryTrustError(dir: string): string | null { return trustError(dir, "directory"); } +/** + * Every directory above the (resolved) plugin directory, up to `/`, must be owned by the user + * or root and not writable by group or others unless it is sticky (like `/tmp`), where others + * cannot rename or replace entries they do not own. With no writable component on the path, no + * other user can swap what the loader checked for something else before it is imported + * (OpenSSH StrictModes applies the same rule). POSIX only. + */ +export function pluginAncestorsTrustError(realDir: string): string | null { + if (process.platform === "win32") return null; + const uid = process.getuid?.(); + let current = dirname(realDir); + for (;;) { + let stats: ReturnType; + try { + stats = lstatSync(current); + } catch (error) { + return `${current}: ${error instanceof Error ? error.message : String(error)}`; + } + if (uid !== undefined && stats.uid !== uid && stats.uid !== 0) return `${current} is owned by another user`; + if ((stats.mode & 0o022) !== 0 && (stats.mode & 0o1000) === 0) { + return `${current} is writable by group or others`; + } + const parent = dirname(current); + if (parent === current) return null; + current = parent; + } +} + function isPlugin(value: unknown): value is OcxPlugin { return typeof value === "object" && value !== null && typeof (value as OcxPlugin).setup === "function"; } @@ -145,8 +173,17 @@ export async function loadOcxPlugins( if (files.length === 0) return []; const dirRefused = pluginDirectoryTrustError(dir); if (dirRefused) return [{ file: dir, name: "plugins directory", loaded: false, error: `refused: ${dirRefused}` }]; + // Check and import through the resolved path, so both refer to the same components. + let realDir: string; + try { + realDir = realpathSync(dir); + } catch (error) { + return [{ file: dir, name: "plugins directory", loaded: false, error: error instanceof Error ? error.message : String(error) }]; + } + const ancestorRefused = pluginAncestorsTrustError(realDir); + if (ancestorRefused) return [{ file: dir, name: "plugins directory", loaded: false, error: `refused: ${ancestorRefused}` }]; const results: PluginLoadResult[] = []; - for (const file of files) { + for (const file of files.map(listed => join(realDir, basename(listed)))) { const fallbackName = basename(file).replace(/\.(ts|js|mjs)$/, ""); const refused = pluginFileTrustError(file); if (refused) { diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index f3d7ab4b2a6..d9ab3098dd9 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -23,7 +23,7 @@ import { codexWsExchange } from "./codex-ws-exchange"; import { CodexWsSession } from "./codex-ws-session"; import { codexWsPool, codexWsReuseIdentity } from "./codex-ws-pool"; import { codexWsCreateFrameExceedsLimit } from "./codex-ws-wire"; -import { rewriteWebSocketDial } from "../../plugins/upstream-hooks"; +import { isLoopbackUrl, rewriteWebSocketDial } from "../../plugins/upstream-hooks"; export { CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, MAX_CODEX_WS_QUEUE_BYTES, MAX_CODEX_WS_CREATE_FRAME_BYTES, CODEX_WS_CREATE_FRAME_LIMIT_BYTES, codexWsCreateFrameExceedsLimit, isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "./codex-ws-wire"; @@ -82,6 +82,31 @@ export function bunSupportsBoundedCodexWsRelay( return comparison !== null && comparison >= 0; } +/** + * Apply plugin rewrites to one Codex WebSocket dial and settle its proxy. `proxy` was resolved + * for the canonical `wsUrl`. A loopback rewrite dials directly; any other rewrite gets its own + * route (scheme and `NO_PROXY` may differ). Null means the rewritten destination needs the SSE + * fallback, exactly as an unusable route for the canonical URL does. + */ +export function planCodexWsDial( + wsUrl: string, + headers: Record, + proxy: string | undefined, + env: Parameters[1] = process.env, +): { url: string; headers: Record; proxy: string | undefined } | null { + const dial = rewriteWebSocketDial(wsUrl, headers, proxy); + if (dial.url === wsUrl || isLoopbackUrl(dial.url)) return dial; + let rewritten: URL; + try { + rewritten = new URL(dial.url); + } catch { + return null; + } + const route = resolveProxyRoute(rewritten, env); + if (route.kind === "fallback") return null; + return { ...dial, proxy: route.kind === "proxy" ? route.proxy : undefined }; +} + export function shouldUseCodexWsUpstream( url: string, init?: RequestInit, @@ -177,10 +202,12 @@ export function codexWsUpstreamFetch( try { // Steering keeps a private physical connection across successor responses; it // must never enter the idle-socket pool or move to a different credential. - // Plugin rewrite runs per exchange, before the pool lookup, and the dialled destination is - // part of the reuse identity: a socket opened to one destination is never reused for another. - const dial = rewriteWebSocketDial(wsUrl, headers, proxy); - const identity = control ? null : codexWsReuseIdentity(url, headers, frameText, dial.proxy, dial.url); + // Plugin rewrite runs per exchange, before the pool lookup. The dialled destination, its + // headers and its proxy are all part of the reuse identity, so a pooled socket is never + // reused for a different destination or with stale plugin headers. + const dial = planCodexWsDial(wsUrl, headers, proxy); + if (!dial) return sseFallback(url, init); + const identity = control ? null : codexWsReuseIdentity(url, dial.headers, frameText, dial.proxy, dial.url); session = (identity ? codexWsPool.acquire(identity, dial.url, dial.headers, dial.proxy) : null) ?? new CodexWsSession(dial.url, dial.headers, false, undefined, dial.proxy); if (!session.busy && !session.reserve()) { diff --git a/structure/ops/plugins.md b/structure/ops/plugins.md index 756cbbeb9bf..5e3c1fc57e7 100644 --- a/structure/ops/plugins.md +++ b/structure/ops/plugins.md @@ -15,7 +15,10 @@ or signs them. - A plugin runs with the operator's credentials, so the loader refuses a plugin file or plugin directory that is a symbolic link (checked with `lstat`), is not a regular file/directory, is owned by another user, or is writable by group or others. This is the same trust boundary as - `config.json`. Owner and mode checks are POSIX-only; on Windows only the file type is checked. + `config.json`. Every ancestor of the resolved plugin directory up to `/` must be owned by the user + or root and not group/other-writable unless sticky (`pluginAncestorsTrustError`), so no other user + can swap a checked path before it is imported; files are imported through the resolved directory. + Owner and mode checks are POSIX-only; on Windows only the file type is checked. - A missing plugin directory means no plugins. Any other read failure (`EACCES`, `ENOTDIR`) is reported as a skipped `plugins directory` entry. - A plugin module default-exports `{ name?, setup(context) }`. An asynchronous `setup` has five @@ -38,15 +41,18 @@ so the request path depends on it without depending on the loader. - It runs synchronously after the transport was chosen: HTTP in `sendWithConnectionPolicy` (`src/server/responses/fetch-helpers.ts`), including `Request` inputs, and the Codex WebSocket in `codexWsUpstreamFetch` (`src/server/responses/ws-upstream.ts`) once per exchange, before the pool - lookup. The dialled destination is part of the reuse identity (`codexWsReuseIdentity` in - `src/server/responses/codex-ws-pool.ts`), so a socket is never reused for another destination. Rewriting any earlier would + lookup. `planCodexWsDial` applies the rewrite and settles the proxy; the dialled destination, rewritten + headers and proxy are part of the reuse identity (`codexWsReuseIdentity` in + `src/server/responses/codex-ws-pool.ts`), so a socket is never reused for another destination or + with stale plugin headers. Rewriting any earlier would hide the ChatGPT origin from the WebSocket selection and push Codex turns onto HTTP. The target carries the URL, mutable headers and the transport (`http` or `websocket`). - `sendWithConnectionPolicy` can run twice for one send (an override handing back to the supplied executor). The outer pass rewrites and marks the init; the inner pass does not rewrite again. - A rewrite onto loopback dials directly on both transports: HTTP sends carry `proxy: false` and mark egress as decided; `rewriteWebSocketDial` drops the proxy the caller chose for the original - destination. Other rewrites follow egress resolved against the rewritten HTTP URL. The pre-dispatch + destination. Other rewrites resolve their route against the rewritten URL on both transports (a + WebSocket route that needs the SSE fallback falls back, as it would for the canonical URL). The pre-dispatch egress refusal in `providerFetch` still validates the provider's configured route against the original URL, so a misconfigured provider fails the same way with or without a plugin. - With no rewriter registered, the send is returned untouched and nothing is allocated. diff --git a/tests/lib/plugin-loader.test.ts b/tests/lib/plugin-loader.test.ts index 7bf78861c60..5a754a19e3e 100644 --- a/tests/lib/plugin-loader.test.ts +++ b/tests/lib/plugin-loader.test.ts @@ -1,7 +1,7 @@ -import { afterEach, beforeEach, expect, test } from "bun:test"; -import { chmodSync, mkdtempSync, rmSync, symlinkSync, writeFileSync } from "node:fs"; +import { afterEach, beforeAll, beforeEach, expect, test } from "bun:test"; +import { chmodSync, mkdirSync, mkdtempSync, realpathSync, rmSync, statSync, symlinkSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; -import { join } from "node:path"; +import { dirname, join } from "node:path"; import { resetOptionalShutdownHooksForTests, runOptionalShutdownHooks } from "../../src/lib/optional-shutdown-hooks"; import { loadOcxPlugins, pluginFileTrustError } from "../../src/plugins/loader"; import { @@ -12,6 +12,24 @@ import { let dir: string; +// The loader refuses a plugin directory with a group- or other-writable, non-sticky ancestor. +// The test runner nests per-process temp roots and creates them with the caller's umask, which is +// group-writable on user-private-group systems (umask 002). Those directories belong to this +// test run, so drop group/other write on every ancestor this user owns. The root-owned sticky +// `/tmp` above them is accepted as is, which also exercises the sticky exception. +beforeAll(() => { + if (process.platform === "win32") return; + for (let current = realpathSync(tmpdir()); ; current = dirname(current)) { + try { + const stats = statSync(current); + if (stats.uid === process.getuid?.() && (stats.mode & 0o022) !== 0 && (stats.mode & 0o1000) === 0) { + chmodSync(current, stats.mode & 0o7755); + } + } catch { /* leave directories this run cannot inspect alone */ } + if (dirname(current) === current) break; + } +}); + beforeEach(() => { dir = mkdtempSync(join(tmpdir(), "ocx-plugins-")); delete process.env["OCX_PLUGINS"]; @@ -96,6 +114,20 @@ test.skipIf(process.platform === "win32")("a plugin directory writable by group expect(hasUpstreamRewriters()).toBe(false); }); +test.skipIf(process.platform === "win32")("a plugin directory under a group-writable, non-sticky parent is refused", async () => { + const parent = join(dir, "shared"); + const nested = join(parent, "plugins"); + mkdirSync(nested, { recursive: true, mode: 0o700 }); + chmodSync(parent, 0o775); + writeFileSync(join(nested, "redirect.ts"), REDIRECT_PLUGIN); + chmodSync(join(nested, "redirect.ts"), 0o600); + const results = await loadOcxPlugins(nested); + expect(results).toHaveLength(1); + expect(results[0]?.loaded).toBe(false); + expect(results[0]?.error).toContain("is writable by group or others"); + expect(hasUpstreamRewriters()).toBe(false); +}); + test("a wrong export shape or a throwing setup is skipped and leaves no hooks behind", async () => { writePlugin("a-shape.ts", "export default { name: 'shape' };"); writePlugin("b-throws.ts", ` diff --git a/tests/lib/plugin-upstream-hooks.test.ts b/tests/lib/plugin-upstream-hooks.test.ts index fa70094beaa..2b1615cf641 100644 --- a/tests/lib/plugin-upstream-hooks.test.ts +++ b/tests/lib/plugin-upstream-hooks.test.ts @@ -8,6 +8,7 @@ import { rewriteWebSocketDial, } from "../../src/plugins/upstream-hooks"; import { sendWithConnectionPolicy } from "../../src/server/responses/fetch-helpers"; +import { planCodexWsDial } from "../../src/server/responses/ws-upstream"; import { codexWsReuseIdentity } from "../../src/server/responses/codex-ws-pool"; import { CODEX_RESPONSES_HTTP_URL } from "../../src/server/responses/codex-ws-request"; @@ -157,6 +158,26 @@ test("the Codex WebSocket reuse identity changes with the dialled destination", expect(local?.key).not.toBe(direct?.key); }); +test("a Codex WebSocket dial resolves its proxy for the rewritten destination", () => { + const ws = "wss://chatgpt.com/backend-api/codex/responses"; + const envProxy = { HTTPS_PROXY: "http://corp:3128", HTTP_PROXY: "http://plain:8080" }; + expect(planCodexWsDial(ws, {}, "http://corp:3128", envProxy)?.proxy).toBe("http://corp:3128"); + + const off = registerUpstreamRewriter("remote-wss", target => { target.url = "wss://relay.example.net/backend-api/codex/responses"; }); + expect(planCodexWsDial(ws, {}, "http://corp:3128", envProxy)?.proxy).toBe("http://corp:3128"); + expect(planCodexWsDial(ws, {}, "http://corp:3128", { ...envProxy, NO_PROXY: "relay.example.net" })?.proxy).toBeUndefined(); + off(); + + const offPlain = registerUpstreamRewriter("remote-ws", target => { target.url = "ws://relay.example.net/backend-api/codex/responses"; }); + expect(planCodexWsDial(ws, {}, "http://corp:3128", envProxy)?.proxy).toBe("http://plain:8080"); + offPlain(); + + registerUpstreamRewriter("loopback", target => { target.url = "ws://127.0.0.1:8787/backend-api/codex/responses"; }); + expect(planCodexWsDial(ws, {}, "http://corp:3128", envProxy)).toEqual({ + url: "ws://127.0.0.1:8787/backend-api/codex/responses", headers: {}, proxy: undefined, + }); +}); + test("unregistering removes the rewriter", () => { const off = registerUpstreamRewriter("temp", target => { target.url = "http://changed/"; }); off(); From 6bfe9af0a01a4904019acb14c91c49c35c0d2b11 Mon Sep 17 00:00:00 2001 From: halysondev Date: Sat, 26 Sep 2026 00:40:17 -0300 Subject: [PATCH 5/6] fix(plugins): keep frame-carried Codex turn headers out of WS rewrites Addresses the fourth CodeRabbit pass on #5896. x-codex-turn-state and x-codex-turn-metadata ride in each WebSocket frame's client_metadata, prepared before the plugin rewrite and authoritative per exchange; the pool deliberately leaves them out of the reuse identity. planCodexWsDial now restores their original values after the rewrite, so a rewriter cannot make the upgrade headers disagree with the frame, and socket reuse keeps working. The header list is exported once as CODEX_WS_FRAME_HEADERS and shared by the frame builder and pool. Co-Authored-By: Claude Opus 5.5 --- .../src/content/docs/guides/local-plugins.md | 4 +++- src/server/responses/codex-ws-pool.ts | 4 ++-- src/server/responses/codex-ws-request.ts | 5 ++++- src/server/responses/ws-upstream.ts | 19 ++++++++++++++----- structure/ops/plugins.md | 5 ++++- tests/lib/plugin-upstream-hooks.test.ts | 18 ++++++++++++++++++ 6 files changed, 45 insertions(+), 10 deletions(-) diff --git a/docs-site/src/content/docs/guides/local-plugins.md b/docs-site/src/content/docs/guides/local-plugins.md index b802edde3ca..bbf5801505b 100644 --- a/docs-site/src/content/docs/guides/local-plugins.md +++ b/docs-site/src/content/docs/guides/local-plugins.md @@ -86,7 +86,9 @@ small interfaces you need locally, as above. cannot reach this machine's loopback. Any other destination follows the normal egress settings (including `NO_PROXY`), evaluated against the rewritten URL, on both transports. - The Codex WebSocket rewriter runs for every turn, before an idle pooled socket is reused, and a - socket is only reused for the same destination and the same rewritten headers. A plugin that starts or stops redirecting takes + socket is only reused for the same destination and the same rewritten headers. On that transport + `x-codex-turn-state` and `x-codex-turn-metadata` travel inside each request frame, so changes a + rewriter makes to those two headers are discarded. A plugin that starts or stops redirecting takes effect on the next turn. - It runs after opencodex has chosen the provider, account and route, so it does not change routing, account selection, retries or request logs. diff --git a/src/server/responses/codex-ws-pool.ts b/src/server/responses/codex-ws-pool.ts index cb11ea643d3..73c99d5a89b 100644 --- a/src/server/responses/codex-ws-pool.ts +++ b/src/server/responses/codex-ws-pool.ts @@ -1,13 +1,13 @@ import { createHmac, randomBytes } from "node:crypto"; import { registerOptionalShutdownHook } from "../../lib/optional-shutdown-hooks"; -import { CODEX_RESPONSES_HTTP_URL } from "./codex-ws-request"; +import { CODEX_RESPONSES_HTTP_URL, CODEX_WS_FRAME_HEADERS } from "./codex-ws-request"; import { CODEX_WS_ID_MAX_BYTES } from "./codex-ws-correlation"; import { CodexWsSession } from "./codex-ws-session"; export const CODEX_WS_POOL_MAX_SESSIONS = 32; export const CODEX_WS_POOL_IDLE_MS = 30_000; export const CODEX_WS_POOL_MAX_AGE_MS = 5 * 60_000; -const MUTABLE_HEADERS = new Set(["x-codex-turn-state", "x-codex-turn-metadata"]); +const MUTABLE_HEADERS = new Set(CODEX_WS_FRAME_HEADERS); let processKey: Buffer | undefined; let poolSequence = 0; diff --git a/src/server/responses/codex-ws-request.ts b/src/server/responses/codex-ws-request.ts index eaa77323e77..2e312b837bf 100644 --- a/src/server/responses/codex-ws-request.ts +++ b/src/server/responses/codex-ws-request.ts @@ -4,6 +4,9 @@ import { CODEX_RESPONSES_LITE_METADATA_KEY, } from "../../codex/forward-transport-headers"; +/** Per-turn headers the WebSocket carries in `client_metadata` of each frame, not in the upgrade. */ +export const CODEX_WS_FRAME_HEADERS = ["x-codex-turn-state", "x-codex-turn-metadata"] as const; + export const CODEX_RESPONSES_HTTP_URL = "https://chatgpt.com/backend-api/codex/responses"; export const CODEX_RESPONSES_WS_URL = "wss://chatgpt.com/backend-api/codex/responses"; export const WS_BETA = "responses_websockets=2026-02-06"; @@ -32,7 +35,7 @@ function applyLiteMetadata(body: Record, headers: Headers): boo body.client_metadata = { ...(metadata as Record | undefined), [CODEX_RESPONSES_LITE_METADATA_KEY]: lite }; } - for (const name of ["x-codex-turn-state", "x-codex-turn-metadata"]) { + for (const name of CODEX_WS_FRAME_HEADERS) { const value = headers.get(name); const current = body.client_metadata as Record | undefined; if (value !== null && !Object.hasOwn(current ?? {}, name)) { diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index d9ab3098dd9..4363a39d24d 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -18,7 +18,7 @@ import type { NativeResponseControl } from "./native-response-control"; import { compareBunVersions } from "../../lib/bun-stream-caps"; import { resolveProxyRoute, socks5ProxyFromEnv } from "../../lib/proxy-env"; import type { CodexWsQuotaObserver } from "./codex-ws-metadata"; -import { CODEX_RESPONSES_HTTP_URL, CODEX_RESPONSES_WS_URL, prepareCodexHttpInit, prepareCodexWsRequest } from "./codex-ws-request"; +import { CODEX_RESPONSES_HTTP_URL, CODEX_RESPONSES_WS_URL, CODEX_WS_FRAME_HEADERS, prepareCodexHttpInit, prepareCodexWsRequest } from "./codex-ws-request"; import { codexWsExchange } from "./codex-ws-exchange"; import { CodexWsSession } from "./codex-ws-session"; import { codexWsPool, codexWsReuseIdentity } from "./codex-ws-pool"; @@ -94,15 +94,24 @@ export function planCodexWsDial( proxy: string | undefined, env: Parameters[1] = process.env, ): { url: string; headers: Record; proxy: string | undefined } | null { - const dial = rewriteWebSocketDial(wsUrl, headers, proxy); + const rewritten = rewriteWebSocketDial(wsUrl, headers, proxy); + // Per-turn headers ride in each frame's client_metadata, which was prepared before the rewrite + // and is authoritative; a pooled socket's upgrade copy is intentionally ignored. A rewriter + // therefore cannot change them here, and the upgrade keeps the values the frame carries. + const dialHeaders = { ...rewritten.headers }; + for (const name of CODEX_WS_FRAME_HEADERS) { + if (Object.hasOwn(headers, name)) dialHeaders[name] = headers[name]!; + else delete dialHeaders[name]; + } + const dial = { ...rewritten, headers: dialHeaders }; if (dial.url === wsUrl || isLoopbackUrl(dial.url)) return dial; - let rewritten: URL; + let destination: URL; try { - rewritten = new URL(dial.url); + destination = new URL(dial.url); } catch { return null; } - const route = resolveProxyRoute(rewritten, env); + const route = resolveProxyRoute(destination, env); if (route.kind === "fallback") return null; return { ...dial, proxy: route.kind === "proxy" ? route.proxy : undefined }; } diff --git a/structure/ops/plugins.md b/structure/ops/plugins.md index 5e3c1fc57e7..6912a501fc4 100644 --- a/structure/ops/plugins.md +++ b/structure/ops/plugins.md @@ -44,7 +44,10 @@ so the request path depends on it without depending on the loader. lookup. `planCodexWsDial` applies the rewrite and settles the proxy; the dialled destination, rewritten headers and proxy are part of the reuse identity (`codexWsReuseIdentity` in `src/server/responses/codex-ws-pool.ts`), so a socket is never reused for another destination or - with stale plugin headers. Rewriting any earlier would + with stale plugin headers. The per-turn headers in `CODEX_WS_FRAME_HEADERS` + (`src/server/responses/codex-ws-request.ts`) ride in each frame's `client_metadata`, prepared + before the rewrite and authoritative, so `planCodexWsDial` restores their original values and + they stay outside the reuse identity. Rewriting any earlier would hide the ChatGPT origin from the WebSocket selection and push Codex turns onto HTTP. The target carries the URL, mutable headers and the transport (`http` or `websocket`). - `sendWithConnectionPolicy` can run twice for one send (an override handing back to the supplied diff --git a/tests/lib/plugin-upstream-hooks.test.ts b/tests/lib/plugin-upstream-hooks.test.ts index 2b1615cf641..4d72b5b292a 100644 --- a/tests/lib/plugin-upstream-hooks.test.ts +++ b/tests/lib/plugin-upstream-hooks.test.ts @@ -178,6 +178,24 @@ test("a Codex WebSocket dial resolves its proxy for the rewritten destination", }); }); +test("a Codex WebSocket rewriter cannot change the per-turn headers carried in the frame", () => { + registerUpstreamRewriter("turn-headers", target => { + target.headers.set("x-codex-turn-state", "rewritten"); + target.headers.set("x-codex-turn-metadata", "added"); + target.headers.set("x-sidecar", "1"); + }); + const dial = planCodexWsDial( + "wss://chatgpt.com/backend-api/codex/responses", + { "x-codex-turn-state": "original", authorization: "Bearer t" }, + undefined, + {}, + ); + expect(dial?.headers["x-codex-turn-state"]).toBe("original"); + expect(Object.hasOwn(dial?.headers ?? {}, "x-codex-turn-metadata")).toBe(false); + expect(dial?.headers["x-sidecar"]).toBe("1"); + expect(dial?.headers.authorization).toBe("Bearer t"); +}); + test("unregistering removes the rewriter", () => { const off = registerUpstreamRewriter("temp", target => { target.url = "http://changed/"; }); off(); From 4e6409ffa4badfb152f4c30ad0dac57ef1ca58f9 Mon Sep 17 00:00:00 2001 From: halysondev Date: Sat, 26 Sep 2026 00:48:28 -0300 Subject: [PATCH 6/6] docs(plugins): exclude per-turn headers from the WebSocket reuse wording Co-Authored-By: Claude Opus 5.5 --- docs-site/src/content/docs/guides/local-plugins.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/docs-site/src/content/docs/guides/local-plugins.md b/docs-site/src/content/docs/guides/local-plugins.md index bbf5801505b..a7864b7d4c1 100644 --- a/docs-site/src/content/docs/guides/local-plugins.md +++ b/docs-site/src/content/docs/guides/local-plugins.md @@ -86,9 +86,10 @@ small interfaces you need locally, as above. cannot reach this machine's loopback. Any other destination follows the normal egress settings (including `NO_PROXY`), evaluated against the rewritten URL, on both transports. - The Codex WebSocket rewriter runs for every turn, before an idle pooled socket is reused, and a - socket is only reused for the same destination and the same rewritten headers. On that transport - `x-codex-turn-state` and `x-codex-turn-metadata` travel inside each request frame, so changes a - rewriter makes to those two headers are discarded. A plugin that starts or stops redirecting takes + socket is only reused for the same destination and the same rewritten headers, apart from the two + per-turn headers `x-codex-turn-state` and `x-codex-turn-metadata`. Those travel inside each + request frame and may differ between exchanges on one socket; changes a rewriter makes to them + are discarded. A plugin that starts or stops redirecting takes effect on the next turn. - It runs after opencodex has chosen the provider, account and route, so it does not change routing, account selection, retries or request logs.