import express from "express"; import path from "node:path"; import { fileURLToPath } from "node:url"; import { Store } from "./store.js"; import type { UIEvent } from "./store.js"; import { handleMcpRequest } from "./mcp-handler.js"; import { generatePrompts } from "./prompts.js"; import type { AgentId } from "./store.js"; import { getWorkflow, isValidWorkflowId, listValidWorkflowIds, summarizeWorkflows, type WorkflowId, } from "./workflows.js"; // ── Port ─────────────────────────────────────────────────────────── const port = (() => { const idx = process.argv.indexOf("--port"); if (idx !== -1 && process.argv[idx + 1]) return parseInt(process.argv[idx + 1], 10); if (process.env.PORT) return parseInt(process.env.PORT, 10); return 4100; })(); // ── App ──────────────────────────────────────────────────────────── const app = express(); app.use(express.json()); const store = new Store(); const __dirname = path.dirname(fileURLToPath(import.meta.url)); // ── MCP endpoints ────────────────────────────────────────────────── app.post("/mcp/claude", async (req, res, next) => { try { await handleMcpRequest("claude", store, req, res); } catch (err) { next(err); } }); app.post("/mcp/codex", async (req, res, next) => { try { await handleMcpRequest("codex", store, req, res); } catch (err) { next(err); } }); // ── REST API ─────────────────────────────────────────────────────── app.get("/api/workflows", (_req, res) => { res.json(summarizeWorkflows()); }); app.get("/api/sessions", (_req, res) => { res.json(store.listSessions()); }); app.post("/api/sessions", (req, res) => { const { goal, initiator, workflow } = req.body as { goal?: string; initiator?: string; workflow?: string; }; if (!goal || typeof goal !== "string") { res.status(400).json({ error: "goal is required" }); return; } let workflowId: WorkflowId = "custom"; if (workflow !== undefined) { if (!isValidWorkflowId(workflow)) { res.status(400).json({ error: `unknown workflow: ${String(workflow)}`, validWorkflows: listValidWorkflowIds(), }); return; } workflowId = workflow; } const template = getWorkflow(workflowId); let effectiveInitiator: AgentId; if (template.forcedInitiator) { effectiveInitiator = template.forcedInitiator; } else { if (initiator !== "claude" && initiator !== "codex") { res.status(400).json({ error: 'initiator must be "claude" or "codex"' }); return; } effectiveInitiator = initiator; } const session = store.createSession(goal, effectiveInitiator, workflowId); res.status(201).json(session); }); app.get("/api/sessions/:id", (req, res) => { const session = store.getSession(req.params.id); if (!session) { res.status(404).json({ error: "session not found" }); return; } res.json(session); }); app.patch("/api/sessions/:id", (req, res) => { const updates = req.body as { autoForward?: boolean; goal?: string; currentPhaseIndex?: number; }; if (updates.currentPhaseIndex !== undefined) { const idx = updates.currentPhaseIndex; if (!Number.isInteger(idx) || idx < 0) { res.status(400).json({ error: "currentPhaseIndex must be a non-negative integer" }); return; } const existing = store.getSession(req.params.id); if (!existing) { res.status(404).json({ error: "session not found" }); return; } const template = getWorkflow(existing.workflow); if (template.phases.length === 0) { res.status(400).json({ error: "this workflow has no phases" }); return; } if (idx >= template.phases.length) { res.status(400).json({ error: `currentPhaseIndex out of range (valid: 0..${template.phases.length - 1})`, }); return; } } const session = store.updateSession(req.params.id, updates); if (!session) { res.status(404).json({ error: "session not found" }); return; } res.json(session); }); app.delete("/api/sessions/:id", (req, res) => { if (!store.deleteSession(req.params.id)) { res.status(404).json({ error: "session not found" }); return; } res.status(204).end(); }); app.post("/api/sessions/:id/close", (req, res) => { const { reason } = req.body as { reason?: string }; const session = store.closeSession(req.params.id, "user", reason); if (!session) { res.status(404).json({ error: "session not found" }); return; } res.json(session); }); // ── Message actions ──────────────────────────────────────────────── function parseMid(raw: string, res: express.Response): number | null { const n = parseInt(raw, 10); if (isNaN(n)) { res.status(400).json({ error: "invalid message id" }); return null; } return n; } app.post("/api/sessions/:id/messages/:mid/deliver", (req, res) => { const mid = parseMid(req.params.mid, res); if (mid === null) return; const msg = store.deliverMessage(req.params.id, mid); if (!msg) { res.status(404).json({ error: "message not found or already delivered" }); return; } res.json(msg); }); app.patch("/api/sessions/:id/messages/:mid", (req, res) => { const mid = parseMid(req.params.mid, res); if (mid === null) return; const { body: newBody } = req.body as { body?: string }; if (!newBody || typeof newBody !== "string") { res.status(400).json({ error: "body is required" }); return; } const msg = store.editMessage(req.params.id, mid, newBody); if (!msg) { res.status(404).json({ error: "message not found or not pending" }); return; } res.json(msg); }); app.delete("/api/sessions/:id/messages/:mid", (req, res) => { const mid = parseMid(req.params.mid, res); if (mid === null) return; if (!store.rejectMessage(req.params.id, mid)) { res.status(404).json({ error: "message not found or not pending" }); return; } res.status(204).end(); }); // ── Prompts ──────────────────────────────────────────────────────── app.get("/api/sessions/:id/prompts", (req, res) => { const session = store.getSession(req.params.id); if (!session) { res.status(404).json({ error: "session not found" }); return; } const proto = req.protocol; const host = req.get("host") || `localhost:${port}`; const baseUrl = `${proto}://${host}`; res.json(generatePrompts(session, baseUrl)); }); // ── SSE for web UI ───────────────────────────────────────────────── app.get("/events", (req, res) => { res.writeHead(200, { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", Connection: "keep-alive", }); res.write(`data: ${JSON.stringify({ type: "connected" })}\n\n`); const onUpdate = (event: UIEvent) => { res.write(`data: ${JSON.stringify(event)}\n\n`); }; store.on("ui-update", onUpdate); req.on("close", () => { store.removeListener("ui-update", onUpdate); }); }); // ── Static ───────────────────────────────────────────────────────── app.get("/", (_req, res) => { res.sendFile(path.resolve(__dirname, "../web/index.html")); }); // ── Error handler ────────────────────────────────────────────────── app.use((err: unknown, _req: express.Request, res: express.Response, _next: express.NextFunction) => { const message = err instanceof Error ? err.message : "Internal server error"; res.status(500).json({ error: message }); }); // ── Start ────────────────────────────────────────────────────────── const server = app.listen(port, "127.0.0.1", () => { console.log(`collab-bridge running at http://localhost:${port}`); console.log(` Web UI: http://localhost:${port}/`); console.log(` Claude MCP: http://localhost:${port}/mcp/claude`); console.log(` Codex MCP: http://localhost:${port}/mcp/codex`); }); // collab_wait_for_reply long-polls for up to 5 minutes (300s). // Ensure the HTTP socket idle timeout doesn't kill it. // (Default is 0 on Node 13+, but was 120s on older versions.) server.timeout = 0;