diff --git a/.env.example b/.env.example index f46d803a2..bce08ab0d 100644 --- a/.env.example +++ b/.env.example @@ -137,8 +137,23 @@ 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. 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 -# 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 +344,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/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 new file mode 100644 index 000000000..458457d92 --- /dev/null +++ b/agent-codex/README.md @@ -0,0 +1,21 @@ +# 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`. + +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. 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/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..faaf64f45 --- /dev/null +++ b/agent-codex/package.json @@ -0,0 +1,21 @@ +{ + "name": "@openbot/agent-codex", + "version": "0.0.2", + "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..1fb5cc351 --- /dev/null +++ b/agent-codex/src/codex-client.ts @@ -0,0 +1,604 @@ +import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; +import { createInterface } from "node:readline"; +import type { CodexDynamicTool, ToolResult } from "./tools"; + +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 ThreadResult = { + thread: { id: string }; +}; + +type TurnStartResult = { + turn: { id: string }; +}; + +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; +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 child = this.spawnAppServer(); + 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.2", + }, + capabilities: { experimentalApi: true }, + }); + this.notify("initialized", {}); + + const result = (await this.request("account/read", { + refreshToken: false, + })) as { + account?: { type?: string; planType?: string | null } | null; + }; + 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, + }; + + 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 { + if (!this.account) { + throw new Error("Codex app-server has not finished starting."); + } + return this.account; + } + + async startThread( + cwd: string, + developerInstructions: string, + dynamicTools: CodexDynamicTool[], + ): Promise { + const result = (await this.request("thread/start", { + cwd, + approvalPolicy: "never", + sandbox: "read-only", + serviceName: "openbot_local_codex", + developerInstructions, + dynamicTools, + config: this.safetyConfig, + ephemeral: false, + })) 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 { + 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 = turnTimeoutMs(); + let settled = false; + let resolveCompletion: (() => void) | undefined; + let rejectCompletion: ((error: Error) => void) | undefined; + const completed = new Promise((resolve, 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(() => { + 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") { + 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") { + const item = isObject(params.item) ? params.item : {}; + if ( + item.type === "agentMessage" && + typeof item.id === "string" && + !streamedItems.has(item.id) && + typeof item.text === "string" && + item.text + ) { + callbacks.onText(item.text); + } + return; + } + + if (message.method === "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 = isObject(params.turn) ? params.turn : {}; + if ( + active.turnId && + typeof turn.id === "string" && + turn.id !== active.turnId + ) + return; + if (turn.status === "completed") { + finish(); + } else { + const error = isObject(turn.error) ? turn.error : {}; + fail( + new Error( + turnError ?? + (typeof error.message === "string" + ? error.message + : undefined) ?? + `Codex turn ended with status ${String(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", networkAccess: false }, + effort: "low", + })) as TurnStartResult; + 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); + } + + 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) { + void 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); + } + + private async answerServerRequest(message: JsonRpcMessage): Promise { + if ( + message.method === "item/commandExecution/requestApproval" || + message.method === "item/fileChange/requestApproval" + ) { + 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: "OpenBot does not expose that Codex-native action.", + }, + }); + } + + private failAll(error: Error): void { + for (const pending of this.pending.values()) { + clearTimeout(pending.timeout); + 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 new file mode 100644 index 000000000..65f4bf11d --- /dev/null +++ b/agent-codex/src/history.ts @@ -0,0 +1,48 @@ +import type { RunAgentInput } from "@ag-ui/core"; + +export type CodexTurnInput = { + developerInstructions: string; + prompt: string; +}; + +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. + * + * 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 + ? `${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 new file mode 100644 index 000000000..9f332bc1b --- /dev/null +++ b/agent-codex/src/index.ts @@ -0,0 +1,246 @@ +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"; +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(); +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 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 threadQueues = new Map>(); + +async function runAgent( + input: RunAgentInput, + requestSignal: AbortSignal, +): 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); + + let messageSequence = 0; + let messageId = ""; + let textOpen = false; + let visibleReceived = false; + const closeText = () => { + if (!textOpen) return; + send({ type: "TEXT_MESSAGE_END", messageId } as BaseEvent); + textOpen = false; + }; + + try { + await serialiseThread(input.threadId, async () => { + const turn = toCodexTurnInput(input); + const dynamicTools = dynamicToolsOf(input); + const allowedToolNames = new Set( + dynamicTools.map((tool) => tool.name), + ); + 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); + + 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, + ); + }); + + closeText(); + if (!visibleReceived) { + throw new Error( + "Codex completed without returning text or calling an OpenBot tool.", + ); + } + send({ + type: "RUN_FINISHED", + threadId: input.threadId, + runId: input.runId, + } as BaseEvent); + } catch (error) { + closeText(); + 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", + }, + }); +} + +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, + 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: "openbot-governed-tools", + threadRecovery: "persistent", + persistedThreads: threadStore.size(), + }); + } + + 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, request.signal); + } + + 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"}; 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 new file mode 100644 index 000000000..546e2e698 --- /dev/null +++ b/agent-codex/tests/history.test.ts @@ -0,0 +1,32 @@ +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("OpenBot dynamic tools"); + expect(result.developerInstructions).toContain("Never run shell commands"); + }); + + test("refuses an empty turn", () => { + expect(() => + toCodexTurnInput(input([{ role: "system", content: "A role" }])), + ).toThrow("needs a user message"); + }); +}); 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/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/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/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/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 74beebab6..f417a9274 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. @@ -146,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 @@ -215,12 +233,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 +247,30 @@ 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")" + 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 + 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" \ + 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 + wait_for "http://localhost:$BOT_PORT/health" "agent-bot" + wait_for "http://localhost:$LANGGRAPH_PORT/health" "agent-langgraph" +fi for table in agent_profiles agent_preferences; do if ! docker compose exec -T postgres \ @@ -251,6 +292,9 @@ green " coworker tables migrated" # grep survives inside `setting()` only because `local v="$(...)"` takes `local`'s own exit status # 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 · persistent threads · OpenBot-governed tools" +fi info "2/4 Server" require_free_or_ours "$SERVER_PORT" server @@ -302,16 +346,18 @@ if identifies_as_openbot "$SERVER_PORT" server; then fi if ! identifies_as_openbot "$SERVER_PORT" server; then if [ "$ONE_COMPUTER_EACH" = "true" ]; then - (cd server && PORT="$SERVER_PORT" \ + (cd server && nohup env \ + PORT="$SERVER_PORT" \ COMPUTER_SUPERVISOR_URL="http://localhost:$SUPERVISOR_PORT" \ SUPERVISOR_TOKEN="$SUPERVISOR_TOKEN" \ COMPUTER_TOKEN="$COMPUTER_TOKEN" \ WORKER_SHARED_SECRET="$WORKER_SHARED_SECRET" \ - bun --env-file=../.env src/index.ts >"$LOGS/server.log" 2>&1 &) + bun --env-file=../.env src/index.ts >"$LOGS/server.log" 2>&1 </dev/null &) else - (cd server && PORT="$SERVER_PORT" \ + (cd server && nohup env \ + PORT="$SERVER_PORT" \ WORKER_SHARED_SECRET="$WORKER_SHARED_SECRET" \ - bun --env-file=../.env src/index.ts >"$LOGS/server.log" 2>&1 &) + bun --env-file=../.env src/index.ts >"$LOGS/server.log" 2>&1 </dev/null &) fi fi wait_for_openbot "$SERVER_PORT" server @@ -335,10 +381,11 @@ wait_for_openbot "$SERVER_PORT" server if ! pgrep -f "bun worker/src/index.ts" >/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 &) info " worker: started (routine sweep loop)" sleep 1 if ! pgrep -f "bun worker/src/index.ts" >/dev/null 2>&1; then @@ -368,7 +415,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 </dev/null &) fi wait_for_openbot "$APP_PORT" app @@ -395,6 +442,7 @@ Try: Logs: $LOGS Routine sweep worker: $LOGS/worker.log Stop the routine worker: pkill -f 'bun worker/src/index.ts' +Stop the Codex coworker: pkill -f 'bun agent-codex/src/index.ts' Stop Docker services: docker compose down A Bot's computer is made by the supervisor rather than by compose, so it keeps running: docker rm -f \$(docker ps -q --filter label=openbot.supervisor=true) 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, + ); } }); });