From e34eaefeb0dd6377201e0fbb386668a16f0dd74e Mon Sep 17 00:00:00 2001 From: ninja Date: Tue, 1 Sep 2026 13:17:42 -0400 Subject: [PATCH 1/2] feat: add local Codex subscription coworker --- .env.example | 14 +- .gitignore | 1 + agent-codex/README.md | 11 ++ agent-codex/bun.lock | 40 ++++ agent-codex/package.json | 21 ++ agent-codex/src/codex-client.ts | 316 ++++++++++++++++++++++++++++++ agent-codex/src/history.ts | 45 +++++ agent-codex/src/index.ts | 142 ++++++++++++++ agent-codex/tests/history.test.ts | 31 +++ agent-codex/tsconfig.json | 11 ++ examples/codex/agents.yaml | 8 + examples/codex/brand.yaml | 3 + examples/codex/channels.yaml | 6 + examples/codex/knowledge.yaml | 1 + examples/codex/model.yaml | 4 + examples/codex/skills.yaml | 1 + scripts/start.sh | 77 ++++++-- 17 files changed, 711 insertions(+), 21 deletions(-) create mode 100644 agent-codex/README.md create mode 100644 agent-codex/bun.lock create mode 100644 agent-codex/package.json create mode 100644 agent-codex/src/codex-client.ts create mode 100644 agent-codex/src/history.ts create mode 100644 agent-codex/src/index.ts create mode 100644 agent-codex/tests/history.test.ts create mode 100644 agent-codex/tsconfig.json create mode 100644 examples/codex/agents.yaml create mode 100644 examples/codex/brand.yaml create mode 100644 examples/codex/channels.yaml create mode 100644 examples/codex/knowledge.yaml create mode 100644 examples/codex/model.yaml create mode 100644 examples/codex/skills.yaml diff --git a/.env.example b/.env.example index f46d803a2..dccf77a84 100644 --- a/.env.example +++ b/.env.example @@ -137,8 +137,19 @@ COPILOTKIT_LICENSE_TOKEN= # 0, or leaving this unset, switches the watchdog off. Nothing is watched and no turn is ever ended. AGENT_STALL_TIMEOUT_MS=60000 +# Local Codex compatibility mode. This runs the host-side AG-UI adapter on port 4202, reuses the +# account already authenticated by `codex login`, and skips the two API-key Bot containers. The +# first compatibility implementation is deliberately text-only and read-only; it does not let +# Codex-native actions bypass OpenBot's gateway or audit trail. +# +# CODEX_AGENT_ENABLED=true +# CODEX_AGENT_PORT=4202 +# CODEX_AGENT_WORKSPACE=.openbot-codex/workspace +# AGENT_ENDPOINT_ALLOWED_HOSTS=localhost:4202 + # Model key. Required by the proof-of-concept Bot, which speaks OpenAI's API directly, and by the -# framework Bot unless you point it at another provider below. +# framework Bot unless you enable the local Codex compatibility mode above or point it at another +# provider below. OPENAI_API_KEY= # Where that key is spent. Unset, it is OpenAI. Set, it is any endpoint speaking the same @@ -329,4 +340,3 @@ AGENT_TOOL_TOKEN= # for a deployment that has not stood up a worker. Set for one that has: openssl rand -base64 32. # Do not accept a default in production. WORKER_SHARED_SECRET= - diff --git a/.gitignore b/.gitignore index bfc1237c7..ecb075f3a 100644 --- a/.gitignore +++ b/.gitignore @@ -18,6 +18,7 @@ node_modules/ app/src/lib/generated/application-config.ts .logs/ .demo-logs/ +.openbot-codex/ **/.impeccable .wave-state.md diff --git a/agent-codex/README.md b/agent-codex/README.md new file mode 100644 index 000000000..c0d971738 --- /dev/null +++ b/agent-codex/README.md @@ -0,0 +1,11 @@ +# Codex coworker spike + +This local-only AG-UI adapter lets OpenBot talk to the installed Codex app-server using the +ChatGPT account already authenticated by `codex login`. + +The spike is intentionally text-only and read-only. It does not translate Codex shell, file, MCP, +app, or browser actions into OpenBot tools yet. Those actions are declined so they cannot bypass +OpenBot's gateway and audit trail. + +Enable it with `CODEX_AGENT_ENABLED=true`. `scripts/start.sh` then runs this adapter on +`CODEX_AGENT_PORT` (default `4202`) and skips the two provider-API-key Bot containers. diff --git a/agent-codex/bun.lock b/agent-codex/bun.lock new file mode 100644 index 000000000..44015f926 --- /dev/null +++ b/agent-codex/bun.lock @@ -0,0 +1,40 @@ +{ + "lockfileVersion": 1, + "configVersion": 1, + "workspaces": { + "": { + "name": "@openbot/agent-codex", + "dependencies": { + "@ag-ui/core": "0.0.57", + "@ag-ui/encoder": "0.0.57", + }, + "devDependencies": { + "@types/bun": "^1.3.3", + "typescript": "^5.9.3", + }, + }, + }, + "packages": { + "@ag-ui/core": ["@ag-ui/core@0.0.57", "", { "dependencies": { "zod": "^3.22.4" } }, "sha512-gho1OWjNE6E3Rl7ZEZ1wr2CEpUHjLFU0FqzCZZk439TicLu+BfLCMkMokB07bMGlRmbJ60hM6LW60iOVauCx+Q=="], + + "@ag-ui/encoder": ["@ag-ui/encoder@0.0.57", "", { "dependencies": { "@ag-ui/core": "0.0.57", "@ag-ui/proto": "0.0.57" } }, "sha512-ifD9NctR4xyPDR58xF9GK1bj/S8oECFkTeDfuYD8tXdbcOstIJ2TOqU2zhiCKnw7Vw+zR9Qv3TbsM9E7Gi9X3Q=="], + + "@ag-ui/proto": ["@ag-ui/proto@0.0.57", "", { "dependencies": { "@ag-ui/core": "0.0.57", "@bufbuild/protobuf": "^2.2.5", "@protobuf-ts/protoc": "^2.11.1" } }, "sha512-pPENOZt0P6ibH8sCTgq05wLYXi5t3P9B5r/1bWYehXjUxtyOdnukSlWM++SsCIwUXsQdm/b3aBgGjEeTF7RenA=="], + + "@bufbuild/protobuf": ["@bufbuild/protobuf@2.14.1", "", {}, "sha512-agRJn3+EJDUe8AvxTx/LnHA/GErvLE62pSaSk7+MwFOtOv8eWBu/qCq2qoZjBjVZ3C2aiFJCveuSs17KMkYGOw=="], + + "@protobuf-ts/protoc": ["@protobuf-ts/protoc@2.11.1", "", { "bin": { "protoc": "protoc.js" } }, "sha512-mUZJaV0daGO6HUX90o/atzQ6A7bbN2RSuHtdwo8SSF2Qoe3zHwa4IHyCN1evftTeHfLmdz+45qo47sL+5P8nyg=="], + + "@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@types/node": ["@types/node@26.4.0", "", { "dependencies": { "undici-types": "~8.3.0" } }, "sha512-faiGnoIrLH/V8cibOMEAZ8pMw6oXqSukl29ra4mN8GdaB2ZewzeaLj+INpV5N+Z1eKWzY+IzaIZH2EIR6YZRNQ=="], + + "bun-types": ["bun-types@1.4.0", "", { "dependencies": { "@types/node": "*" } }, "sha512-iIKw23BspnQQYd3prITOBxeUsxBHnwzX6YJfGMuNOZzeNcMmVqzIIVGRm1l69ogaPQmb4wB6BN8mA5bE9YuC5Q=="], + + "typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="], + + "undici-types": ["undici-types@8.3.0", "", {}, "sha512-j375ScV60dom+YkPFIfTLcOiPxkN/buHz5GobjLhixFuANaNs3C9l4GmrWqejgXWJ7BbJcFYpTEUkS1Ge8bpZQ=="], + + "zod": ["zod@3.25.76", "", {}, "sha512-gzUt/qt81nXsFGKIFcC3YnfEAx5NkunCfnDlvuBSSFS02bcXu4Lmea0AFIUwbLWxWPx3d9p8S5QoaujKcNQxcQ=="], + } +} diff --git a/agent-codex/package.json b/agent-codex/package.json new file mode 100644 index 000000000..a896dd9b3 --- /dev/null +++ b/agent-codex/package.json @@ -0,0 +1,21 @@ +{ + "name": "@openbot/agent-codex", + "version": "0.0.1", + "license": "MIT", + "private": true, + "type": "module", + "scripts": { + "start": "bun src/index.ts", + "dev": "bun --watch src/index.ts", + "test": "bun test", + "typecheck": "tsc --noEmit" + }, + "dependencies": { + "@ag-ui/core": "0.0.57", + "@ag-ui/encoder": "0.0.57" + }, + "devDependencies": { + "@types/bun": "^1.3.3", + "typescript": "^5.9.3" + } +} diff --git a/agent-codex/src/codex-client.ts b/agent-codex/src/codex-client.ts new file mode 100644 index 000000000..57de2ef48 --- /dev/null +++ b/agent-codex/src/codex-client.ts @@ -0,0 +1,316 @@ +import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; +import { createInterface } from "node:readline"; + +type JsonObject = Record; +type JsonRpcMessage = { + id?: number | string; + method?: string; + params?: JsonObject; + result?: unknown; + error?: { code?: number; message?: string }; +}; + +type AccountSummary = { + authMode: string; + planType: string | null; +}; + +type PendingRequest = { + resolve: (value: unknown) => void; + reject: (error: Error) => void; + timeout: ReturnType; +}; + +type ThreadStartResult = { + thread: { id: string }; +}; + +type TurnStartResult = { + turn: { id: string }; +}; + +type TurnCallbacks = { + onText(delta: string): void; +}; + +const REQUEST_TIMEOUT_MS = 30_000; +const DEFAULT_TURN_TIMEOUT_MS = 180_000; + +/** Minimal JSON-RPC client for the local Codex app-server stdio transport. */ +export class CodexAppServerClient { + private child: ChildProcessWithoutNullStreams | undefined; + private nextId = 1; + private pending = new Map(); + private listeners = new Set<(message: JsonRpcMessage) => void>(); + private account: AccountSummary | undefined; + + async start(): Promise { + if (this.child) return; + + const binary = process.env.CODEX_BINARY?.trim() || "codex"; + const child = spawn(binary, ["app-server", "--stdio"], { + env: process.env, + stdio: ["pipe", "pipe", "pipe"], + }); + this.child = child; + + const lines = createInterface({ input: child.stdout }); + lines.on("line", (line) => this.receive(line)); + child.stderr.setEncoding("utf8"); + child.stderr.on("data", (chunk: string) => { + const text = chunk.trim(); + if (text) console.error(`[codex app-server] ${text}`); + }); + child.once("error", (error) => this.failAll(error)); + child.once("exit", (code, signal) => { + this.child = undefined; + this.failAll( + new Error( + `Codex app-server exited (${signal ?? `status ${code ?? "unknown"}`}).`, + ), + ); + }); + + await this.request("initialize", { + clientInfo: { + name: "openbot_local_codex", + title: "OpenBot local Codex coworker", + version: "0.0.1", + }, + capabilities: null, + }); + this.notify("initialized", {}); + + const result = (await this.request("account/read", { + refreshToken: false, + })) as { + account?: { type?: string; planType?: string | null } | null; + requiresOpenaiAuth?: boolean; + }; + if (result.account?.type !== "chatgpt") { + throw new Error( + "Codex is not logged in with ChatGPT. Run `codex login` on this Mac first.", + ); + } + this.account = { + authMode: result.account.type, + planType: result.account.planType ?? null, + }; + } + + accountSummary(): AccountSummary { + if (!this.account) { + throw new Error("Codex app-server has not finished starting."); + } + return this.account; + } + + async startThread( + cwd: string, + developerInstructions: string, + ): Promise { + const result = (await this.request("thread/start", { + cwd, + approvalPolicy: "never", + sandbox: "read-only", + serviceName: "openbot_local_codex", + developerInstructions, + ephemeral: false, + })) as ThreadStartResult; + if (!result.thread?.id) { + throw new Error("Codex app-server did not return a thread id."); + } + return result.thread.id; + } + + async runTurn( + threadId: string, + cwd: string, + prompt: string, + callbacks: TurnCallbacks, + ): Promise { + let turnId: string | undefined; + let turnError: string | undefined; + const streamedItems = new Set(); + const timeoutMs = Number.parseInt( + process.env.CODEX_AGENT_TURN_TIMEOUT_MS ?? `${DEFAULT_TURN_TIMEOUT_MS}`, + 10, + ); + + let finish: (() => void) | undefined; + let fail: ((error: Error) => void) | undefined; + const completed = new Promise((resolve, reject) => { + finish = resolve; + fail = reject; + }); + const timeout = setTimeout(() => { + fail?.(new Error(`Codex did not finish within ${timeoutMs}ms.`)); + }, timeoutMs); + + const unsubscribe = this.onMessage((message) => { + const params = message.params ?? {}; + if (params.threadId !== threadId) return; + + if (message.method === "item/agentMessage/delta") { + if (turnId && params.turnId !== turnId) return; + const itemId = typeof params.itemId === "string" ? params.itemId : ""; + const delta = typeof params.delta === "string" ? params.delta : ""; + if (itemId) streamedItems.add(itemId); + if (delta) callbacks.onText(delta); + return; + } + + if (message.method === "item/completed") { + if (turnId && params.turnId !== turnId) return; + const item = params.item as + | { type?: string; id?: string; text?: string } + | undefined; + if ( + item?.type === "agentMessage" && + item.id && + !streamedItems.has(item.id) && + item.text + ) { + callbacks.onText(item.text); + } + return; + } + + if (message.method === "error") { + const error = params.error as { message?: string } | undefined; + turnError = error?.message ?? "Codex reported an unknown error."; + return; + } + + if (message.method === "turn/completed") { + const turn = params.turn as + | { + id?: string; + status?: string; + error?: { message?: string } | null; + } + | undefined; + if (turnId && turn?.id !== turnId) return; + if (turn?.status === "completed") { + finish?.(); + } else { + fail?.( + new Error( + turnError ?? + turn?.error?.message ?? + `Codex turn ended with status ${turn?.status ?? "unknown"}.`, + ), + ); + } + } + }); + + try { + const result = (await this.request("turn/start", { + threadId, + input: [{ type: "text", text: prompt, text_elements: [] }], + cwd, + approvalPolicy: "never", + sandboxPolicy: { type: "readOnly" }, + effort: "low", + })) as TurnStartResult; + turnId = result.turn?.id; + if (!turnId) + throw new Error("Codex app-server did not return a turn id."); + await completed; + } finally { + clearTimeout(timeout); + unsubscribe(); + } + } + + private onMessage(listener: (message: JsonRpcMessage) => void): () => void { + this.listeners.add(listener); + return () => this.listeners.delete(listener); + } + + private request(method: string, params: JsonObject): Promise { + const id = this.nextId++; + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => { + this.pending.delete(id); + reject(new Error(`Codex app-server request ${method} timed out.`)); + }, REQUEST_TIMEOUT_MS); + this.pending.set(id, { resolve, reject, timeout }); + this.write({ method, id, params }); + }); + } + + private notify(method: string, params: JsonObject): void { + this.write({ method, params }); + } + + private write(message: JsonRpcMessage): void { + const child = this.child; + if (!child?.stdin.writable) { + throw new Error("Codex app-server is not running."); + } + child.stdin.write(`${JSON.stringify(message)}\n`); + } + + private receive(line: string): void { + let message: JsonRpcMessage; + try { + message = JSON.parse(line) as JsonRpcMessage; + } catch { + console.error("Codex app-server returned a non-JSON line."); + return; + } + + if (message.id !== undefined && message.method) { + this.answerServerRequest(message); + return; + } + + if (message.id !== undefined) { + const pending = this.pending.get(message.id); + if (!pending) return; + this.pending.delete(message.id); + clearTimeout(pending.timeout); + if (message.error) { + pending.reject( + new Error( + message.error.message ?? "Codex app-server request failed.", + ), + ); + } else { + pending.resolve(message.result); + } + return; + } + + for (const listener of this.listeners) listener(message); + } + + /** The spike never authorizes Codex-native actions; OpenBot tool bridging comes later. */ + private answerServerRequest(message: JsonRpcMessage): void { + if ( + message.method === "item/commandExecution/requestApproval" || + message.method === "item/fileChange/requestApproval" + ) { + this.write({ id: message.id, result: { decision: "decline" } }); + return; + } + this.write({ + id: message.id, + error: { + code: -32601, + message: + "This OpenBot compatibility spike does not expose that action.", + }, + }); + } + + private failAll(error: Error): void { + for (const pending of this.pending.values()) { + clearTimeout(pending.timeout); + pending.reject(error); + } + this.pending.clear(); + } +} diff --git a/agent-codex/src/history.ts b/agent-codex/src/history.ts new file mode 100644 index 000000000..595ad5748 --- /dev/null +++ b/agent-codex/src/history.ts @@ -0,0 +1,45 @@ +import type { RunAgentInput } from "@ag-ui/core"; + +export type CodexTurnInput = { + developerInstructions: string; + prompt: string; +}; + +const SPIKE_INSTRUCTIONS = `You are a Codex coworker inside a local OpenBot compatibility test. +Respond with text only. Do not run shell commands, modify files, browse the web, use MCP servers, +invoke apps, spawn subagents, or call tools. The host intentionally does not expose those actions +during this first compatibility test. Be concise and follow the coworker's standing role.`; + +/** + * Reduce the AG-UI history to the two inputs Codex needs for this turn. + * + * OpenBot sends the full durable transcript on every run. Codex owns its own durable thread once the + * adapter creates it, so replaying that transcript would duplicate every earlier message. The latest + * user message is the new turn; standing system/developer messages become thread instructions. + */ +export function toCodexTurnInput(input: RunAgentInput): CodexTurnInput { + const standingRole = input.messages + .filter( + (message) => message.role === "system" || message.role === "developer", + ) + .map((message) => String(message.content ?? "").trim()) + .filter(Boolean) + .join("\n\n"); + + const latestUser = [...input.messages] + .reverse() + .find((message) => message.role === "user"); + const prompt = String(latestUser?.content ?? "").trim(); + if (!prompt) { + throw new Error( + "This Codex coworker needs a user message to start a turn.", + ); + } + + return { + developerInstructions: standingRole + ? `${SPIKE_INSTRUCTIONS}\n\nStanding role from OpenBot:\n${standingRole}` + : SPIKE_INSTRUCTIONS, + prompt, + }; +} diff --git a/agent-codex/src/index.ts b/agent-codex/src/index.ts new file mode 100644 index 000000000..6dbf00aa4 --- /dev/null +++ b/agent-codex/src/index.ts @@ -0,0 +1,142 @@ +import type { BaseEvent, RunAgentInput } from "@ag-ui/core"; +import { EventEncoder } from "@ag-ui/encoder"; +import { mkdir } from "node:fs/promises"; +import { resolve } from "node:path"; +import { hasManagedAgentToken } from "../../shared/agent-authorisation"; +import { CodexAppServerClient } from "./codex-client"; +import { toCodexTurnInput } from "./history"; + +const PORT = Number.parseInt(process.env.PORT ?? "4202", 10); +const MANAGED_AGENT_TOKEN = process.env.MANAGED_AGENT_TOKEN?.trim(); +if (!MANAGED_AGENT_TOKEN) { + console.error( + "MANAGED_AGENT_TOKEN is not set. The Codex coworker will not start without OpenBot authentication.", + ); + process.exit(1); +} + +const WORKSPACE = resolve( + process.env.CODEX_AGENT_WORKSPACE?.trim() || ".openbot-codex/workspace", +); +await mkdir(WORKSPACE, { recursive: true }); + +const codex = new CodexAppServerClient(); +await codex.start(); + +const codexThreads = new Map(); + +async function runAgent(input: RunAgentInput): Promise { + const encoder = new EventEncoder(); + const stream = new ReadableStream({ + async start(controller) { + const utf8 = new TextEncoder(); + const send = (event: BaseEvent) => + controller.enqueue(utf8.encode(encoder.encodeSSE(event))); + + send({ + type: "RUN_STARTED", + threadId: input.threadId, + runId: input.runId, + } as BaseEvent); + + const messageId = `msg_${input.runId}`; + let textOpen = false; + let textReceived = false; + + try { + const turn = toCodexTurnInput(input); + let codexThreadId = codexThreads.get(input.threadId); + if (!codexThreadId) { + codexThreadId = await codex.startThread( + WORKSPACE, + turn.developerInstructions, + ); + codexThreads.set(input.threadId, codexThreadId); + } + + await codex.runTurn(codexThreadId, WORKSPACE, turn.prompt, { + onText(delta) { + if (!textOpen) { + send({ + type: "TEXT_MESSAGE_START", + messageId, + role: "assistant", + } as BaseEvent); + textOpen = true; + } + textReceived = true; + send({ + type: "TEXT_MESSAGE_CONTENT", + messageId, + delta, + } as BaseEvent); + }, + }); + + if (!textReceived) { + throw new Error("Codex completed without returning a text message."); + } + if (textOpen) { + send({ type: "TEXT_MESSAGE_END", messageId } as BaseEvent); + } + send({ + type: "RUN_FINISHED", + threadId: input.threadId, + runId: input.runId, + } as BaseEvent); + } catch (error) { + if (textOpen) { + send({ type: "TEXT_MESSAGE_END", messageId } as BaseEvent); + } + send({ + type: "RUN_ERROR", + message: + error instanceof Error + ? error.message + : "The Codex coworker could not answer.", + } as BaseEvent); + } finally { + controller.close(); + } + }, + }); + + return new Response(stream, { + headers: { + "content-type": encoder.getContentType(), + "cache-control": "no-cache", + connection: "keep-alive", + }, + }); +} + +Bun.serve({ + port: PORT, + idleTimeout: 255, + async fetch(request) { + const url = new URL(request.url); + if (url.pathname === "/health") { + const account = codex.accountSummary(); + return Response.json({ + status: "ok", + authMode: account.authMode, + planType: account.planType, + safety: "read-only-text-spike", + }); + } + + if (url.pathname === "/ag-ui" && request.method === "POST") { + if (!hasManagedAgentToken(request, MANAGED_AGENT_TOKEN)) { + return Response.json({ error: "Unauthorized." }, { status: 401 }); + } + return runAgent((await request.json()) as RunAgentInput); + } + + return Response.json({ error: "Not found." }, { status: 404 }); + }, +}); + +const account = codex.accountSummary(); +console.info( + `agent-codex listening on http://localhost:${PORT}/ag-ui (${account.authMode}, ${account.planType ?? "unknown plan"})`, +); diff --git a/agent-codex/tests/history.test.ts b/agent-codex/tests/history.test.ts new file mode 100644 index 000000000..193e11a73 --- /dev/null +++ b/agent-codex/tests/history.test.ts @@ -0,0 +1,31 @@ +import { describe, expect, test } from "bun:test"; +import type { RunAgentInput } from "@ag-ui/core"; +import { toCodexTurnInput } from "../src/history"; + +const input = (messages: unknown[]): RunAgentInput => + ({ messages }) as RunAgentInput; + +describe("toCodexTurnInput", () => { + test("uses the latest user message and carries the standing role", () => { + const result = toCodexTurnInput( + input([ + { role: "system", content: "You are the finance coworker." }, + { role: "user", content: "First question" }, + { role: "assistant", content: "First answer" }, + { role: "user", content: "Follow-up question" }, + ]), + ); + + expect(result.prompt).toBe("Follow-up question"); + expect(result.developerInstructions).toContain( + "You are the finance coworker.", + ); + expect(result.developerInstructions).toContain("Respond with text only"); + }); + + test("refuses an empty turn", () => { + expect(() => + toCodexTurnInput(input([{ role: "system", content: "A role" }])), + ).toThrow("needs a user message"); + }); +}); diff --git a/agent-codex/tsconfig.json b/agent-codex/tsconfig.json new file mode 100644 index 000000000..edc4271f9 --- /dev/null +++ b/agent-codex/tsconfig.json @@ -0,0 +1,11 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "ESNext", + "moduleResolution": "Bundler", + "strict": true, + "noEmit": true, + "types": ["bun"] + }, + "include": ["src/**/*.ts", "tests/**/*.ts"] +} diff --git a/examples/codex/agents.yaml b/examples/codex/agents.yaml new file mode 100644 index 000000000..98ee96793 --- /dev/null +++ b/examples/codex/agents.yaml @@ -0,0 +1,8 @@ +agents: + - id: codex-assistant + name: Codex + title: Subscription-powered coworker + role_description: Help with general questions using the locally signed-in Codex account. + avatar_seed: codex-assistant + type: remote-ag-ui + endpoint: ${MANAGED_AGENT_AG_UI_URL:-} diff --git a/examples/codex/brand.yaml b/examples/codex/brand.yaml new file mode 100644 index 000000000..5313ee0ae --- /dev/null +++ b/examples/codex/brand.yaml @@ -0,0 +1,3 @@ +tenant: + id: openbot-codex-local + product_name: OpenBot + Codex diff --git a/examples/codex/channels.yaml b/examples/codex/channels.yaml new file mode 100644 index 000000000..a0cd305c9 --- /dev/null +++ b/examples/codex/channels.yaml @@ -0,0 +1,6 @@ +channels: + - id: codex + name: Codex + description: Test a Codex coworker backed by the local ChatGPT subscription. + permitted_agents: [codex-assistant] + allowed_groups: [all] diff --git a/examples/codex/knowledge.yaml b/examples/codex/knowledge.yaml new file mode 100644 index 000000000..b87707414 --- /dev/null +++ b/examples/codex/knowledge.yaml @@ -0,0 +1 @@ +sources: [] diff --git a/examples/codex/model.yaml b/examples/codex/model.yaml new file mode 100644 index 000000000..ba7c60b7e --- /dev/null +++ b/examples/codex/model.yaml @@ -0,0 +1,4 @@ +model: + provider: openai + credential_secret_ref: unused-in-codex-spike + default_model: gpt-5.6-terra diff --git a/examples/codex/skills.yaml b/examples/codex/skills.yaml new file mode 100644 index 000000000..ada13e6ed --- /dev/null +++ b/examples/codex/skills.yaml @@ -0,0 +1 @@ +skills: [] diff --git a/scripts/start.sh b/scripts/start.sh index 74beebab6..0c26da181 100755 --- a/scripts/start.sh +++ b/scripts/start.sh @@ -30,8 +30,17 @@ SERVER_PORT="$(setting SERVER_PORT 3001)" COMPUTER_PORT="$(setting COMPUTER_PORT 4100)" BOT_PORT="$(setting BOT_PORT 4200)" LANGGRAPH_PORT="$(setting LANGGRAPH_PORT 4201)" +CODEX_AGENT_PORT="$(setting CODEX_AGENT_PORT 4202)" SUPERVISOR_PORT="$(setting SUPERVISOR_PORT 4500)" ONE_COMPUTER_EACH="${OPENBOT_ONE_COMPUTER_EACH:-true}" +CODEX_AGENT_ENABLED="$(setting CODEX_AGENT_ENABLED false)" +case "$CODEX_AGENT_ENABLED" in + true|false) ;; + *) + printf '\033[31m%s\033[0m\n' "CODEX_AGENT_ENABLED must be true or false." + exit 1 + ;; +esac export APP_PORT SERVER_PORT SUPERVISOR_TOKEN="$(setting SUPERVISOR_TOKEN openbot-dev-supervisor-token)" COMPUTER_TOKEN="$(setting COMPUTER_TOKEN openbot-dev-computer-token)" @@ -54,10 +63,15 @@ WORKER_SHARED_SECRET="$(setting WORKER_SHARED_SECRET openbot-dev-worker-secret)" # # Written into .env rather than exported for this run alone, so `docker compose up` by hand later # sees the same value the script used. -# The laptop stack runs agent-langgraph on LANGGRAPH_PORT. The one-container image does not, so -# this default stays in the script rather than in .env: a `docker run --env-file .env` must not -# inherit a URL that points at a process the image does not contain. -MANAGED_AGENT_AG_UI_URL="$(setting MANAGED_AGENT_AG_UI_URL "http://localhost:${LANGGRAPH_PORT}/ag-ui")" +# The laptop stack normally runs agent-langgraph on LANGGRAPH_PORT. The local Codex compatibility +# mode instead runs a host process so it can reuse the person's existing `codex login` session +# without copying ChatGPT credentials into Docker. +if [ "$CODEX_AGENT_ENABLED" = "true" ]; then + DEFAULT_MANAGED_AGENT_URL="http://localhost:${CODEX_AGENT_PORT}/ag-ui" +else + DEFAULT_MANAGED_AGENT_URL="http://localhost:${LANGGRAPH_PORT}/ag-ui" +fi +MANAGED_AGENT_AG_UI_URL="$(setting MANAGED_AGENT_AG_UI_URL "$DEFAULT_MANAGED_AGENT_URL")" export MANAGED_AGENT_AG_UI_URL # Whether this run minted a secret that something already running may not have. @@ -215,12 +229,13 @@ fi # `docker compose up -d` is declarative and does nothing for a service whose configuration has not # changed, so naming them all costs a comparison and buys the guarantee that what is running is what # this run configured. -for svc in agent-computer agent-bot agent-langgraph; do - SERVICES+=("$svc") -done +SERVICES+=(agent-computer) +if [ "$CODEX_AGENT_ENABLED" != "true" ]; then + SERVICES+=(agent-bot agent-langgraph) +fi export SUPERVISOR_TOKEN COMPUTER_TOKEN WORKER_SHARED_SECRET -export COMPUTER_PORT BOT_PORT LANGGRAPH_PORT SUPERVISOR_PORT +export COMPUTER_PORT BOT_PORT LANGGRAPH_PORT CODEX_AGENT_PORT SUPERVISOR_PORT docker compose up -d --build "${SERVICES[@]}" >/dev/null if ! docker compose run --rm --build migrate >"$LOGS/migrate.log" 2>&1; then red " Migrations did not apply. The database is not the schema this server expects." @@ -228,8 +243,25 @@ if ! docker compose run --rm --build migrate >"$LOGS/migrate.log" 2>&1; then exit 1 fi wait_for "http://localhost:$COMPUTER_PORT/health" "agent-computer" -wait_for "http://localhost:$BOT_PORT/health" "agent-bot" -wait_for "http://localhost:$LANGGRAPH_PORT/health" "agent-langgraph" +if [ "$CODEX_AGENT_ENABLED" = "true" ]; then + CODEX_AGENT_WORKSPACE="$(setting CODEX_AGENT_WORKSPACE "$ROOT/.openbot-codex/workspace")" + mkdir -p "$CODEX_AGENT_WORKSPACE" + require_free_or_ours "$CODEX_AGENT_PORT" "agent-codex" + if identifies_as_openbot "$CODEX_AGENT_PORT" "agent-codex"; then + pkill -f "bun agent-codex/src/index.ts" >/dev/null 2>&1 || true + sleep 1 + fi + (cd "$ROOT" && \ + nohup env \ + PORT="$CODEX_AGENT_PORT" \ + MANAGED_AGENT_TOKEN="$MANAGED_AGENT_TOKEN" \ + CODEX_AGENT_WORKSPACE="$CODEX_AGENT_WORKSPACE" \ + bun agent-codex/src/index.ts >"$LOGS/agent-codex.log" 2>&1 "$LOGS/server.log" 2>&1 &) + bun --env-file=../.env src/index.ts >"$LOGS/server.log" 2>&1 "$LOGS/server.log" 2>&1 &) + bun --env-file=../.env src/index.ts >"$LOGS/server.log" 2>&1 /dev/null 2>&1; then WORKER_DATABASE_URL="$(setting DATABASE_URL postgres://openbot:openbot@localhost:5432/openbot)" (cd "$ROOT" && \ - DATABASE_URL="$WORKER_DATABASE_URL" \ - SERVER_INTERNAL_URL="http://localhost:$SERVER_PORT" \ - WORKER_SHARED_SECRET="$WORKER_SHARED_SECRET" \ - bun worker/src/index.ts >"$LOGS/worker.log" 2>&1 &) + nohup env \ + DATABASE_URL="$WORKER_DATABASE_URL" \ + SERVER_INTERNAL_URL="http://localhost:$SERVER_PORT" \ + WORKER_SHARED_SECRET="$WORKER_SHARED_SECRET" \ + bun worker/src/index.ts >"$LOGS/worker.log" 2>&1 /dev/null 2>&1; then @@ -368,7 +406,7 @@ PY info "4/4 App" require_free_or_ours "$APP_PORT" app if ! identifies_as_openbot "$APP_PORT" app; then - (cd app && bun run dev --port "$APP_PORT" --strictPort >"$LOGS/app.log" 2>&1 &) + (cd app && nohup bun run dev --port "$APP_PORT" --strictPort >"$LOGS/app.log" 2>&1 Date: Tue, 1 Sep 2026 13:55:05 -0400 Subject: [PATCH 2/2] feat: recover Codex threads through governed tools --- .env.example | 10 +- CHANGELOG.md | 13 + agent-codex/README.md | 24 +- agent-codex/package.json | 2 +- agent-codex/src/codex-client.ts | 396 +++++++++++++++++++++---- agent-codex/src/history.ts | 15 +- agent-codex/src/index.ts | 184 +++++++++--- agent-codex/src/thread-store.ts | 116 ++++++++ agent-codex/src/tools.ts | 145 +++++++++ agent-codex/tests/codex-client.test.ts | 315 ++++++++++++++++++++ agent-codex/tests/history.test.ts | 3 +- agent-codex/tests/thread-store.test.ts | 70 +++++ agent-codex/tests/tools.test.ts | 142 +++++++++ bun.lock | 14 + package.json | 1 + scripts/start.sh | 11 +- tests/workspace.test.ts | 13 +- 17 files changed, 1358 insertions(+), 116 deletions(-) create mode 100644 agent-codex/src/thread-store.ts create mode 100644 agent-codex/src/tools.ts create mode 100644 agent-codex/tests/codex-client.test.ts create mode 100644 agent-codex/tests/thread-store.test.ts create mode 100644 agent-codex/tests/tools.test.ts diff --git a/.env.example b/.env.example index dccf77a84..bce08ab0d 100644 --- a/.env.example +++ b/.env.example @@ -138,13 +138,17 @@ COPILOTKIT_LICENSE_TOKEN= AGENT_STALL_TIMEOUT_MS=60000 # Local Codex compatibility mode. This runs the host-side AG-UI adapter on port 4202, reuses the -# account already authenticated by `codex login`, and skips the two API-key Bot containers. The -# first compatibility implementation is deliberately text-only and read-only; it does not let -# Codex-native actions bypass OpenBot's gateway or audit trail. +# account already authenticated by `codex login`, and skips the two API-key Bot containers. Its +# Codex threads survive adapter restarts. Every side-effecting tool call returns through OpenBot's +# signed gateway, where the current grant, policy and audit trail are applied. MCP, app, plugin and +# web capabilities are disabled; turns are sandboxed read-only without network access, and the +# adapter interrupts any native shell or file action Codex still attempts. # # CODEX_AGENT_ENABLED=true # CODEX_AGENT_PORT=4202 # CODEX_AGENT_WORKSPACE=.openbot-codex/workspace +# CODEX_AGENT_STATE=.openbot-codex/threads.json +# OPENBOT_TOOL_URL=http://localhost:3001/api/agent-tools/call # AGENT_ENDPOINT_ALLOWED_HOSTS=localhost:4202 # Model key. Required by the proof-of-concept Bot, which speaks OpenAI's API directly, and by the diff --git a/CHANGELOG.md b/CHANGELOG.md index efa50060a..7276df747 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,19 @@ Newest first. `Unreleased` is what is on `main` and not yet tagged. ## Unreleased +### A local Codex coworker keeps its conversation and uses OpenBot's governed tools + +OpenBot can now run a coworker through the Codex app already signed in on the host, without an API +key. Its conversation survives adapter and app-server restarts: OpenBot records the Codex thread it +owns, resumes that exact thread, and refuses to silently replace unreadable recovery state. + +Connector calls do not go around the deployment. Codex sees only the tools OpenBot assigned to that +coworker, and every call returns with the deployment's signed run assertion through the existing +grant, policy and audit gateway. Native Codex shell, file, web, app, plugin and MCP paths are disabled +or interrupted for this mode, and the remaining turn is sandboxed read-only without network access. +Rejected, stale and duplicate callbacks are reported as failed tool results rather than being +carried out. + ### Coworkers are made in a wizard and managed in a dialog Creating a coworker is now a three-step wizard — who it is, who may see it, then where it runs, diff --git a/agent-codex/README.md b/agent-codex/README.md index c0d971738..458457d92 100644 --- a/agent-codex/README.md +++ b/agent-codex/README.md @@ -1,11 +1,21 @@ -# Codex coworker spike +# Codex coworker -This local-only AG-UI adapter lets OpenBot talk to the installed Codex app-server using the -ChatGPT account already authenticated by `codex login`. +This local-only AG-UI adapter lets OpenBot talk to the installed Codex app-server using the ChatGPT +account already authenticated by `codex login`. -The spike is intentionally text-only and read-only. It does not translate Codex shell, file, MCP, -app, or browser actions into OpenBot tools yet. Those actions are declined so they cannot bypass -OpenBot's gateway and audit trail. +The adapter records the join between each OpenBot Intelligence thread and its persistent Codex +thread in `.openbot-codex/threads.json`. On every later run—including after the adapter restarts—it +resumes that Codex thread and refreshes its standing instructions. Codex restores the dynamic-tool +catalog persisted with the thread; OpenBot still rechecks the current grant and policy on every call. + +Only tools that OpenBot marks as deployment-owned are exposed as Codex dynamic tools. Calls return +to `/api/agent-tools/call` with OpenBot's signed run assertion and agent token, so the deployment +rechecks the Bot's grant and policy and writes the normal audit events. Codex-native shell, file, +MCP, app, web and multi-agent paths are disabled. The adapter interrupts native shell and file +attempts; turns run in the read-only sandbox without network access, so an action that races the +interrupt cannot write or reach the network. Enable it with `CODEX_AGENT_ENABLED=true`. `scripts/start.sh` then runs this adapter on -`CODEX_AGENT_PORT` (default `4202`) and skips the two provider-API-key Bot containers. +`CODEX_AGENT_PORT` (default `4202`) and skips the two provider-API-key Bot containers. The start +script supplies `AGENT_TOOL_TOKEN`, `OPENBOT_TOOL_URL` and `CODEX_AGENT_STATE`; set those explicitly +when starting `agent-codex` by hand. diff --git a/agent-codex/package.json b/agent-codex/package.json index a896dd9b3..faaf64f45 100644 --- a/agent-codex/package.json +++ b/agent-codex/package.json @@ -1,6 +1,6 @@ { "name": "@openbot/agent-codex", - "version": "0.0.1", + "version": "0.0.2", "license": "MIT", "private": true, "type": "module", diff --git a/agent-codex/src/codex-client.ts b/agent-codex/src/codex-client.ts index 57de2ef48..1fb5cc351 100644 --- a/agent-codex/src/codex-client.ts +++ b/agent-codex/src/codex-client.ts @@ -1,5 +1,6 @@ import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; import { createInterface } from "node:readline"; +import type { CodexDynamicTool, ToolResult } from "./tools"; type JsonObject = Record; type JsonRpcMessage = { @@ -21,7 +22,7 @@ type PendingRequest = { timeout: ReturnType; }; -type ThreadStartResult = { +type ThreadResult = { thread: { id: string }; }; @@ -29,29 +30,53 @@ type TurnStartResult = { turn: { id: string }; }; -type TurnCallbacks = { +export type TurnCallbacks = { onText(delta: string): void; + onToolCall( + callId: string, + name: string, + args: Record, + ): Promise; }; +type ActiveTurn = { + turnId?: string; + callbacks: TurnCallbacks; + toolCallIds: Set; + fail(error: Error): void; +}; + +type SpawnAppServer = () => ChildProcessWithoutNullStreams; + const REQUEST_TIMEOUT_MS = 30_000; const DEFAULT_TURN_TIMEOUT_MS = 180_000; - -/** Minimal JSON-RPC client for the local Codex app-server stdio transport. */ +const BLOCKED_ITEM_TYPES = new Set([ + "commandExecution", + "fileChange", + "mcpToolCall", + "collabAgentToolCall", + "subAgentActivity", + "webSearch", + "imageView", + "imageGeneration", +]); + +/** JSON-RPC client for a local Codex app-server with OpenBot as its only tool boundary. */ export class CodexAppServerClient { private child: ChildProcessWithoutNullStreams | undefined; private nextId = 1; private pending = new Map(); private listeners = new Set<(message: JsonRpcMessage) => void>(); + private activeTurns = new Map(); private account: AccountSummary | undefined; + private safetyConfig: JsonObject = safetyConfigFor([]); + + constructor(private readonly spawnAppServer: SpawnAppServer = launchCodex) {} async start(): Promise { if (this.child) return; - const binary = process.env.CODEX_BINARY?.trim() || "codex"; - const child = spawn(binary, ["app-server", "--stdio"], { - env: process.env, - stdio: ["pipe", "pipe", "pipe"], - }); + const child = this.spawnAppServer(); this.child = child; const lines = createInterface({ input: child.stdout }); @@ -75,9 +100,9 @@ export class CodexAppServerClient { clientInfo: { name: "openbot_local_codex", title: "OpenBot local Codex coworker", - version: "0.0.1", + version: "0.0.2", }, - capabilities: null, + capabilities: { experimentalApi: true }, }); this.notify("initialized", {}); @@ -85,7 +110,6 @@ export class CodexAppServerClient { refreshToken: false, })) as { account?: { type?: string; planType?: string | null } | null; - requiresOpenaiAuth?: boolean; }; if (result.account?.type !== "chatgpt") { throw new Error( @@ -96,6 +120,15 @@ export class CodexAppServerClient { authMode: result.account.type, planType: result.account.planType ?? null, }; + + const configResult = (await this.request("config/read", { + includeLayers: false, + })) as { config?: { mcp_servers?: unknown } }; + this.safetyConfig = safetyConfigFor( + isObject(configResult.config?.mcp_servers) + ? Object.keys(configResult.config.mcp_servers) + : [], + ); } accountSummary(): AccountSummary { @@ -108,6 +141,7 @@ export class CodexAppServerClient { async startThread( cwd: string, developerInstructions: string, + dynamicTools: CodexDynamicTool[], ): Promise { const result = (await this.request("thread/start", { cwd, @@ -115,44 +149,120 @@ export class CodexAppServerClient { sandbox: "read-only", serviceName: "openbot_local_codex", developerInstructions, + dynamicTools, + config: this.safetyConfig, ephemeral: false, - })) as ThreadStartResult; + })) as ThreadResult; if (!result.thread?.id) { throw new Error("Codex app-server did not return a thread id."); } return result.thread.id; } + async resumeThread( + threadId: string, + cwd: string, + developerInstructions: string, + ): Promise { + const result = (await this.request("thread/resume", { + threadId, + cwd, + approvalPolicy: "never", + sandbox: "read-only", + developerInstructions, + config: this.safetyConfig, + })) as ThreadResult; + if (result.thread?.id !== threadId) { + throw new Error( + `Codex resumed ${result.thread?.id ?? "no thread"} instead of ${threadId}.`, + ); + } + } + async runTurn( threadId: string, cwd: string, prompt: string, callbacks: TurnCallbacks, + signal?: AbortSignal, ): Promise { - let turnId: string | undefined; + if (this.activeTurns.has(threadId)) { + throw new Error(`Codex thread ${threadId} already has an active turn.`); + } + let turnError: string | undefined; const streamedItems = new Set(); - const timeoutMs = Number.parseInt( - process.env.CODEX_AGENT_TURN_TIMEOUT_MS ?? `${DEFAULT_TURN_TIMEOUT_MS}`, - 10, - ); - - let finish: (() => void) | undefined; - let fail: ((error: Error) => void) | undefined; + const timeoutMs = turnTimeoutMs(); + let settled = false; + let resolveCompletion: (() => void) | undefined; + let rejectCompletion: ((error: Error) => void) | undefined; const completed = new Promise((resolve, reject) => { - finish = resolve; - fail = reject; + resolveCompletion = resolve; + rejectCompletion = reject; }); + const finish = () => { + if (settled) return; + settled = true; + resolveCompletion?.(); + }; + const fail = (error: Error) => { + if (settled) return; + settled = true; + rejectCompletion?.(error); + }; + const active: ActiveTurn = { callbacks, toolCallIds: new Set(), fail }; + this.activeTurns.set(threadId, active); + + const interrupt = () => { + if (!active.turnId) return; + void this.request("turn/interrupt", { + threadId, + turnId: active.turnId, + }).catch(() => {}); + }; + const abort = () => { + interrupt(); + fail( + new Error("OpenBot ended the request before the Codex turn finished."), + ); + }; + if (signal?.aborted) abort(); + else signal?.addEventListener("abort", abort, { once: true }); + const timeout = setTimeout(() => { - fail?.(new Error(`Codex did not finish within ${timeoutMs}ms.`)); + interrupt(); + fail(new Error(`Codex did not finish within ${timeoutMs}ms.`)); }, timeoutMs); const unsubscribe = this.onMessage((message) => { const params = message.params ?? {}; if (params.threadId !== threadId) return; + if ( + active.turnId && + typeof params.turnId === "string" && + params.turnId !== active.turnId + ) + return; + if (!active.turnId && typeof params.turnId === "string") { + active.turnId = params.turnId; + } + + if (message.method === "item/started") { + const item = isObject(params.item) ? params.item : {}; + if ( + typeof item.type === "string" && + BLOCKED_ITEM_TYPES.has(item.type) + ) { + const error = new Error( + `Codex attempted the native ${item.type} path. OpenBot refused it because side effects must use a governed OpenBot tool.`, + ); + interrupt(); + fail(error); + } + return; + } if (message.method === "item/agentMessage/delta") { - if (turnId && params.turnId !== turnId) return; const itemId = typeof params.itemId === "string" ? params.itemId : ""; const delta = typeof params.delta === "string" ? params.delta : ""; if (itemId) streamedItems.add(itemId); @@ -161,14 +271,12 @@ export class CodexAppServerClient { } if (message.method === "item/completed") { - if (turnId && params.turnId !== turnId) return; - const item = params.item as - | { type?: string; id?: string; text?: string } - | undefined; + const item = isObject(params.item) ? params.item : {}; if ( - item?.type === "agentMessage" && - item.id && + item.type === "agentMessage" && + typeof item.id === "string" && !streamedItems.has(item.id) && + typeof item.text === "string" && item.text ) { callbacks.onText(item.text); @@ -177,28 +285,33 @@ export class CodexAppServerClient { } if (message.method === "error") { - const error = params.error as { message?: string } | undefined; - turnError = error?.message ?? "Codex reported an unknown error."; + const error = isObject(params.error) ? params.error : {}; + turnError = + typeof error.message === "string" + ? error.message + : "Codex reported an unknown error."; return; } if (message.method === "turn/completed") { - const turn = params.turn as - | { - id?: string; - status?: string; - error?: { message?: string } | null; - } - | undefined; - if (turnId && turn?.id !== turnId) return; - if (turn?.status === "completed") { - finish?.(); + const turn = isObject(params.turn) ? params.turn : {}; + if ( + active.turnId && + typeof turn.id === "string" && + turn.id !== active.turnId + ) + return; + if (turn.status === "completed") { + finish(); } else { - fail?.( + const error = isObject(turn.error) ? turn.error : {}; + fail( new Error( turnError ?? - turn?.error?.message ?? - `Codex turn ended with status ${turn?.status ?? "unknown"}.`, + (typeof error.message === "string" + ? error.message + : undefined) ?? + `Codex turn ended with status ${String(turn.status ?? "unknown")}.`, ), ); } @@ -211,19 +324,34 @@ export class CodexAppServerClient { input: [{ type: "text", text: prompt, text_elements: [] }], cwd, approvalPolicy: "never", - sandboxPolicy: { type: "readOnly" }, + sandboxPolicy: { type: "readOnly", networkAccess: false }, effort: "low", })) as TurnStartResult; - turnId = result.turn?.id; - if (!turnId) + const startedTurnId = result.turn?.id; + if (!startedTurnId) { throw new Error("Codex app-server did not return a turn id."); + } + if (active.turnId && active.turnId !== startedTurnId) { + throw new Error( + `Codex sent events for turn ${active.turnId} before starting ${startedTurnId}.`, + ); + } + active.turnId = startedTurnId; + if (signal?.aborted) interrupt(); await completed; } finally { clearTimeout(timeout); + signal?.removeEventListener("abort", abort); unsubscribe(); + this.activeTurns.delete(threadId); } } + stop(): void { + this.child?.kill(); + this.child = undefined; + } + private onMessage(listener: (message: JsonRpcMessage) => void): () => void { this.listeners.add(listener); return () => this.listeners.delete(listener); @@ -263,7 +391,7 @@ export class CodexAppServerClient { } if (message.id !== undefined && message.method) { - this.answerServerRequest(message); + void this.answerServerRequest(message); return; } @@ -287,8 +415,7 @@ export class CodexAppServerClient { for (const listener of this.listeners) listener(message); } - /** The spike never authorizes Codex-native actions; OpenBot tool bridging comes later. */ - private answerServerRequest(message: JsonRpcMessage): void { + private async answerServerRequest(message: JsonRpcMessage): Promise { if ( message.method === "item/commandExecution/requestApproval" || message.method === "item/fileChange/requestApproval" @@ -296,12 +423,121 @@ export class CodexAppServerClient { this.write({ id: message.id, result: { decision: "decline" } }); return; } + + if (message.method === "item/permissions/requestApproval") { + this.write({ + id: message.id, + error: { + code: -32000, + message: + "OpenBot denied the requested native permission. Use an OpenBot dynamic tool instead.", + }, + }); + return; + } + + if (message.method === "item/tool/call") { + const params = message.params ?? {}; + const threadId = + typeof params.threadId === "string" ? params.threadId : ""; + const active = this.activeTurns.get(threadId); + const turnId = typeof params.turnId === "string" ? params.turnId : ""; + const callId = typeof params.callId === "string" ? params.callId : ""; + const name = typeof params.tool === "string" ? params.tool : ""; + if ( + !active || + !turnId || + (active.turnId !== undefined && turnId !== active.turnId) || + !callId || + !name || + params.namespace !== null + ) { + this.write({ + id: message.id, + result: { + contentItems: [ + { + type: "inputText", + text: "OpenBot refused this tool call because it does not belong to the active turn.", + }, + ], + success: false, + }, + }); + return; + } + active.turnId = turnId; + + if (!isObject(params.arguments)) { + this.write({ + id: message.id, + result: { + contentItems: [ + { + type: "inputText", + text: "OpenBot refused this tool call because its arguments were not a JSON object.", + }, + ], + success: false, + }, + }); + return; + } + + if (active.toolCallIds.has(callId)) { + this.write({ + id: message.id, + result: { + contentItems: [ + { + type: "inputText", + text: "OpenBot refused a duplicate tool call id so the action could not run twice.", + }, + ], + success: false, + }, + }); + return; + } + active.toolCallIds.add(callId); + + try { + const result = await active.callbacks.onToolCall( + callId, + name, + params.arguments, + ); + this.write({ + id: message.id, + result: { + contentItems: [{ type: "inputText", text: result.text }], + success: result.success, + }, + }); + } catch (error) { + this.write({ + id: message.id, + result: { + contentItems: [ + { + type: "inputText", + text: `OpenBot's governed tool callback failed: ${ + error instanceof Error ? error.message : "unknown error" + }`, + }, + ], + success: false, + }, + }); + } + return; + } + this.write({ id: message.id, error: { code: -32601, - message: - "This OpenBot compatibility spike does not expose that action.", + message: "OpenBot does not expose that Codex-native action.", }, }); } @@ -312,5 +548,57 @@ export class CodexAppServerClient { pending.reject(error); } this.pending.clear(); + for (const turn of this.activeTurns.values()) turn.fail(error); } } + +function launchCodex(): ChildProcessWithoutNullStreams { + const binary = process.env.CODEX_BINARY?.trim() || "codex"; + return spawn(binary, ["app-server", "--stdio"], { + env: process.env, + stdio: ["pipe", "pipe", "pipe"], + }); +} + +function safetyConfigFor(mcpServerNames: string[]): JsonObject { + return { + mcp_servers: Object.fromEntries( + mcpServerNames.map((name) => [name, { enabled: false }]), + ), + features: { + apps: false, + plugins: false, + multi_agent: false, + hooks: false, + memories: false, + goals: false, + code_mode: { enabled: false }, + }, + web_search: "disabled", + apps: { + _default: { + enabled: false, + destructive_enabled: false, + open_world_enabled: false, + }, + }, + tools: { + web_search: false, + view_image: false, + }, + }; +} + +function turnTimeoutMs(): number { + const configured = Number.parseInt( + process.env.CODEX_AGENT_TURN_TIMEOUT_MS ?? `${DEFAULT_TURN_TIMEOUT_MS}`, + 10, + ); + return Number.isFinite(configured) && configured > 0 + ? configured + : DEFAULT_TURN_TIMEOUT_MS; +} + +function isObject(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} diff --git a/agent-codex/src/history.ts b/agent-codex/src/history.ts index 595ad5748..65f4bf11d 100644 --- a/agent-codex/src/history.ts +++ b/agent-codex/src/history.ts @@ -5,10 +5,13 @@ export type CodexTurnInput = { prompt: string; }; -const SPIKE_INSTRUCTIONS = `You are a Codex coworker inside a local OpenBot compatibility test. -Respond with text only. Do not run shell commands, modify files, browse the web, use MCP servers, -invoke apps, spawn subagents, or call tools. The host intentionally does not expose those actions -during this first compatibility test. Be concise and follow the coworker's standing role.`; +const OPENBOT_INSTRUCTIONS = `You are a Codex coworker inside OpenBot. +You may call the OpenBot dynamic tools provided for this thread. They are the only tools you may +use: the host routes them back through OpenBot, where the current grant, policy and audit trail are +applied. Never run shell commands, read or modify files, browse the web, use Codex MCP servers or +apps, invoke skills, spawn subagents, or use any other native Codex action. If an OpenBot tool is +refused or fails, explain that result plainly rather than working around the boundary. Be concise +and follow the coworker's standing role.`; /** * Reduce the AG-UI history to the two inputs Codex needs for this turn. @@ -38,8 +41,8 @@ export function toCodexTurnInput(input: RunAgentInput): CodexTurnInput { return { developerInstructions: standingRole - ? `${SPIKE_INSTRUCTIONS}\n\nStanding role from OpenBot:\n${standingRole}` - : SPIKE_INSTRUCTIONS, + ? `${OPENBOT_INSTRUCTIONS}\n\nStanding role from OpenBot:\n${standingRole}` + : OPENBOT_INSTRUCTIONS, prompt, }; } diff --git a/agent-codex/src/index.ts b/agent-codex/src/index.ts index 6dbf00aa4..9f332bc1b 100644 --- a/agent-codex/src/index.ts +++ b/agent-codex/src/index.ts @@ -5,6 +5,8 @@ import { resolve } from "node:path"; import { hasManagedAgentToken } from "../../shared/agent-authorisation"; import { CodexAppServerClient } from "./codex-client"; import { toCodexTurnInput } from "./history"; +import { CodexThreadStore } from "./thread-store"; +import { dynamicToolsOf, OpenBotToolGateway, runAssertionOf } from "./tools"; const PORT = Number.parseInt(process.env.PORT ?? "4202", 10); const MANAGED_AGENT_TOKEN = process.env.MANAGED_AGENT_TOKEN?.trim(); @@ -15,17 +17,36 @@ if (!MANAGED_AGENT_TOKEN) { process.exit(1); } +const TOOL_TOKEN = process.env.AGENT_TOOL_TOKEN?.trim(); +if (!TOOL_TOKEN) { + console.error( + "AGENT_TOOL_TOKEN is not set. The Codex coworker only runs tools through OpenBot's governance gateway.", + ); + process.exit(1); +} + const WORKSPACE = resolve( process.env.CODEX_AGENT_WORKSPACE?.trim() || ".openbot-codex/workspace", ); +const STATE_PATH = resolve( + process.env.CODEX_AGENT_STATE?.trim() || ".openbot-codex/threads.json", +); +const TOOL_URL = + process.env.OPENBOT_TOOL_URL?.trim() || + "http://localhost:3001/api/agent-tools/call"; await mkdir(WORKSPACE, { recursive: true }); +const threadStore = await CodexThreadStore.open(STATE_PATH); +const gateway = new OpenBotToolGateway({ url: TOOL_URL, token: TOOL_TOKEN }); const codex = new CodexAppServerClient(); await codex.start(); -const codexThreads = new Map(); +const threadQueues = new Map>(); -async function runAgent(input: RunAgentInput): Promise { +async function runAgent( + input: RunAgentInput, + requestSignal: AbortSignal, +): Promise { const encoder = new EventEncoder(); const stream = new ReadableStream({ async start(controller) { @@ -39,45 +60,107 @@ async function runAgent(input: RunAgentInput): Promise { runId: input.runId, } as BaseEvent); - const messageId = `msg_${input.runId}`; + let messageSequence = 0; + let messageId = ""; let textOpen = false; - let textReceived = false; + let visibleReceived = false; + const closeText = () => { + if (!textOpen) return; + send({ type: "TEXT_MESSAGE_END", messageId } as BaseEvent); + textOpen = false; + }; try { - const turn = toCodexTurnInput(input); - let codexThreadId = codexThreads.get(input.threadId); - if (!codexThreadId) { - codexThreadId = await codex.startThread( - WORKSPACE, - turn.developerInstructions, + await serialiseThread(input.threadId, async () => { + const turn = toCodexTurnInput(input); + const dynamicTools = dynamicToolsOf(input); + const allowedToolNames = new Set( + dynamicTools.map((tool) => tool.name), ); - codexThreads.set(input.threadId, codexThreadId); - } + const runAssertion = runAssertionOf(input); + let codexThreadId = threadStore.get(input.threadId); + if (codexThreadId) { + await codex.resumeThread( + codexThreadId, + WORKSPACE, + turn.developerInstructions, + ); + } else { + codexThreadId = await codex.startThread( + WORKSPACE, + turn.developerInstructions, + dynamicTools, + ); + // Persist before the first turn. A crash can orphan an empty Codex thread, but it can + // never produce conversation state that OpenBot subsequently forgets how to resume. + await threadStore.remember(input.threadId, codexThreadId); + } + + await codex.runTurn( + codexThreadId, + WORKSPACE, + turn.prompt, + { + onText(delta) { + if (!textOpen) { + messageId = `msg_${input.runId}_${messageSequence++}`; + send({ + type: "TEXT_MESSAGE_START", + messageId, + role: "assistant", + } as BaseEvent); + textOpen = true; + } + visibleReceived = true; + send({ + type: "TEXT_MESSAGE_CONTENT", + messageId, + delta, + } as BaseEvent); + }, + async onToolCall(callId, name, args) { + closeText(); + visibleReceived = true; + send({ + type: "TOOL_CALL_START", + toolCallId: callId, + toolCallName: name, + } as BaseEvent); + send({ + type: "TOOL_CALL_ARGS", + toolCallId: callId, + delta: JSON.stringify(args), + } as BaseEvent); + send({ + type: "TOOL_CALL_END", + toolCallId: callId, + } as BaseEvent); - await codex.runTurn(codexThreadId, WORKSPACE, turn.prompt, { - onText(delta) { - if (!textOpen) { - send({ - type: "TEXT_MESSAGE_START", - messageId, - role: "assistant", - } as BaseEvent); - textOpen = true; - } - textReceived = true; - send({ - type: "TEXT_MESSAGE_CONTENT", - messageId, - delta, - } as BaseEvent); - }, + const result = allowedToolNames.has(name) + ? await gateway.call(runAssertion, name, args, requestSignal) + : { + text: `Refused. ${name} was not granted to this Codex turn by OpenBot.`, + success: false, + }; + send({ + type: "TOOL_CALL_RESULT", + messageId: `${callId}-result`, + toolCallId: callId, + content: result.text, + role: "tool", + } as BaseEvent); + return result; + }, + }, + requestSignal, + ); }); - if (!textReceived) { - throw new Error("Codex completed without returning a text message."); - } - if (textOpen) { - send({ type: "TEXT_MESSAGE_END", messageId } as BaseEvent); + closeText(); + if (!visibleReceived) { + throw new Error( + "Codex completed without returning text or calling an OpenBot tool.", + ); } send({ type: "RUN_FINISHED", @@ -85,9 +168,7 @@ async function runAgent(input: RunAgentInput): Promise { runId: input.runId, } as BaseEvent); } catch (error) { - if (textOpen) { - send({ type: "TEXT_MESSAGE_END", messageId } as BaseEvent); - } + closeText(); send({ type: "RUN_ERROR", message: @@ -110,6 +191,27 @@ async function runAgent(input: RunAgentInput): Promise { }); } +async function serialiseThread( + threadId: string, + operation: () => Promise, +): Promise { + const previous = threadQueues.get(threadId) ?? Promise.resolve(); + let release: (() => void) | undefined; + const gate = new Promise((resolveGate) => { + release = resolveGate; + }); + const queued = previous.catch(() => {}).then(() => gate); + threadQueues.set(threadId, queued); + + await previous.catch(() => {}); + try { + return await operation(); + } finally { + release?.(); + if (threadQueues.get(threadId) === queued) threadQueues.delete(threadId); + } +} + Bun.serve({ port: PORT, idleTimeout: 255, @@ -121,7 +223,9 @@ Bun.serve({ status: "ok", authMode: account.authMode, planType: account.planType, - safety: "read-only-text-spike", + safety: "openbot-governed-tools", + threadRecovery: "persistent", + persistedThreads: threadStore.size(), }); } @@ -129,7 +233,7 @@ Bun.serve({ if (!hasManagedAgentToken(request, MANAGED_AGENT_TOKEN)) { return Response.json({ error: "Unauthorized." }, { status: 401 }); } - return runAgent((await request.json()) as RunAgentInput); + return runAgent((await request.json()) as RunAgentInput, request.signal); } return Response.json({ error: "Not found." }, { status: 404 }); @@ -138,5 +242,5 @@ Bun.serve({ const account = codex.accountSummary(); console.info( - `agent-codex listening on http://localhost:${PORT}/ag-ui (${account.authMode}, ${account.planType ?? "unknown plan"})`, + `agent-codex listening on http://localhost:${PORT}/ag-ui (${account.authMode}, ${account.planType ?? "unknown plan"}; persistent threads; OpenBot-governed tools)`, ); diff --git a/agent-codex/src/thread-store.ts b/agent-codex/src/thread-store.ts new file mode 100644 index 000000000..33364ae42 --- /dev/null +++ b/agent-codex/src/thread-store.ts @@ -0,0 +1,116 @@ +import { randomUUID } from "node:crypto"; +import { mkdir, readFile, rename, unlink, writeFile } from "node:fs/promises"; +import { dirname } from "node:path"; + +type ThreadState = { + version: 1; + threads: Record; +}; + +const EMPTY_STATE: ThreadState = { version: 1, threads: {} }; + +/** + * The durable join between an OpenBot Intelligence thread and its Codex app-server thread. + * + * Codex persists its own rollout, but it cannot know which OpenBot thread owns it. This small file is + * the missing half. Writes replace the file atomically, so killing the adapter during a write leaves + * either the previous complete mapping or the next one, never half a JSON document. + */ +export class CodexThreadStore { + private readonly threads: Map; + private writes = Promise.resolve(); + + private constructor( + private readonly path: string, + state: ThreadState, + ) { + this.threads = new Map(Object.entries(state.threads)); + } + + static async open(path: string): Promise { + let raw: string; + try { + raw = await readFile(path, "utf8"); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") { + return new CodexThreadStore(path, EMPTY_STATE); + } + throw error; + } + + return new CodexThreadStore(path, parseState(raw, path)); + } + + get(openbotThreadId: string): string | undefined { + return this.threads.get(openbotThreadId); + } + + size(): number { + return this.threads.size; + } + + async remember( + openbotThreadId: string, + codexThreadId: string, + ): Promise { + if (!openbotThreadId || !codexThreadId) { + throw new Error("Thread ids must not be empty."); + } + const write = this.writes + .catch(() => {}) + .then(async () => { + await mkdir(dirname(this.path), { recursive: true }); + const temporary = `${this.path}.${process.pid}.${randomUUID()}.tmp`; + const nextThreads = new Map(this.threads); + nextThreads.set(openbotThreadId, codexThreadId); + const state: ThreadState = { + version: 1, + threads: Object.fromEntries(nextThreads), + }; + try { + await writeFile(temporary, `${JSON.stringify(state, null, 2)}\n`, { + encoding: "utf8", + mode: 0o600, + }); + await rename(temporary, this.path); + this.threads.set(openbotThreadId, codexThreadId); + } catch (error) { + await unlink(temporary).catch(() => {}); + throw error; + } + }); + this.writes = write; + await write; + } +} + +function parseState(raw: string, path: string): ThreadState { + let value: unknown; + try { + value = JSON.parse(raw); + } catch { + throw new Error( + `Codex thread state at ${path} is not valid JSON. Refusing to forget existing conversations.`, + ); + } + + if ( + !isObject(value) || + value.version !== 1 || + !isObject(value.threads) || + Object.entries(value.threads).some( + ([openbotThreadId, codexThreadId]) => + !openbotThreadId || typeof codexThreadId !== "string" || !codexThreadId, + ) + ) { + throw new Error( + `Codex thread state at ${path} has an unsupported shape. Refusing to forget existing conversations.`, + ); + } + + return value as ThreadState; +} + +function isObject(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} diff --git a/agent-codex/src/tools.ts b/agent-codex/src/tools.ts new file mode 100644 index 000000000..14bc38e56 --- /dev/null +++ b/agent-codex/src/tools.ts @@ -0,0 +1,145 @@ +import type { RunAgentInput } from "@ag-ui/core"; + +export type CodexDynamicTool = { + type: "function"; + name: string; + description: string; + inputSchema: Record; +}; + +export type ToolResult = { + text: string; + success: boolean; +}; + +type ToolGatewayOptions = { + url: string; + token: string; + fetch?: typeof fetch; +}; + +const TOOL_NAME = /^[A-Za-z0-9_-]{1,64}$/; +const DEFAULT_PARAMETERS = { type: "object", properties: {} }; + +/** + * Only tools OpenBot says this deployment executes become Codex dynamic tools. + * + * AG-UI's `tools` array also contains browser-owned components and human-input tools. Routing those + * to the deployment gateway would turn a chart into a failed MCP call, so the signed run metadata's + * deployment-owned allowlist is the authority for which descriptions cross this boundary. + */ +export function dynamicToolsOf(input: RunAgentInput): CodexDynamicTool[] { + const deploymentTools = deploymentToolNames(input); + const seen = new Set(); + const tools: CodexDynamicTool[] = []; + + for (const tool of input.tools ?? []) { + if (!deploymentTools.has(tool.name) || seen.has(tool.name)) continue; + if (!TOOL_NAME.test(tool.name)) { + throw new Error( + `OpenBot tool ${tool.name} cannot be exposed to Codex because its name is not Responses-API safe.`, + ); + } + seen.add(tool.name); + tools.push({ + type: "function", + name: tool.name, + description: tool.description || "An OpenBot-governed tool.", + inputSchema: isObject(tool.parameters) + ? tool.parameters + : DEFAULT_PARAMETERS, + }); + } + + return tools; +} + +export function runAssertionOf(input: RunAgentInput): string { + const props = input.forwardedProps as { openbotRun?: unknown } | undefined; + return typeof props?.openbotRun === "string" ? props.openbotRun : ""; +} + +function deploymentToolNames(input: RunAgentInput): Set { + const props = input.forwardedProps as + | { openbotDeploymentTools?: unknown } + | undefined; + const names = props?.openbotDeploymentTools; + return new Set( + Array.isArray(names) + ? names.filter((name): name is string => typeof name === "string") + : [], + ); +} + +/** Calls an OpenBot tool through the deployment that owns its grant, policy and audit row. */ +export class OpenBotToolGateway { + private readonly fetch: typeof fetch; + + constructor(private readonly options: ToolGatewayOptions) { + this.fetch = options.fetch ?? fetch; + } + + async call( + run: string, + name: string, + args: Record, + signal?: AbortSignal, + ): Promise { + if (!this.options.token) { + return { + text: "Refused. This Codex coworker has no credential for OpenBot's tool gateway.", + success: false, + }; + } + if (!run) { + return { + text: "Refused. This run carried no signed statement of which Bot and person it is for.", + success: false, + }; + } + + try { + const response = await this.fetch(this.options.url, { + method: "POST", + headers: { + "content-type": "application/json", + "x-openbot-agent-token": this.options.token, + }, + body: JSON.stringify({ name, args, run }), + signal: signal + ? AbortSignal.any([signal, AbortSignal.timeout(65_000)]) + : AbortSignal.timeout(65_000), + }); + const body = (await response.json().catch(() => null)) as { + text?: unknown; + isError?: unknown; + error?: unknown; + } | null; + + if (!response.ok) { + const reason = + typeof body?.error === "string" + ? body.error + : `OpenBot's tool gateway returned HTTP ${response.status}.`; + return { text: `Refused. ${reason}`, success: false }; + } + + const text = + typeof body?.text === "string" + ? body.text + : "The OpenBot tool returned nothing."; + return { text, success: body?.isError !== true }; + } catch (error) { + return { + text: `That tool could not be called through OpenBot: ${ + error instanceof Error ? error.message : "unknown error" + }`, + success: false, + }; + } + } +} + +function isObject(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} diff --git a/agent-codex/tests/codex-client.test.ts b/agent-codex/tests/codex-client.test.ts new file mode 100644 index 000000000..9b79c1adf --- /dev/null +++ b/agent-codex/tests/codex-client.test.ts @@ -0,0 +1,315 @@ +import { describe, expect, test } from "bun:test"; +import type { ChildProcessWithoutNullStreams } from "node:child_process"; +import { EventEmitter } from "node:events"; +import { createInterface } from "node:readline"; +import { PassThrough } from "node:stream"; +import { CodexAppServerClient } from "../src/codex-client"; +import type { CodexDynamicTool } from "../src/tools"; + +type Message = { + id?: number | string; + method?: string; + params?: Record; + result?: Record; +}; + +const dynamicTool: CodexDynamicTool = { + type: "function", + name: "search_files", + description: "Search files", + inputSchema: { + type: "object", + properties: { query: { type: "string" } }, + }, +}; + +class FakeAppServer { + readonly messages: Message[] = []; + readonly process: ChildProcessWithoutNullStreams; + toolResponse: Record | undefined; + toolResponses: Record[] = []; + + private readonly stdin = new PassThrough(); + private readonly stdout = new PassThrough(); + private readonly stderr = new PassThrough(); + + constructor( + private readonly turn: + | "tool" + | "native" + | "duplicate" + | "malformed" + | "waiting" = "tool", + ) { + const child = new EventEmitter() as EventEmitter & { + stdin: PassThrough; + stdout: PassThrough; + stderr: PassThrough; + kill(): boolean; + }; + child.stdin = this.stdin; + child.stdout = this.stdout; + child.stderr = this.stderr; + child.kill = () => { + child.emit("exit", 0, null); + return true; + }; + this.process = child as unknown as ChildProcessWithoutNullStreams; + + const lines = createInterface({ input: this.stdin }); + lines.on("line", (line) => this.receive(JSON.parse(line) as Message)); + } + + private receive(message: Message): void { + this.messages.push(message); + if ( + typeof message.id === "string" && + message.id.startsWith("dynamic-tool-call") && + message.result + ) { + this.toolResponse = message.result; + this.toolResponses.push(message.result); + if (this.turn === "duplicate" && this.toolResponses.length === 1) { + this.sendToolCall("dynamic-tool-call-duplicate"); + return; + } + this.send({ + method: "item/agentMessage/delta", + params: { + threadId: "codex-thread", + turnId: "turn-1", + itemId: "message-1", + delta: "I found three files.", + }, + }); + this.send({ + method: "turn/completed", + params: { + threadId: "codex-thread", + turn: { id: "turn-1", status: "completed" }, + }, + }); + return; + } + if (message.id === undefined || !message.method) return; + + switch (message.method) { + case "initialize": + this.result(message.id, {}); + break; + case "account/read": + this.result(message.id, { + account: { type: "chatgpt", planType: "plus" }, + }); + break; + case "config/read": + this.result(message.id, { + config: { mcp_servers: { existing_server: { command: "unsafe" } } }, + }); + break; + case "thread/start": + this.result(message.id, { thread: { id: "codex-thread" } }); + break; + case "thread/resume": + this.result(message.id, { + thread: { id: String(message.params?.threadId) }, + }); + break; + case "turn/start": + this.result(message.id, { turn: { id: "turn-1" } }); + queueMicrotask(() => { + if (this.turn === "tool" || this.turn === "duplicate") { + this.sendToolCall("dynamic-tool-call"); + } else if (this.turn === "malformed") { + this.sendToolCall("dynamic-tool-call", []); + } else if (this.turn === "native") { + this.send({ + method: "item/started", + params: { + threadId: "codex-thread", + turnId: "turn-1", + item: { id: "command-1", type: "commandExecution" }, + }, + }); + } + }); + break; + case "turn/interrupt": + this.result(message.id, {}); + break; + default: + this.result(message.id, {}); + } + } + + private result(id: number | string, result: Record): void { + this.send({ id, result }); + } + + private sendToolCall(id: string, args: unknown = { query: "budget" }): void { + this.send({ + id, + method: "item/tool/call", + params: { + threadId: "codex-thread", + turnId: "turn-1", + callId: "call-1", + namespace: null, + tool: "search_files", + arguments: args, + }, + }); + } + + private send(message: Message): void { + this.stdout.write(`${JSON.stringify(message)}\n`); + } +} + +describe("CodexAppServerClient", () => { + test("resumes persistent threads and answers dynamic tool calls", async () => { + const server = new FakeAppServer(); + const client = new CodexAppServerClient(() => server.process); + await client.start(); + + expect(client.accountSummary()).toEqual({ + authMode: "chatgpt", + planType: "plus", + }); + await expect( + client.startThread("/workspace", "governed only", [dynamicTool]), + ).resolves.toBe("codex-thread"); + await client.resumeThread("codex-thread", "/workspace", "governed only"); + + const text: string[] = []; + await client.runTurn("codex-thread", "/workspace", "Find budget", { + onText(delta) { + text.push(delta); + }, + async onToolCall(callId, name, args) { + expect(callId).toBe("call-1"); + expect(name).toBe("search_files"); + expect(args).toEqual({ query: "budget" }); + return { text: "three files", success: true }; + }, + }); + + expect(text).toEqual(["I found three files."]); + expect(server.toolResponse).toEqual({ + contentItems: [{ type: "inputText", text: "three files" }], + success: true, + }); + const initialize = server.messages.find( + (message) => message.method === "initialize", + ); + expect(initialize?.params?.capabilities).toEqual({ experimentalApi: true }); + const start = server.messages.find( + (message) => message.method === "thread/start", + ); + expect(start?.params?.dynamicTools).toEqual([dynamicTool]); + expect(start?.params?.config).toMatchObject({ + mcp_servers: { existing_server: { enabled: false } }, + features: { apps: false, plugins: false, multi_agent: false }, + web_search: "disabled", + apps: { _default: { enabled: false } }, + tools: { web_search: false, view_image: false }, + }); + const turn = server.messages.find( + (message) => message.method === "turn/start", + ); + expect(turn?.params?.sandboxPolicy).toEqual({ + type: "readOnly", + networkAccess: false, + }); + client.stop(); + }); + + test("interrupts a turn that attempts a Codex-native action", async () => { + const server = new FakeAppServer("native"); + const client = new CodexAppServerClient(() => server.process); + await client.start(); + await client.startThread("/workspace", "governed only", []); + + await expect( + client.runTurn("codex-thread", "/workspace", "Run a command", { + onText() {}, + async onToolCall() { + return { text: "unreachable", success: false }; + }, + }), + ).rejects.toThrow("OpenBot refused it"); + expect( + server.messages.some((message) => message.method === "turn/interrupt"), + ).toBe(true); + client.stop(); + }); + + test("executes a repeated dynamic-tool call id only once", async () => { + const server = new FakeAppServer("duplicate"); + const client = new CodexAppServerClient(() => server.process); + await client.start(); + await client.startThread("/workspace", "governed only", [dynamicTool]); + let calls = 0; + + await client.runTurn("codex-thread", "/workspace", "Search once", { + onText() {}, + async onToolCall() { + calls += 1; + return { text: "three files", success: true }; + }, + }); + + expect(calls).toBe(1); + expect(server.toolResponses).toHaveLength(2); + expect(server.toolResponses[1]).toMatchObject({ success: false }); + client.stop(); + }); + + test("interrupts Codex when OpenBot aborts the request", async () => { + const server = new FakeAppServer("waiting"); + const client = new CodexAppServerClient(() => server.process); + await client.start(); + await client.startThread("/workspace", "governed only", []); + const abort = new AbortController(); + + const turn = client.runTurn( + "codex-thread", + "/workspace", + "Wait forever", + { + onText() {}, + async onToolCall() { + return { text: "unreachable", success: false }; + }, + }, + abort.signal, + ); + abort.abort(); + + await expect(turn).rejects.toThrow("OpenBot ended the request"); + expect( + server.messages.some((message) => message.method === "turn/interrupt"), + ).toBe(true); + client.stop(); + }); + + test("does not coerce malformed tool arguments into an action", async () => { + const server = new FakeAppServer("malformed"); + const client = new CodexAppServerClient(() => server.process); + await client.start(); + await client.startThread("/workspace", "governed only", [dynamicTool]); + let calls = 0; + + await client.runTurn("codex-thread", "/workspace", "Search once", { + onText() {}, + async onToolCall() { + calls += 1; + return { text: "should not run", success: true }; + }, + }); + + expect(calls).toBe(0); + expect(server.toolResponses[0]).toMatchObject({ success: false }); + client.stop(); + }); +}); diff --git a/agent-codex/tests/history.test.ts b/agent-codex/tests/history.test.ts index 193e11a73..546e2e698 100644 --- a/agent-codex/tests/history.test.ts +++ b/agent-codex/tests/history.test.ts @@ -20,7 +20,8 @@ describe("toCodexTurnInput", () => { expect(result.developerInstructions).toContain( "You are the finance coworker.", ); - expect(result.developerInstructions).toContain("Respond with text only"); + expect(result.developerInstructions).toContain("OpenBot dynamic tools"); + expect(result.developerInstructions).toContain("Never run shell commands"); }); test("refuses an empty turn", () => { diff --git a/agent-codex/tests/thread-store.test.ts b/agent-codex/tests/thread-store.test.ts new file mode 100644 index 000000000..d334a7c48 --- /dev/null +++ b/agent-codex/tests/thread-store.test.ts @@ -0,0 +1,70 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { mkdtemp, readFile, rm, stat } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { CodexThreadStore } from "../src/thread-store"; + +const temporaryDirectories: string[] = []; + +afterEach(async () => { + await Promise.all( + temporaryDirectories.splice(0).map((path) => rm(path, { recursive: true })), + ); +}); + +async function statePath(): Promise { + const directory = await mkdtemp(join(tmpdir(), "openbot-codex-state-")); + temporaryDirectories.push(directory); + return join(directory, "nested", "threads.json"); +} + +describe("CodexThreadStore", () => { + test("recovers OpenBot-to-Codex joins after reopening", async () => { + const path = await statePath(); + const store = await CodexThreadStore.open(path); + await store.remember("openbot-1", "codex-1"); + + const recovered = await CodexThreadStore.open(path); + expect(recovered.get("openbot-1")).toBe("codex-1"); + expect(recovered.size()).toBe(1); + + const mode = (await stat(path)).mode & 0o777; + expect(mode).toBe(0o600); + }); + + test("serialises concurrent atomic writes without losing a mapping", async () => { + const path = await statePath(); + const store = await CodexThreadStore.open(path); + + await Promise.all([ + store.remember("openbot-1", "codex-1"), + store.remember("openbot-2", "codex-2"), + ]); + + const state = JSON.parse(await readFile(path, "utf8")) as { + threads: Record; + }; + expect(state.threads).toEqual({ + "openbot-1": "codex-1", + "openbot-2": "codex-2", + }); + }); + + test("refuses to silently replace corrupt recovery state", async () => { + const path = await statePath(); + await Bun.write(path, "not json"); + + await expect(CodexThreadStore.open(path)).rejects.toThrow( + "Refusing to forget existing conversations", + ); + }); + + test("rejects unsupported state shapes", async () => { + const path = await statePath(); + await Bun.write(path, JSON.stringify({ version: 2, threads: {} })); + + await expect(CodexThreadStore.open(path)).rejects.toThrow( + "unsupported shape", + ); + }); +}); diff --git a/agent-codex/tests/tools.test.ts b/agent-codex/tests/tools.test.ts new file mode 100644 index 000000000..df5bc3d4f --- /dev/null +++ b/agent-codex/tests/tools.test.ts @@ -0,0 +1,142 @@ +import { describe, expect, test } from "bun:test"; +import type { RunAgentInput } from "@ag-ui/core"; +import { + dynamicToolsOf, + OpenBotToolGateway, + runAssertionOf, +} from "../src/tools"; + +function input(overrides: Partial = {}): RunAgentInput { + return { + threadId: "thread-1", + runId: "run-1", + state: {}, + messages: [], + context: [], + tools: [], + forwardedProps: {}, + ...overrides, + }; +} + +describe("Codex dynamic tools", () => { + test("exposes only deployment-owned tools and preserves their schema", () => { + const tools = dynamicToolsOf( + input({ + tools: [ + { + name: "search_files", + description: "Search Drive", + parameters: { + type: "object", + properties: { query: { type: "string" } }, + required: ["query"], + }, + }, + { + name: "render_chart", + description: "Browser-owned chart", + parameters: { type: "object" }, + }, + ], + forwardedProps: { + openbotDeploymentTools: ["search_files", "missing", "search_files"], + openbotRun: "signed-run", + }, + }), + ); + + expect(tools).toEqual([ + { + type: "function", + name: "search_files", + description: "Search Drive", + inputSchema: { + type: "object", + properties: { query: { type: "string" } }, + required: ["query"], + }, + }, + ]); + }); + + test("rejects tool names the Codex dynamic-tool protocol cannot represent", () => { + expect(() => + dynamicToolsOf( + input({ + tools: [{ name: "unsafe tool", description: "No", parameters: {} }], + forwardedProps: { openbotDeploymentTools: ["unsafe tool"] }, + }), + ), + ).toThrow("Responses-API safe"); + }); + + test("reads the opaque signed run assertion without interpreting it", () => { + expect( + runAssertionOf( + input({ forwardedProps: { openbotRun: "opaque.signature" } }), + ), + ).toBe("opaque.signature"); + expect(runAssertionOf(input())).toBe(""); + }); +}); + +describe("OpenBotToolGateway", () => { + test("sends the signed run and agent credential to OpenBot", async () => { + let request: { url: string; init?: RequestInit } | undefined; + const fakeFetch = (async ( + url: string | URL | Request, + init?: RequestInit, + ) => { + request = { url: String(url), init }; + return Response.json({ text: "three files", isError: false }); + }) as typeof fetch; + const gateway = new OpenBotToolGateway({ + url: "http://openbot.test/api/agent-tools/call", + token: "agent-secret", + fetch: fakeFetch, + }); + + await expect( + gateway.call("signed-run", "search_files", { query: "budget" }), + ).resolves.toEqual({ text: "three files", success: true }); + expect(request?.url).toBe("http://openbot.test/api/agent-tools/call"); + expect(request?.init?.headers).toEqual({ + "content-type": "application/json", + "x-openbot-agent-token": "agent-secret", + }); + expect(JSON.parse(String(request?.init?.body))).toEqual({ + name: "search_files", + args: { query: "budget" }, + run: "signed-run", + }); + }); + + test("returns refusals to Codex as failed tool results", async () => { + const deniedFetch = (async () => + Response.json( + { error: "Policy denied this call." }, + { status: 403 }, + )) as unknown as typeof fetch; + const denied = new OpenBotToolGateway({ + url: "http://openbot.test/tools", + token: "token", + fetch: deniedFetch, + }); + + await expect(denied.call("run", "delete_file", {})).resolves.toEqual({ + text: "Refused. Policy denied this call.", + success: false, + }); + await expect( + new OpenBotToolGateway({ + url: "http://openbot.test/tools", + token: "", + fetch: deniedFetch, + }).call("run", "delete_file", {}), + ).resolves.toMatchObject({ success: false }); + await expect(denied.call("", "delete_file", {})).resolves.toMatchObject({ + success: false, + }); + }); +}); diff --git a/bun.lock b/bun.lock index 415b9ce62..77e2bbd39 100644 --- a/bun.lock +++ b/bun.lock @@ -13,6 +13,18 @@ "yaml": "^2.9.0", }, }, + "agent-codex": { + "name": "@openbot/agent-codex", + "version": "0.0.2", + "dependencies": { + "@ag-ui/core": "0.0.57", + "@ag-ui/encoder": "0.0.57", + }, + "devDependencies": { + "@types/bun": "^1.3.3", + "typescript": "^5.9.3", + }, + }, "app": { "name": "app", "version": "0.0.0", @@ -457,6 +469,8 @@ "@nodelib/fs.walk": ["@nodelib/fs.walk@1.2.8", "", { "dependencies": { "@nodelib/fs.scandir": "2.1.5", "fastq": "^1.6.0" } }, "sha512-oGB+UxlgWcgQkgwo8GcEGwemoTFt3FIO9ababBmaGwXIoBKZ+GTy0pP185beGg7Llih/NSHSV2XAs1lnznocSg=="], + "@openbot/agent-codex": ["@openbot/agent-codex@workspace:agent-codex"], + "@opentelemetry/api": ["@opentelemetry/api@1.9.1", "", {}, "sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q=="], "@opentelemetry/semantic-conventions": ["@opentelemetry/semantic-conventions@1.43.0", "", {}, "sha512-eSYWTm620tTk45EKSedaUL8MFYI8hW164hIXsgIHyxu3VobUB3fFCu5t0hQby6OoWRPsG1KkKUG2M5UadiLiVg=="], diff --git a/package.json b/package.json index 84c4ee66f..b07afb510 100644 --- a/package.json +++ b/package.json @@ -6,6 +6,7 @@ "packageManager": "bun@1.3.14", "workspaces": [ "app", + "agent-codex", "server", "worker" ], diff --git a/scripts/start.sh b/scripts/start.sh index 0c26da181..f417a9274 100755 --- a/scripts/start.sh +++ b/scripts/start.sh @@ -160,6 +160,10 @@ identifies_as_openbot() { curl -fsS --max-time 3 "http://localhost:$port/" 2>/dev/null \ | grep -qi '[^<]*OpenBot' ;; + agent-codex) + curl -fsS --max-time 3 "http://localhost:$port/health" 2>/dev/null \ + | grep -q '"safety":"openbot-governed-tools"' + ;; # Compose services on dedicated loopback ports, answering a route named for this stack. *) curl -fsS --max-time 3 "http://localhost:$port/health" >/dev/null 2>&1 @@ -245,6 +249,8 @@ fi wait_for "http://localhost:$COMPUTER_PORT/health" "agent-computer" if [ "$CODEX_AGENT_ENABLED" = "true" ]; then CODEX_AGENT_WORKSPACE="$(setting CODEX_AGENT_WORKSPACE "$ROOT/.openbot-codex/workspace")" + CODEX_AGENT_STATE="$(setting CODEX_AGENT_STATE "$ROOT/.openbot-codex/threads.json")" + OPENBOT_TOOL_URL="$(setting OPENBOT_TOOL_URL "http://localhost:${SERVER_PORT}/api/agent-tools/call")" mkdir -p "$CODEX_AGENT_WORKSPACE" require_free_or_ours "$CODEX_AGENT_PORT" "agent-codex" if identifies_as_openbot "$CODEX_AGENT_PORT" "agent-codex"; then @@ -255,7 +261,10 @@ if [ "$CODEX_AGENT_ENABLED" = "true" ]; then nohup env \ PORT="$CODEX_AGENT_PORT" \ MANAGED_AGENT_TOKEN="$MANAGED_AGENT_TOKEN" \ + AGENT_TOOL_TOKEN="$AGENT_TOOL_TOKEN" \ + OPENBOT_TOOL_URL="$OPENBOT_TOOL_URL" \ CODEX_AGENT_WORKSPACE="$CODEX_AGENT_WORKSPACE" \ + CODEX_AGENT_STATE="$CODEX_AGENT_STATE" \ bun agent-codex/src/index.ts >"$LOGS/agent-codex.log" 2>&1 </dev/null &) wait_for "http://localhost:$CODEX_AGENT_PORT/health" "agent-codex" 60 else @@ -284,7 +293,7 @@ green " coworker tables migrated" # and masks it. So: report what was resolved, and do not re-read the file. green " managed coworker endpoint: $MANAGED_AGENT_AG_UI_URL" if [ "$CODEX_AGENT_ENABLED" = "true" ]; then - green " Codex coworker: ChatGPT login · read-only text compatibility mode" + green " Codex coworker: ChatGPT login · persistent threads · OpenBot-governed tools" fi info "2/4 Server" diff --git a/tests/workspace.test.ts b/tests/workspace.test.ts index aacea898d..a18db1b11 100644 --- a/tests/workspace.test.ts +++ b/tests/workspace.test.ts @@ -13,16 +13,23 @@ function packageManifest(path: string) { } describe("OpenBot workspace", () => { - test("defines the app, server, and worker packages", () => { + test("defines every package installed by the root workspace", () => { const rootManifest = JSON.parse( readFileSync(join(repositoryRoot, "package.json"), "utf8"), ) as { workspaces: string[] }; - expect(rootManifest.workspaces).toEqual(["app", "server", "worker"]); + expect(rootManifest.workspaces).toEqual([ + "app", + "agent-codex", + "server", + "worker", + ]); for (const packageName of rootManifest.workspaces) { expect(existsSync(join(repositoryRoot, packageName))).toBe(true); - expect(packageManifest(packageName).name).toBe(packageName); + expect(packageManifest(packageName).name).toBe( + packageName === "agent-codex" ? "@openbot/agent-codex" : packageName, + ); } }); });