repo2 / api /index.js
agent
port opencode-proxy into repo2
f5eeff1
Raw
History Blame Contribute Delete
12.2 kB
import express from "express";
import { createOpencodeClient } from "@opencode-ai/sdk";
import { spawn } from "child_process";
import { fileURLToPath } from "url";
import path from "path";
const __dirname = path.dirname(fileURLToPath(import.meta.url));
process.env.PATH = `${path.join(__dirname, "..", "bin")}:${process.env.PATH ?? ""}`;
process.env.HOME = "/tmp";
if (process.env.OPENCODE_AUTH && !process.env.OPENCODE_AUTH_CONTENT) {
process.env.OPENCODE_AUTH_CONTENT = Buffer.from(process.env.OPENCODE_AUTH, "base64").toString("utf8");
}
const DEFAULT_MODEL = "opencode/big-pickle";
const POOL = parseInt(process.env.OPENCODE_WORKERS || "4", 10);
const L = (...a) => console.log(`[PROXY] ${new Date().toISOString()}`, ...a);
const E = (...a) => console.error(`[PROXY] ${new Date().toISOString()}`, ...a);
function toModel(m) {
if (!m) return undefined;
if (typeof m === "object") return m;
const i = m.indexOf("/");
return i === -1 ? { providerID: "opencode", modelID: m } : { providerID: m.slice(0, i), modelID: m.slice(i + 1) };
}
/* ── Worker pool ──────────────────────────────────────────────── */
let _workers = [];
let _next = 0;
let _ready = null;
async function _spawn(i) {
const port = 4096 + i;
const proc = spawn("opencode", ["serve", "--hostname=127.0.0.1", `--port=${port}`, "--print-logs"], {
env: { ...process.env, OPENCODE_CONFIG_CONTENT: "{}" },
});
let log = "";
let resolved = false;
const cap = (c) => { log += c.toString(); };
proc.stdout?.on("data", cap);
proc.stderr?.on("data", cap);
proc.on("spawn", () => L("pool: worker", i, "pid", proc.pid));
proc.on("exit", (code) => {
E("pool: worker", i, "exited code", code, "- respawning in 2s");
if (_workers[i]) _workers[i].alive = false;
setTimeout(() => _spawn(i), 2000).unref();
});
const url = await new Promise((resolve, reject) => {
const timer = setTimeout(() => { proc.kill("SIGKILL"); reject(new Error("timeout\n" + log.slice(-2000))); }, 90000);
const check = () => {
if (resolved) return;
const m = log.match(/on\s+(https?:\/\/[^\s]+)/);
if (m) { resolved = true; clearTimeout(timer); L("pool: worker", i, "ready at", m[1]); resolve(m[1]); }
};
proc.stdout?.on("data", check);
proc.stderr?.on("data", check);
proc.on("exit", (code) => { clearTimeout(timer); if (!resolved) reject(new Error(`exited ${code}\n` + log.slice(-2000))); });
proc.on("error", (e) => { clearTimeout(timer); reject(e); });
});
const client = createOpencodeClient({ baseUrl: url });
_workers[i] = { client, proc, index: i, alive: true };
L("pool: worker", i, "online");
}
async function _boot() {
L("pool: spawning", POOL, "workers");
await Promise.allSettled(Array.from({ length: POOL }, (_, i) => _spawn(i).catch(e => E("pool: worker", i, "spawn failed:", e.message))));
const n = _workers.filter(w => w?.alive).length;
L("pool: ready", n, "/", POOL);
}
async function _pool() {
if (!_ready) _ready = _boot();
await _ready;
}
async function getClient() {
await _pool();
for (let t = 0; t < POOL * 2; t++) {
const i = _next++ % POOL;
const w = _workers[i];
if (w && w.alive && !w.proc.killed) return w.client;
}
await new Promise(r => setTimeout(r, 1000));
return getClient();
}
async function withOpencode(fn) {
const client = await getClient();
try {
return await fn(client);
} catch (e) {
const isConnErr = e?.code === "ECONNREFUSED" || e?.code === "ECONNRESET" || (e?.message && e.message.includes("fetch failed"));
if (isConnErr) {
for (const w of _workers) {
if (w && w.client === client) { w.alive = false; break; }
}
L("withOpencode: retrying with another worker");
const client2 = await getClient();
return await fn(client2);
}
throw e;
}
}
/* ── Express ──────────────────────────────────────────────────── */
const app = express();
app.use(express.json({ limit: "10mb" }));
app.use((req, res, next) => {
res.setHeader("Access-Control-Allow-Origin", "*");
res.setHeader("Access-Control-Allow-Methods", "GET,POST,OPTIONS");
res.setHeader("Access-Control-Allow-Headers", req.headers["access-control-request-headers"] || "*");
res.setHeader("Access-Control-Max-Age", "86400");
if (req.method === "OPTIONS") return res.status(204).end();
next();
});
app.use((req, _res, next) => {
const size = req.headers["content-length"] ? `${req.headers["content-length"]}b` : "unknown";
L("request:", req.method, req.originalUrl, "from", req.headers.origin || req.ip, "body", size);
next();
});
app.get("/health", async (_req, res) => {
L("GET /health");
try { await withOpencode(() => ({ ok: true })); res.json({ ok: true }); }
catch (e) { E("GET /health:", e); res.status(500).json({ error: String(e) }); }
});
app.post("/prompt", async (req, res) => {
const { parts, model, sessionId, system, tools, agent, noReply } = req.body || {};
if (!parts) return res.status(400).json({ error: "parts required" });
L("POST /prompt:", model || DEFAULT_MODEL, "session", sessionId || "(new)");
try {
const out = await withOpencode(async (c) => {
const created = sessionId ? { data: { id: sessionId } } : await c.session.create({ body: {} });
const result = await c.session.prompt({
path: { id: created.data.id },
body: { parts, model: toModel(model || DEFAULT_MODEL), system, tools, agent, noReply },
});
return { sessionId: created.data.id, ok: !result.error, data: result.data, error: result.error };
});
res.json(out);
} catch (e) { E("POST /prompt:", e); res.status(500).json({ error: String(e) }); }
});
function openaiText(data) {
return (data?.parts ?? []).filter((p) => p.type === "text").map((p) => p.text).join("");
}
function openaiReasoning(data) {
return (data?.parts ?? []).filter((p) => p.type === "reasoning").map((p) => p.text).join("");
}
function openaiUsage(info) {
const t = info?.tokens ?? {};
return { prompt_tokens: t.input ?? 0, completion_tokens: t.output ?? 0, total_tokens: t.total ?? (t.input ?? 0) + (t.output ?? 0) };
}
app.post("/v1/chat/completions", async (req, res) => {
const { model, messages, stream, system: sysArg, sessionId } = req.body || {};
if (!messages) return res.status(400).json({ error: { message: "messages required", type: "invalid_request_error" } });
const flat = (c) => typeof c === "string" ? c : (c ?? []).map((x) => x.text ?? x).join("");
const system = [sysArg, ...messages.filter((m) => m.role === "system").map((m) => flat(m.content))].filter(Boolean).join("\n");
const convo = messages.filter((m) => m.role !== "system");
if (!convo.some((m) => m.role === "user")) return res.status(400).json({ error: { message: "no user message", type: "invalid_request_error" } });
const text = convo.length === 1
? flat(convo[0].content)
: convo.map((m) => `${m.role === "assistant" ? "Assistant" : "User"}: ${flat(m.content)}`).join("\n\n");
const parts = [{ type: "text", text }];
const modelId = model || DEFAULT_MODEL;
const id = "chatcmpl-" + Math.random().toString(36).slice(2);
const created = Date.now() / 1000 | 0;
L("POST /v1/chat/completions:", "model", modelId, "stream", Boolean(stream), "messages", messages.length);
if (!stream) {
try {
const out = await withOpencode(async (c) => {
const s = sessionId ? { data: { id: sessionId } } : await c.session.create({ body: {} });
const result = await c.session.prompt({ path: { id: s.data.id }, body: { parts, model: toModel(modelId), system } });
return result.error ? { error: result.error } : { data: result.data };
});
if (out.error) return res.status(500).json({ error: { message: JSON.stringify(out.error), type: "server_error" } });
const data = out.data, content = openaiText(data), reasoning = openaiReasoning(data), usage = openaiUsage(data.info);
L("chat: ok, chars", content.length);
const message = { role: "assistant", content, ...(reasoning ? { reasoning_content: reasoning } : {}) };
return res.json({ id, object: "chat.completion", created, model: modelId, choices: [{ index: 0, message, finish_reason: "stop" }], usage });
} catch (e) { E("chat:", e); return res.status(500).json({ error: { message: String(e), type: "server_error" } }); }
}
res.setHeader("Content-Type", "text/event-stream");
res.setHeader("Cache-Control", "no-cache");
res.setHeader("Connection", "keep-alive");
L("chat(stream): starting");
let client;
try { client = await getClient(); }
catch (e) { E("chat(stream): get client failed:", e); return res.status(500).json({ error: { message: String(e), type: "server_error" } }); }
let sentRole = false;
const chunk = (delta) => {
if (!sentRole) delta = { role: "assistant", ...delta };
sentRole = true;
res.write(`data: ${JSON.stringify({ id, object: "chat.completion.chunk", created, model: modelId, choices: [{ index: 0, delta }] })}\n\n`);
};
let sid;
try {
const made = sessionId ? { data: { id: sessionId } } : await client.session.create({ body: {} });
sid = made.data.id;
L("chat(stream): session", sid);
} catch (e) { E("chat(stream): session create failed:", e); return res.status(500).json({ error: { message: String(e), type: "server_error" } }); }
const ac = new AbortController();
req.on("close", () => { ac.abort(); L("chat(stream): client disconnected"); });
let readyResolve;
const ready = new Promise((r) => { readyResolve = r; });
let fullText = "";
let deltaCount = 0;
const pump = (async () => {
try {
const sse = await client.global.event({ signal: ac.signal });
for await (const ev of sse.stream) {
readyResolve();
const p = ev?.payload ?? ev;
const t = p?.type;
if (t === "message.part.delta") {
if (p.properties?.sessionID === sid && p.properties?.delta) {
if (p.properties.field === "text") chunk({ content: p.properties.delta });
else if (p.properties.field === "reasoning") chunk({ reasoning_content: p.properties.delta });
deltaCount++;
}
continue;
}
if (t === "message.part.updated") {
const part = p.properties?.part;
if (!part || part.sessionID !== sid || (part.type !== "text" && part.type !== "reasoning")) continue;
let delta = p.properties?.delta;
if (!delta) delta = (part.text ?? "").slice(fullText.length);
if (delta) { fullText += delta; chunk(part.type === "reasoning" ? { reasoning_content: delta } : { content: delta }); deltaCount++; }
continue;
}
if (t === "message.updated") {
const info = p.properties?.info;
if (info?.role === "assistant" && info?.finish) { ac.abort(); return; }
}
}
} catch (e) { E("chat(stream): pump error:", e); }
})();
await Promise.race([ready, new Promise((r) => setTimeout(r, 1000))]);
try {
await client.session.prompt({ path: { id: sid }, body: { parts, model: toModel(modelId), system } });
L("chat(stream): prompt sent");
} catch (e) {
E("chat(stream): prompt failed:", e);
chunk({ content: "\n[error] " + String(e) });
} finally {
await Promise.race([pump, new Promise((r) => setTimeout(r, 2000))]);
L("chat(stream): done, deltas", deltaCount, "chars", fullText.length);
ac.abort();
await pump.catch(() => {});
}
res.write(`data: ${JSON.stringify({ choices: [{ index: 0, delta: {}, finish_reason: "stop" }] })}\n\n`);
res.write("data: [DONE]\n\n");
res.end();
});
app.all("*", (req, res) => { L("404:", req.method, req.originalUrl); res.status(404).json({ error: "not found" }); });
app.use((err, _req, res, _next) => {
E("unhandled error:", err);
if (res.headersSent) return;
res.status(500).json({ error: { message: String(err?.message ?? err), type: "server_error" } });
});
export default app;