From d050a9b997d8e9968d9af9275ff8f8f8650ff058 Mon Sep 17 00:00:00 2001 From: Zhifei Li Date: Thu, 1 Oct 2026 00:06:16 -0700 Subject: [PATCH] feat(agent): add Codex chat backend and deployment support --- deploy/CODEX.md | 61 +++++++++ deploy/deploy.sh | 4 +- deploy/pixelrag-agent-codex.conf | 14 ++ web/agent-server.mjs | 214 ++++++++----------------------- web/lib/agent-tools.mjs | 119 +++++++++++++++++ web/lib/codex-backend.mjs | 118 +++++++++++++++++ web/lib/codex-backend.test.mjs | 101 +++++++++++++++ web/lib/codex-mcp.mjs | 29 +++++ 8 files changed, 496 insertions(+), 164 deletions(-) create mode 100644 deploy/CODEX.md create mode 100644 deploy/pixelrag-agent-codex.conf create mode 100644 web/lib/agent-tools.mjs create mode 100644 web/lib/codex-backend.mjs create mode 100644 web/lib/codex-backend.test.mjs create mode 100644 web/lib/codex-mcp.mjs diff --git a/deploy/CODEX.md b/deploy/CODEX.md new file mode 100644 index 0000000..472d3dd --- /dev/null +++ b/deploy/CODEX.md @@ -0,0 +1,61 @@ +# Codex chat backend + +The standalone agent server supports `CHAT_BACKEND=claude` (the existing +default) and `CHAT_BACKEND=codex`. Both use the same search and screenshot +handlers and expose the existing `/chat` SSE contract to the Next.js proxy. +The Codex path uses the logged-in Codex CLI via `codex exec --json`, with a +per-conversation stdio MCP adapter for `pixelrag_search` and `pixelrag_tile`. +It supports text history and an uploaded image, search/gallery events, +disconnect cancellation, tool-call limits, and a wall-time timeout. + +Prerequisites: Node 22+, the existing web dependencies, Codex CLI 0.159.1 +(the tested version), and `codex login` under the service's OS account. +No OpenAI API key is required when using that account's ChatGPT login. + +```sh +CHAT_BACKEND=codex CODEX_BIN=/home/zlab-ssd1/.local/bin/codex \ + PIXELRAG_SEARCH_URL=http://localhost:30001 AGENT_PORT=30011 \ + /home/yichuan/.nix-profile/bin/node web/agent-server.mjs +``` + +`CHAT_CODEX_MODEL` optionally selects a model; otherwise Codex chooses its +default. `CHAT_CODEX_TIMEOUT_MS` defaults to 180000. `CHAT_MAX_TURNS` bounds +MCP tool calls for Codex. Existing request/concurrency limits apply to both +providers. `CHAT_MAX_BUDGET_USD` and Claude thinking-token limits apply only +to Claude; Codex subscription use is bounded by request, tool, and time limits. + +The model runs in an empty temporary directory with a read-only sandbox, +shell, web search, browser, external apps, plugins, and image-file reading +disabled. Only the two PixelRAG MCP tools are preapproved. A random local +bearer token isolates each conversation's tool endpoint. Tool arguments +are validated with the same Zod schemas as Claude's tools. + +The production drop-in `pixelrag-agent-codex.conf` runs the original +`/home/yichuan/visrag` checkout under `zlab-ssd1`, the account with the +working Codex login. It preserves the existing port, rate limits, CORS, +and search-backend configuration. Adjust its paths and account before +using it on another host. Do not copy authentication tokens between users. + +An administrator can apply it with: + +```sh +sudo install -m 0644 deploy/pixelrag-agent-codex.conf \ + /etc/systemd/system/pixelrag-agent.service.d/codex.conf +sudo systemctl daemon-reload +sudo systemctl restart pixelrag-agent.service +curl --fail http://localhost:30010/health +curl --fail http://localhost:30010/ready +``` + +The health response includes `backend: "codex"`. Removing the drop-in and +restarting the service restores the original Claude deployment. + +Validation: + +```sh +node --test web/lib/*.test.mjs +node --check web/agent-server.mjs +``` + +Official documentation: [Non-interactive execution](https://learn.chatgpt.com/docs/non-interactive-mode) +and [MCP tool configuration](https://learn.chatgpt.com/docs/config-file/config-reference). diff --git a/deploy/deploy.sh b/deploy/deploy.sh index 0894429..4fcbdcb 100755 --- a/deploy/deploy.sh +++ b/deploy/deploy.sh @@ -45,8 +45,8 @@ if changed '^uv\.lock$'; then fi # 2. Agent backend — cheap restart, safe to automate. -if changed '^web/agent-server\.mjs$'; then - say "agent-server.mjs changed -> restart pixelrag-agent" +if changed '^web/(agent-server\.mjs|lib/(agent-tools|codex-backend|codex-mcp|retrieval-health)\.mjs)$'; then + say "agent backend changed -> restart pixelrag-agent" sudo systemctl restart pixelrag-agent.service sleep 2 say "pixelrag-agent: $(systemctl is-active pixelrag-agent.service)" diff --git a/deploy/pixelrag-agent-codex.conf b/deploy/pixelrag-agent-codex.conf new file mode 100644 index 0000000..7f47f94 --- /dev/null +++ b/deploy/pixelrag-agent-codex.conf @@ -0,0 +1,14 @@ +# Install as /etc/systemd/system/pixelrag-agent.service.d/codex.conf. +# This machine's working Codex login belongs to zlab-ssd1. The service must +# run under that account to use it; do not copy auth tokens between users. +[Service] +User=zlab-ssd1 +Group=zlab-ssd1 +WorkingDirectory=/home/yichuan/visrag/web +Environment=HOME=/home/zlab-ssd1 +Environment=CODEX_HOME=/home/zlab-ssd1/.codex +Environment=CHAT_BACKEND=codex +Environment=CODEX_BIN=/home/zlab-ssd1/.local/bin/codex +Environment=CHAT_CODEX_TIMEOUT_MS=180000 +ExecStart= +ExecStart=/home/yichuan/.nix-profile/bin/node /home/yichuan/visrag/web/agent-server.mjs diff --git a/web/agent-server.mjs b/web/agent-server.mjs index 9127e93..dbfd2ae 100644 --- a/web/agent-server.mjs +++ b/web/agent-server.mjs @@ -2,16 +2,21 @@ /** * PixelRAG Agent backend — standalone SSE server. * - * Runs the Claude Agent SDK with subscription auth (uses the logged-in + * Select CHAT_BACKEND=codex for the logged-in Codex CLI, or claude for the + * Claude Agent SDK with subscription auth (uses the logged-in * `claude` CLI on this machine — no ANTHROPIC_API_KEY needed). Exposes the * same agent loop + pixelrag tools as the Next.js /api/chat route, so the * deployed Vercel frontend can proxy to it instead of running the SDK in * serverless (where the native CLI binary and credentials don't exist). * - * Run on a machine where `claude` is logged in: - * node deploy/agent-server.mjs + * Run on a machine where the selected CLI is logged in: + * CHAT_BACKEND=codex node web/agent-server.mjs * * Env: + * CHAT_BACKEND claude (default) or codex + * CODEX_BIN Codex CLI executable (default codex on PATH) + * CHAT_CODEX_MODEL optional model override; otherwise CLI default + * CHAT_CODEX_TIMEOUT_MS per-request wall time (default 180000) * AGENT_PORT listen port (default 30010) * PIXELRAG_SEARCH_URL search API base (default http://localhost:30001) * CHAT_MAX_BUDGET_USD per-conversation budget cap (default 0.50) @@ -19,21 +24,15 @@ */ import http from "node:http" -import { query, tool, createSdkMcpServer } from "@anthropic-ai/claude-agent-sdk" -import { z } from "zod" +import { createTools } from "./lib/agent-tools.mjs" +import { runCodex, codexTool } from "./lib/codex-backend.mjs" -import { - RETRIEVAL_OK, - RETRIEVAL_BACKEND_DOWN, - classifyStatus, - classifyThrow, - createRetrievalHealth, - recordRetrieval, - isUngrounded, - backendDownInstruction, - UNGROUNDED_CLIENT_MESSAGE, -} from "./lib/retrieval-health.mjs" +import { createRetrievalHealth, isUngrounded, UNGROUNDED_CLIENT_MESSAGE } from "./lib/retrieval-health.mjs" +const BACKEND = process.env.CHAT_BACKEND || "claude" +if (!["claude", "codex"].includes(BACKEND)) throw new Error("CHAT_BACKEND must be claude or codex") +const { query, tool, createSdkMcpServer } = BACKEND === "claude" + ? await import("@anthropic-ai/claude-agent-sdk") : {} const PORT = parseInt(process.env.AGENT_PORT || "30010", 10) const SEARCH_URL = process.env.PIXELRAG_SEARCH_URL || "https://api.pixelrag.ai" const MAX_BUDGET = parseFloat(process.env.CHAT_MAX_BUDGET_USD || "2.00") @@ -90,122 +89,6 @@ function log(...args) { console.log(new Date().toISOString(), ...args) } -// An unreachable index is not information for the model to work around: it -// ends the turn. Recorded on the turn's health so the stream can close as -// degraded instead of done. -function retrievalDown(health, onEvent, label, detail) { - recordRetrieval(health, RETRIEVAL_BACKEND_DOWN, detail) - onEvent("search_unavailable", { query: label, reason: detail }) - return { content: [{ type: "text", text: backendDownInstruction(detail) }], isError: true } -} - -function safeDecodeURIComponent(str) { - try { - return decodeURIComponent(str) - } catch { - return str - } -} - -function createTools(onEvent, uploadedImage, health) { - const searchTool = tool( - "pixelrag_search", - "Search the visual Wikipedia index by text, by the user's uploaded image, or BOTH combined. When the user uploaded an image, you MUST set use_uploaded_image=true AND provide a text query to get joint image+text retrieval — this gives the best results. Returns ranked results with article URLs, tile positions, and `pages` — the article's valid tile:chunk ranges (e.g. '0:0-7,1:0-4' = tile 0 has chunks 0-7, tile 1 has chunks 0-4). Use this first, then pixelrag_tile to view tiles.", - { - query: z.string().optional().describe("Natural language search query. Omit only when searching purely by an uploaded image."), - use_uploaded_image: z.boolean().optional().describe("Set true to include the user's uploaded image in the search (visual similarity). ALWAYS combine with a text query for best results — set this AND provide a query string in the same call."), - n_results: z.number().int().min(1).max(20).optional().describe("Number of results (default 5)"), - }, - async (args) => { - if (args.use_uploaded_image && !uploadedImage) { - return { content: [{ type: "text", text: "No image was uploaded in this conversation — use a text query instead." }] } - } - const searchByImage = Boolean(args.use_uploaded_image && uploadedImage) - if (!searchByImage && !args.query) { - return { content: [{ type: "text", text: "Provide a text query, or set use_uploaded_image:true when the user uploaded an image." }] } - } - // Text and image can be combined in one query for joint image+text retrieval. - const queryObj = {} - if (searchByImage) queryObj.image = uploadedImage - if (args.query) queryObj.text = args.query - const label = searchByImage && args.query ? `${args.query} + uploaded image` : args.query || "uploaded image" - onEvent("searching", { query: label }) - let resp - try { - resp = await fetch(`${SEARCH_URL}/search`, { - method: "POST", - headers: { "Content-Type": "application/json", "X-Source": "chat" }, - body: JSON.stringify({ queries: [queryObj], n_docs: args.n_results ?? 5, articles_only: true }), - signal: AbortSignal.timeout(30000), - }) - } catch (err) { - return retrievalDown(health, onEvent, label, String(err)) - } - const outcome = classifyStatus(resp.status) - if (outcome === RETRIEVAL_BACKEND_DOWN) { - return retrievalDown(health, onEvent, label, `HTTP ${resp.status}`) - } - if (outcome !== RETRIEVAL_OK) { - recordRetrieval(health, outcome, `HTTP ${resp.status}`) - return { content: [{ type: "text", text: `Search API error: ${resp.status}` }] } - } - recordRetrieval(health, RETRIEVAL_OK) - const data = await resp.json() - const hits = data.results?.[0]?.hits ?? [] - const results = hits.map((h) => { - const slug = h.url.includes("/wiki/") ? h.url.split("/wiki/").pop() : h.url - return { - title: safeDecodeURIComponent(slug || "").replace(/_/g, " "), - url: h.url.startsWith("http") ? h.url : `https://en.wikipedia.org/wiki/${slug}`, - score: Math.round(h.score * 1000) / 1000, - article_id: h.article_id, - tile_index: h.tile_index, - chunk_index: h.chunk_index, - pages: h.article_pages, - } - }) - onEvent("search_results", { query: label, hits }) - return { - content: [{ type: "text", text: JSON.stringify({ query: label, results, count: results.length }, null, 2) }], - } - } - ) - - const tileTool = tool( - "pixelrag_tile", - "View a Wikipedia screenshot tile by its coordinates. Returns the tile as an image so you can read the visual content. Only request coordinates within the article's `pages` ranges from search results (e.g. pages '0:0-7,1:0-4' means tile 1 ends at chunk 4) — coordinates beyond them do not exist.", - { - article_id: z.number().int().describe("Article ID from search results"), - tile_index: z.number().int().describe("Tile index from search results"), - chunk_index: z.number().int().describe("Chunk index from search results"), - }, - async (args) => { - const tileUrl = `${SEARCH_URL}/tile/${args.article_id}/${args.tile_index}/${args.chunk_index}` - try { - const resp = await fetch(tileUrl, { signal: AbortSignal.timeout(30000) }) - // The agent pages through articles by guessing chunk coordinates, so - // 404s are normal exploration — only surface tiles that actually load, - // otherwise the chat gallery renders broken images. - if (!resp.ok) { - if (classifyStatus(resp.status) === RETRIEVAL_BACKEND_DOWN) { - recordRetrieval(health, RETRIEVAL_BACKEND_DOWN, `tile HTTP ${resp.status}`) - } - return { content: [{ type: "text", text: `Tile not found: ${resp.status}` }] } - } - onEvent("viewing_tile", { article_id: args.article_id, tile_index: args.tile_index, chunk_index: args.chunk_index }) - const buffer = await resp.arrayBuffer() - const base64 = Buffer.from(buffer).toString("base64") - const mimeType = resp.headers.get("content-type") || "image/png" - return { content: [{ type: "image", data: base64, mimeType }] } - } catch (err) { - recordRetrieval(health, classifyThrow(err), `tile ${err}`) - return { content: [{ type: "text", text: `Failed to fetch tile: ${err}` }] } - } - } - ) - - return [searchTool, tileTool] -} function sse(event, data) { return `event: ${event}\ndata: ${JSON.stringify(data)}\n\n` @@ -220,7 +103,7 @@ const server = http.createServer(async (req, res) => { if (req.method === "GET" && req.url === "/health") { res.writeHead(200, { "Content-Type": "application/json" }) - res.end(JSON.stringify({ status: "ok" })) + res.end(JSON.stringify({ status: "ok", backend: BACKEND })) return } @@ -316,40 +199,47 @@ const server = http.createServer(async (req, res) => { const send = (event, data) => res.write(sse(event, data)) const uploadedImage = last?.image && typeof last.image === "string" ? last.image : null const health = createRetrievalHealth() - const tools = createTools(send, uploadedImage, health) - const mcpServer = createSdkMcpServer({ name: "pixelrag", version: "1.0.0", tools }) + const tools = createTools(send, uploadedImage, health, BACKEND === "codex" ? codexTool : tool, SEARCH_URL) + const controller = new AbortController() + res.on("close", () => controller.abort()) inFlight++ let sentText = false try { - for await (const message of query({ - prompt, - options: { - systemPrompt: SYSTEM_PROMPT, - mcpServers: { pixelrag: mcpServer }, - allowedTools: ["mcp__pixelrag__pixelrag_search", "mcp__pixelrag__pixelrag_tile"], - maxTurns: MAX_TURNS, - maxBudgetUsd: MAX_BUDGET, - maxThinkingTokens: THINKING_TOKENS, - includePartialMessages: true, - model: "sonnet", - }, - })) { - // Stream extended-thinking deltas (Claude Code-style reasoning trace) - if (message.type === "stream_event") { - const ev = message.event - if (ev?.type === "content_block_delta" && ev.delta?.type === "thinking_delta") { - send("thinking", { text: ev.delta.thinking }) + if (BACKEND === "codex") { + await runCodex({ textPrompt, uploadedImage, systemPrompt: SYSTEM_PROMPT, tools, send, signal: controller.signal, maxToolCalls: MAX_TURNS }) + sentText = true + } else { + const mcpServer = createSdkMcpServer({ name: "pixelrag", version: "1.0.0", tools }) + for await (const message of query({ + prompt, + options: { + systemPrompt: SYSTEM_PROMPT, + mcpServers: { pixelrag: mcpServer }, + allowedTools: ["mcp__pixelrag__pixelrag_search", "mcp__pixelrag__pixelrag_tile"], + maxTurns: MAX_TURNS, + maxBudgetUsd: MAX_BUDGET, + maxThinkingTokens: THINKING_TOKENS, + includePartialMessages: true, + model: "sonnet", + }, + })) { + // Stream extended-thinking deltas (Claude Code-style reasoning trace) + if (message.type === "stream_event") { + const ev = message.event + if (ev?.type === "content_block_delta" && ev.delta?.type === "thinking_delta") { + send("thinking", { text: ev.delta.thinking }) + } + continue } - continue - } - if (message.type === "assistant" && message.message) { - for (const block of message.message.content) { - if (block.type === "text" && block.text) { send("text", { text: block.text }); sentText = true } + if (message.type === "assistant" && message.message) { + for (const block of message.message.content) { + if (block.type === "text" && block.text) { send("text", { text: block.text }); sentText = true } + } + } + if (message.type === "result" && message.subtype === "success" && !sentText) { + send("text", { text: message.result }) } - } - if (message.type === "result" && message.subtype === "success" && !sentText) { - send("text", { text: message.result }) } } const elapsed = ((Date.now() - t0) / 1000).toFixed(1) @@ -387,5 +277,5 @@ const server = http.createServer(async (req, res) => { }) server.listen(PORT, () => { - log(`PixelRAG agent server on :${PORT} → search ${SEARCH_URL}, budget $${MAX_BUDGET}/conv`) + log(`PixelRAG ${BACKEND} agent server on :${PORT} → search ${SEARCH_URL}, ${BACKEND === "codex" ? "Codex timeout " + (process.env.CHAT_CODEX_TIMEOUT_MS || 180000) + "ms" : "budget $" + MAX_BUDGET + "/conv"}`) }) diff --git a/web/lib/agent-tools.mjs b/web/lib/agent-tools.mjs new file mode 100644 index 0000000..d947223 --- /dev/null +++ b/web/lib/agent-tools.mjs @@ -0,0 +1,119 @@ +import { z } from "zod" +import { RETRIEVAL_OK, RETRIEVAL_BACKEND_DOWN, classifyStatus, classifyThrow, recordRetrieval, backendDownInstruction } from "./retrieval-health.mjs" + +// An unreachable index is not information for the model to work around: it +// ends the turn. Recorded on the turn's health so the stream can close as +// degraded instead of done. +function retrievalDown(health, onEvent, label, detail) { + recordRetrieval(health, RETRIEVAL_BACKEND_DOWN, detail) + onEvent("search_unavailable", { query: label, reason: detail }) + return { content: [{ type: "text", text: backendDownInstruction(detail) }], isError: true } +} + +function safeDecodeURIComponent(str) { + try { + return decodeURIComponent(str) + } catch { + return str + } +} + +export function createTools(onEvent, uploadedImage, health, tool, SEARCH_URL) { + const searchTool = tool( + "pixelrag_search", + "Search the visual Wikipedia index by text, by the user's uploaded image, or BOTH combined. When the user uploaded an image, you MUST set use_uploaded_image=true AND provide a text query to get joint image+text retrieval — this gives the best results. Returns ranked results with article URLs, tile positions, and `pages` — the article's valid tile:chunk ranges (e.g. '0:0-7,1:0-4' = tile 0 has chunks 0-7, tile 1 has chunks 0-4). Use this first, then pixelrag_tile to view tiles.", + { + query: z.string().optional().describe("Natural language search query. Omit only when searching purely by an uploaded image."), + use_uploaded_image: z.boolean().optional().describe("Set true to include the user's uploaded image in the search (visual similarity). ALWAYS combine with a text query for best results — set this AND provide a query string in the same call."), + n_results: z.number().int().min(1).max(20).optional().describe("Number of results (default 5)"), + }, + async (args) => { + if (args.use_uploaded_image && !uploadedImage) { + return { content: [{ type: "text", text: "No image was uploaded in this conversation — use a text query instead." }] } + } + const searchByImage = Boolean(args.use_uploaded_image && uploadedImage) + if (!searchByImage && !args.query) { + return { content: [{ type: "text", text: "Provide a text query, or set use_uploaded_image:true when the user uploaded an image." }] } + } + // Text and image can be combined in one query for joint image+text retrieval. + const queryObj = {} + if (searchByImage) queryObj.image = uploadedImage + if (args.query) queryObj.text = args.query + const label = searchByImage && args.query ? `${args.query} + uploaded image` : args.query || "uploaded image" + onEvent("searching", { query: label }) + let resp + try { + resp = await fetch(`${SEARCH_URL}/search`, { + method: "POST", + headers: { "Content-Type": "application/json", "X-Source": "chat" }, + body: JSON.stringify({ queries: [queryObj], n_docs: args.n_results ?? 5, articles_only: true }), + signal: AbortSignal.timeout(30000), + }) + } catch (err) { + return retrievalDown(health, onEvent, label, String(err)) + } + const outcome = classifyStatus(resp.status) + if (outcome === RETRIEVAL_BACKEND_DOWN) { + return retrievalDown(health, onEvent, label, `HTTP ${resp.status}`) + } + if (outcome !== RETRIEVAL_OK) { + recordRetrieval(health, outcome, `HTTP ${resp.status}`) + return { content: [{ type: "text", text: `Search API error: ${resp.status}` }] } + } + recordRetrieval(health, RETRIEVAL_OK) + const data = await resp.json() + const hits = data.results?.[0]?.hits ?? [] + const results = hits.map((h) => { + const slug = h.url.includes("/wiki/") ? h.url.split("/wiki/").pop() : h.url + return { + title: safeDecodeURIComponent(slug || "").replace(/_/g, " "), + url: h.url.startsWith("http") ? h.url : `https://en.wikipedia.org/wiki/${slug}`, + score: Math.round(h.score * 1000) / 1000, + article_id: h.article_id, + tile_index: h.tile_index, + chunk_index: h.chunk_index, + pages: h.article_pages, + } + }) + onEvent("search_results", { query: label, hits }) + return { + content: [{ type: "text", text: JSON.stringify({ query: label, results, count: results.length }, null, 2) }], + } + } + ) + + const tileTool = tool( + "pixelrag_tile", + "View a Wikipedia screenshot tile by its coordinates. Returns the tile as an image so you can read the visual content. Only request coordinates within the article's `pages` ranges from search results (e.g. pages '0:0-7,1:0-4' means tile 1 ends at chunk 4) — coordinates beyond them do not exist.", + { + article_id: z.number().int().describe("Article ID from search results"), + tile_index: z.number().int().describe("Tile index from search results"), + chunk_index: z.number().int().describe("Chunk index from search results"), + }, + async (args) => { + const tileUrl = `${SEARCH_URL}/tile/${args.article_id}/${args.tile_index}/${args.chunk_index}` + try { + const resp = await fetch(tileUrl, { signal: AbortSignal.timeout(30000) }) + // The agent pages through articles by guessing chunk coordinates, so + // 404s are normal exploration — only surface tiles that actually load, + // otherwise the chat gallery renders broken images. + if (!resp.ok) { + if (classifyStatus(resp.status) === RETRIEVAL_BACKEND_DOWN) { + recordRetrieval(health, RETRIEVAL_BACKEND_DOWN, `tile HTTP ${resp.status}`) + } + return { content: [{ type: "text", text: `Tile not found: ${resp.status}` }] } + } + onEvent("viewing_tile", { article_id: args.article_id, tile_index: args.tile_index, chunk_index: args.chunk_index }) + const buffer = await resp.arrayBuffer() + const base64 = Buffer.from(buffer).toString("base64") + const mimeType = resp.headers.get("content-type") || "image/png" + return { content: [{ type: "image", data: base64, mimeType }] } + } catch (err) { + recordRetrieval(health, classifyThrow(err), `tile ${err}`) + return { content: [{ type: "text", text: `Failed to fetch tile: ${err}` }] } + } + } + ) + + return [searchTool, tileTool] +} diff --git a/web/lib/codex-backend.mjs b/web/lib/codex-backend.mjs new file mode 100644 index 0000000..bb879f3 --- /dev/null +++ b/web/lib/codex-backend.mjs @@ -0,0 +1,118 @@ +import http from "node:http" +import { spawn } from "node:child_process" +import { randomBytes } from "node:crypto" +import { mkdtemp, writeFile, rm } from "node:fs/promises" +import { tmpdir } from "node:os" +import { join } from "node:path" +import { fileURLToPath } from "node:url" +import readline from "node:readline" +import { z } from "zod" + +export function codexTool(name, description, shape, handler) { + const schema = z.object(shape) + return { name, description, inputSchema: z.toJSONSchema(schema), annotations: { readOnlyHint: true, destructiveHint: false, openWorldHint: false }, call: (args) => handler(schema.parse(args)) } +} + +export async function runCodex({ textPrompt, uploadedImage, systemPrompt, tools, send, signal, maxToolCalls = 12, onDiagnostic = () => {} }) { + const dir = await mkdtemp(join(tmpdir(), "pixelrag-codex-")) + const token = randomBytes(32).toString("hex") + let calls = 0 + let child + let timedOut = false + let exited = false + const bridge = http.createServer(async (req, res) => { + const respond = (status, data) => { res.writeHead(status, { "Content-Type": "application/json" }); res.end(JSON.stringify(data)) } + if (req.method !== "POST" || req.headers.authorization !== `Bearer ${token}`) return respond(403, {}) + try { + let body = "" + for await (const chunk of req) { + body += chunk + if (body.length > 65536) return respond(413, {}) + } + const request = JSON.parse(body) + if (request.method === "tools/list") return respond(200, { tools: tools.map(({ name, description, inputSchema, annotations }) => ({ name, description, inputSchema, annotations })) }) + if (request.method !== "tools/call") return respond(400, {}) + if (++calls > maxToolCalls) return respond(200, { isError: true, content: [{ type: "text", text: "Tool limit reached. Finish the answer using the tiles already read." }] }) + const tool = tools.find((t) => t.name === request.params?.name) + if (!tool) throw new Error("Unknown tool") + respond(200, await tool.call(request.params.arguments || {})) + } catch (error) { + respond(200, { isError: true, content: [{ type: "text", text: String(error) }] }) + } + }) + let timer + let forceKill + const cancel = () => { + if (child && !exited) { + child.kill("SIGTERM") + if (!forceKill) forceKill = setTimeout(() => { if (!exited) child.kill("SIGKILL") }, 2000) + } + } + try { + await new Promise((resolve, reject) => { bridge.once("error", reject); bridge.listen(0, "127.0.0.1", resolve) }) + const endpoint = `http://127.0.0.1:${bridge.address().port}` + const args = ["exec", "--ignore-user-config", "--ephemeral", "--skip-git-repo-check", "--json", "--sandbox", "read-only", "--cd", dir, + "--disable", "shell_tool", "--disable", "code_mode", "--enable", "code_mode_host", "--disable", "multi_agent", "--disable", "apps", "--disable", "plugins", "--disable", "view_image", "--disable", "browser_use_external", "--disable", "browser_use_full_cdp_access", "--disable", "browser_use", "--disable", "computer_use", "--disable", "image_generation", + "-c", "approval_policy=\"never\"", "-c", "web_search=\"disabled\"", + "-c", `developer_instructions=${JSON.stringify(systemPrompt)}`, + "-c", `mcp_servers.pixelrag.command=${JSON.stringify(process.execPath)}`, + "-c", `mcp_servers.pixelrag.args=${JSON.stringify([fileURLToPath(new URL("./codex-mcp.mjs", import.meta.url)), endpoint, token])}`, + "-c", "mcp_servers.pixelrag.required=true", + "-c", 'mcp_servers.pixelrag.tools.pixelrag_search.approval_mode="approve"', + "-c", 'mcp_servers.pixelrag.tools.pixelrag_tile.approval_mode="approve"'] + if (process.env.CHAT_CODEX_MODEL) args.push("--model", process.env.CHAT_CODEX_MODEL) + if (uploadedImage) { + const match = uploadedImage.match(/^data:(image\/(?:png|jpeg|webp|gif));base64,(.+)$/is) + if (!match) throw new Error("Unsupported uploaded image; use a PNG, JPEG, WebP, or GIF data URL") + const path = join(dir, "uploaded-image") + await writeFile(path, Buffer.from(match[2], "base64"), { mode: 0o600 }) + args.push("--image", path) + } + args.push("-") + if (signal?.aborted) throw new Error("Chat cancelled") + child = spawn(process.env.CODEX_BIN || "codex", args, { cwd: dir, stdio: ["pipe", "pipe", "pipe"] }) + let diagnostics = "" + child.stderr.on("data", (chunk) => { + const safe = String(chunk).replaceAll(token, "[redacted]") + diagnostics = (diagnostics + safe).slice(-4000) + onDiagnostic(safe) + }) + const completed = new Promise((resolve, reject) => { + child.once("error", reject) + child.once("close", (code) => { + exited = true + clearTimeout(forceKill) + code === 0 ? resolve() : reject(new Error(signal?.aborted ? "Chat cancelled" : timedOut ? "Codex request timed out" : `Codex failed (${code}): ${diagnostics}`)) + }) + }) + // Attach a rejection handler immediately, before consuming stdout. + completed.catch(() => {}) + signal?.addEventListener("abort", cancel, { once: true }) + timer = setTimeout(() => { timedOut = true; cancel() }, Number(process.env.CHAT_CODEX_TIMEOUT_MS || 180000)) + child.stdin.on("error", () => {}) + child.stdin.end(textPrompt) + let sentText = false + let failure + for await (const line of readline.createInterface({ input: child.stdout })) { + let event + try { event = JSON.parse(line) } catch { continue } + if (event.type === "item.completed" && event.item?.type === "agent_message" && event.item.text) { + send("text", { text: event.item.text }); sentText = true + } else if (event.type === "item.completed" && event.item?.type === "reasoning" && event.item.text) { + send("thinking", { text: event.item.text }) + } else if (event.type === "turn.failed" || event.type === "error") { + failure = event.error?.message || event.message || "Codex turn failed" + } + } + await completed + if (failure) throw new Error(failure) + if (!sentText) throw new Error("Codex returned no answer") + } finally { + clearTimeout(timer) + signal?.removeEventListener("abort", cancel) + cancel() + bridge.closeAllConnections() + await new Promise((resolve) => bridge.close(resolve)) + await rm(dir, { recursive: true, force: true }) + } +} diff --git a/web/lib/codex-backend.test.mjs b/web/lib/codex-backend.test.mjs new file mode 100644 index 0000000..cb1df5c --- /dev/null +++ b/web/lib/codex-backend.test.mjs @@ -0,0 +1,101 @@ +import test from "node:test" +import assert from "node:assert/strict" +import { mkdtemp, writeFile, rm } from "node:fs/promises" +import { join } from "node:path" +import { tmpdir } from "node:os" +import { z } from "zod" +import { runCodex, codexTool } from "./codex-backend.mjs" + +// A fake CLI uses the actual stdio adapter, making this test cover the MCP +// protocol, tool validation, and SSE event conversion without model calls. +const fakeCli = ` +import {spawn} from 'node:child_process'; +import readline from 'node:readline'; +import assert from 'node:assert/strict'; +import {readFile} from 'node:fs/promises'; +const args=process.argv.slice(2); +assert(args.includes('read-only')); +assert(args.includes('shell_tool')); +if(args.includes('--image')) assert((await readFile(args[args.indexOf('--image')+1])).length>0); +const config=args.find(x=>x.startsWith('mcp_servers.pixelrag.args=')); +const bridgeArgs=JSON.parse(config.slice(config.indexOf('=')+1)); +const child=spawn(process.execPath,bridgeArgs,{stdio:['pipe','pipe','inherit']}); +const lines=readline.createInterface({input:child.stdout}); +const iterator=lines[Symbol.asyncIterator](); +let id=0; +async function call(method,params={}) { + child.stdin.write(JSON.stringify({jsonrpc:'2.0',id:++id,method,params})+'\\n'); + return JSON.parse((await iterator.next()).value).result; +} +try { + const init=await call('initialize',{protocolVersion:'2024-11-05'}); + assert.equal(init.serverInfo.name,'pixelrag'); + child.stdin.write(JSON.stringify({jsonrpc:'2.0',method:'notifications/initialized'})+'\\n'); + const list=await call('tools/list'); + assert.equal(list.tools[0].name,'pixelrag_search'); + const result=await call('tools/call',{name:'pixelrag_search',arguments:{query:'Taj Mahal'}}); + console.log(JSON.stringify({type:'item.completed',item:{type:'reasoning',text:'Reading tiles'}})); + console.log(JSON.stringify({type:'item.completed',item:{type:'agent_message',text:JSON.stringify(result)}})); + if(process.env.PIXELRAG_TEST_FAILURE) console.log(JSON.stringify({type:'turn.failed',error:{message:'test failure'}})); +} finally {child.stdin.end();lines.close();} +` + +test("Codex CLI adapter calls MCP tools and emits SSE-compatible events", async () => { + const dir = await mkdtemp(join(tmpdir(), "codex-test-")) + const cli = join(dir, "cli.mjs") + const previous = process.env.CODEX_BIN + try { + await writeFile(cli, `#!${process.execPath}\n${fakeCli}`, { mode: 0o700 }) + process.env.CODEX_BIN = cli + const events = [] + const options = { + textPrompt: "What is the Taj Mahal?", systemPrompt: "Read tiles.", + tools: [codexTool("pixelrag_search", "Search", { query: z.string() }, async (args) => { + events.push(["searching", args]) + return { content: [{ type: "text", text: "Retrieved " + args.query }] } + })], + send: (event, data) => events.push([event, data]), + } + await runCodex(options) + assert.equal(events[0][0], "searching") + assert.equal(events[1][0], "thinking") + assert.match(events[2][1].text, /Retrieved Taj Mahal/) + process.env.PIXELRAG_TEST_FAILURE = "1" + await assert.rejects(runCodex(options), /test failure/) + delete process.env.PIXELRAG_TEST_FAILURE + const controller = new AbortController() + controller.abort() + await assert.rejects(runCodex({ ...options, signal: controller.signal }), /cancelled/) + await assert.rejects(runCodex({ ...options, uploadedImage: "invalid image" }), /Unsupported uploaded image/) + await runCodex({ ...options, uploadedImage: "data:image/png;base64,iVBORw0KGgo=" }) + events.length = 0 + await runCodex({ ...options, maxToolCalls: 0 }) + assert.match(events.at(-1)[1].text, /Tool limit reached/) + assert.equal(events.filter(([event]) => event === "searching").length, 0) + // Exercise cancellation of an already-running subprocess and the timeout. + await writeFile(cli, `#!${process.execPath}\nsetInterval(() => {}, 1000)`, { mode: 0o700 }) + const running = new AbortController() + const cancelTimer = setTimeout(() => running.abort(), 100) + await assert.rejects(runCodex({ ...options, signal: running.signal }), /cancelled/) + clearTimeout(cancelTimer) + process.env.CHAT_CODEX_TIMEOUT_MS = "100" + await assert.rejects(runCodex(options), /timed out/) + delete process.env.CHAT_CODEX_TIMEOUT_MS + process.env.CODEX_BIN = join(dir, "nonexistent-cli") + await assert.rejects(runCodex(options), /ENOENT/) + } finally { + if (previous === undefined) delete process.env.CODEX_BIN + else process.env.CODEX_BIN = previous + delete process.env.PIXELRAG_TEST_FAILURE + delete process.env.CHAT_CODEX_TIMEOUT_MS + await rm(dir, { recursive: true, force: true }) + } +}) + +test("MCP schemas reject invalid tile coordinates before invoking retrieval", async () => { + let called = false + const tool = codexTool("pixelrag_tile", "Read tile", { article_id: z.number().int() }, async () => { called = true }) + assert.equal(tool.inputSchema.type, "object") + assert.throws(() => tool.call({ article_id: "wrong" })) + assert.equal(called, false) +}) diff --git a/web/lib/codex-mcp.mjs b/web/lib/codex-mcp.mjs new file mode 100644 index 0000000..d9dc95d --- /dev/null +++ b/web/lib/codex-mcp.mjs @@ -0,0 +1,29 @@ +// A stdio MCP adapter for a single conversation's local tool endpoint. +// No filesystem, shell, or arbitrary URL tools are exposed to the model. +import readline from "node:readline" + +const [endpoint, token] = process.argv.slice(2) +const reply = (id, result) => process.stdout.write(JSON.stringify({ jsonrpc: "2.0", id, result }) + "\n") +for await (const line of readline.createInterface({ input: process.stdin })) { + let request + try { + request = JSON.parse(line) + if (request.id === undefined) continue + if (request.method === "initialize") { + reply(request.id, { protocolVersion: request.params.protocolVersion, capabilities: { tools: {} }, serverInfo: { name: "pixelrag", version: "1.0.0" } }) + } else if (request.method === "ping") { + reply(request.id, {}) + } else if (["tools/list", "tools/call"].includes(request.method)) { + const response = await fetch(endpoint, { + method: "POST", headers: { "Content-Type": "application/json", Authorization: `Bearer ${token}` }, + body: JSON.stringify(request), signal: AbortSignal.timeout(65000), + }) + if (!response.ok) throw new Error(`Tool endpoint HTTP ${response.status}`) + reply(request.id, await response.json()) + } else { + process.stdout.write(JSON.stringify({ jsonrpc: "2.0", id: request.id, error: { code: -32601, message: "Method not found" } }) + "\n") + } + } catch (error) { + if (request?.id !== undefined) process.stdout.write(JSON.stringify({ jsonrpc: "2.0", id: request.id, error: { code: -32603, message: String(error) } }) + "\n") + } +}