Skip to content
File

Blob: src/index.ts

typescript271 lines
1import express from "express";
2import path from "node:path";
3import { fileURLToPath } from "node:url";
4import { Store } from "./store.js";
5import type { UIEvent } from "./store.js";
6import { handleMcpRequest } from "./mcp-handler.js";
7import { generatePrompts } from "./prompts.js";
8import type { AgentId } from "./store.js";
9import {
10 getWorkflow,
11 isValidWorkflowId,
12 listValidWorkflowIds,
13 summarizeWorkflows,
14 type WorkflowId,
15} from "./workflows.js";
16 
17// ── Port ───────────────────────────────────────────────────────────
18 
19const port = (() => {
20 const idx = process.argv.indexOf("--port");
21 if (idx !== -1 && process.argv[idx + 1]) return parseInt(process.argv[idx + 1], 10);
22 if (process.env.PORT) return parseInt(process.env.PORT, 10);
23 return 4100;
24})();
25 
26// ── App ────────────────────────────────────────────────────────────
27 
28const app = express();
29app.use(express.json());
30 
31const store = new Store();
32 
33const __dirname = path.dirname(fileURLToPath(import.meta.url));
34 
35// ── MCP endpoints ──────────────────────────────────────────────────
36 
37app.post("/mcp/claude", async (req, res, next) => {
38 try {
39 await handleMcpRequest("claude", store, req, res);
40 } catch (err) {
41 next(err);
42 }
43});
44 
45app.post("/mcp/codex", async (req, res, next) => {
46 try {
47 await handleMcpRequest("codex", store, req, res);
48 } catch (err) {
49 next(err);
50 }
51});
52 
53// ── REST API ───────────────────────────────────────────────────────
54 
55app.get("/api/workflows", (_req, res) => {
56 res.json(summarizeWorkflows());
57});
58 
59app.get("/api/sessions", (_req, res) => {
60 res.json(store.listSessions());
61});
62 
63app.post("/api/sessions", (req, res) => {
64 const { goal, initiator, workflow } = req.body as {
65 goal?: string;
66 initiator?: string;
67 workflow?: string;
68 };
69 if (!goal || typeof goal !== "string") {
70 res.status(400).json({ error: "goal is required" });
71 return;
72 }
73 let workflowId: WorkflowId = "custom";
74 if (workflow !== undefined) {
75 if (!isValidWorkflowId(workflow)) {
76 res.status(400).json({
77 error: `unknown workflow: ${String(workflow)}`,
78 validWorkflows: listValidWorkflowIds(),
79 });
80 return;
81 }
82 workflowId = workflow;
83 }
84 const template = getWorkflow(workflowId);
85 let effectiveInitiator: AgentId;
86 if (template.forcedInitiator) {
87 effectiveInitiator = template.forcedInitiator;
88 } else {
89 if (initiator !== "claude" && initiator !== "codex") {
90 res.status(400).json({ error: 'initiator must be "claude" or "codex"' });
91 return;
92 }
93 effectiveInitiator = initiator;
94 }
95 const session = store.createSession(goal, effectiveInitiator, workflowId);
96 res.status(201).json(session);
97});
98 
99app.get("/api/sessions/:id", (req, res) => {
100 const session = store.getSession(req.params.id);
101 if (!session) {
102 res.status(404).json({ error: "session not found" });
103 return;
104 }
105 res.json(session);
106});
107 
108app.patch("/api/sessions/:id", (req, res) => {
109 const updates = req.body as {
110 autoForward?: boolean;
111 goal?: string;
112 currentPhaseIndex?: number;
113 };
114 if (updates.currentPhaseIndex !== undefined) {
115 const idx = updates.currentPhaseIndex;
116 if (!Number.isInteger(idx) || idx < 0) {
117 res.status(400).json({ error: "currentPhaseIndex must be a non-negative integer" });
118 return;
119 }
120 const existing = store.getSession(req.params.id);
121 if (!existing) {
122 res.status(404).json({ error: "session not found" });
123 return;
124 }
125 const template = getWorkflow(existing.workflow);
126 if (template.phases.length === 0) {
127 res.status(400).json({ error: "this workflow has no phases" });
128 return;
129 }
130 if (idx >= template.phases.length) {
131 res.status(400).json({
132 error: `currentPhaseIndex out of range (valid: 0..${template.phases.length - 1})`,
133 });
134 return;
135 }
136 }
137 const session = store.updateSession(req.params.id, updates);
138 if (!session) {
139 res.status(404).json({ error: "session not found" });
140 return;
141 }
142 res.json(session);
143});
144 
145app.delete("/api/sessions/:id", (req, res) => {
146 if (!store.deleteSession(req.params.id)) {
147 res.status(404).json({ error: "session not found" });
148 return;
149 }
150 res.status(204).end();
151});
152 
153app.post("/api/sessions/:id/close", (req, res) => {
154 const { reason } = req.body as { reason?: string };
155 const session = store.closeSession(req.params.id, "user", reason);
156 if (!session) {
157 res.status(404).json({ error: "session not found" });
158 return;
159 }
160 res.json(session);
161});
162 
163// ── Message actions ────────────────────────────────────────────────
164 
165function parseMid(raw: string, res: express.Response): number | null {
166 const n = parseInt(raw, 10);
167 if (isNaN(n)) {
168 res.status(400).json({ error: "invalid message id" });
169 return null;
170 }
171 return n;
172}
173 
174app.post("/api/sessions/:id/messages/:mid/deliver", (req, res) => {
175 const mid = parseMid(req.params.mid, res);
176 if (mid === null) return;
177 const msg = store.deliverMessage(req.params.id, mid);
178 if (!msg) {
179 res.status(404).json({ error: "message not found or already delivered" });
180 return;
181 }
182 res.json(msg);
183});
184 
185app.patch("/api/sessions/:id/messages/:mid", (req, res) => {
186 const mid = parseMid(req.params.mid, res);
187 if (mid === null) return;
188 const { body: newBody } = req.body as { body?: string };
189 if (!newBody || typeof newBody !== "string") {
190 res.status(400).json({ error: "body is required" });
191 return;
192 }
193 const msg = store.editMessage(req.params.id, mid, newBody);
194 if (!msg) {
195 res.status(404).json({ error: "message not found or not pending" });
196 return;
197 }
198 res.json(msg);
199});
200 
201app.delete("/api/sessions/:id/messages/:mid", (req, res) => {
202 const mid = parseMid(req.params.mid, res);
203 if (mid === null) return;
204 if (!store.rejectMessage(req.params.id, mid)) {
205 res.status(404).json({ error: "message not found or not pending" });
206 return;
207 }
208 res.status(204).end();
209});
210 
211// ── Prompts ────────────────────────────────────────────────────────
212 
213app.get("/api/sessions/:id/prompts", (req, res) => {
214 const session = store.getSession(req.params.id);
215 if (!session) {
216 res.status(404).json({ error: "session not found" });
217 return;
218 }
219 const proto = req.protocol;
220 const host = req.get("host") || `localhost:${port}`;
221 const baseUrl = `${proto}://${host}`;
222 res.json(generatePrompts(session, baseUrl));
223});
224 
225// ── SSE for web UI ─────────────────────────────────────────────────
226 
227app.get("/events", (req, res) => {
228 res.writeHead(200, {
229 "Content-Type": "text/event-stream",
230 "Cache-Control": "no-cache",
231 Connection: "keep-alive",
232 });
233 res.write(`data: ${JSON.stringify({ type: "connected" })}\n\n`);
234 
235 const onUpdate = (event: UIEvent) => {
236 res.write(`data: ${JSON.stringify(event)}\n\n`);
237 };
238 store.on("ui-update", onUpdate);
239 
240 req.on("close", () => {
241 store.removeListener("ui-update", onUpdate);
242 });
243});
244 
245// ── Static ─────────────────────────────────────────────────────────
246 
247app.get("/", (_req, res) => {
248 res.sendFile(path.resolve(__dirname, "../web/index.html"));
249});
250 
251// ── Error handler ──────────────────────────────────────────────────
252 
253app.use((err: unknown, _req: express.Request, res: express.Response, _next: express.NextFunction) => {
254 const message = err instanceof Error ? err.message : "Internal server error";
255 res.status(500).json({ error: message });
256});
257 
258// ── Start ──────────────────────────────────────────────────────────
259 
260const server = app.listen(port, "127.0.0.1", () => {
261 console.log(`collab-bridge running at http://localhost:${port}`);
262 console.log(` Web UI: http://localhost:${port}/`);
263 console.log(` Claude MCP: http://localhost:${port}/mcp/claude`);
264 console.log(` Codex MCP: http://localhost:${port}/mcp/codex`);
265});
266 
267// collab_wait_for_reply long-polls for up to 5 minutes (300s).
268// Ensure the HTTP socket idle timeout doesn't kill it.
269// (Default is 0 on Node 13+, but was 120s on older versions.)
270server.timeout = 0;