import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { StreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/streamableHttp.js"; import { z } from "zod"; import type { IncomingMessage, ServerResponse } from "node:http"; import type { AgentId, Session, Store } from "./store.js"; import { getWorkflow, getPhase } from "./workflows.js"; interface PhaseHint { index: number; total: number; id: string; name: string; instructions: string; } function computePhaseHint(session: Session, receiver: AgentId): PhaseHint | undefined { const template = getWorkflow(session.workflow); if (template.phases.length === 0) return undefined; const phase = getPhase(template, session.currentPhaseIndex); if (!phase) return undefined; const instructions = phase.reminderForReceiver[receiver] ?? phase.instructions[receiver]; if (!instructions) return undefined; return { index: session.currentPhaseIndex, total: template.phases.length, id: phase.id, name: phase.name, instructions, }; } function createMcpServer(agent: AgentId, store: Store): McpServer { const server = new McpServer( { name: "collab-bridge", version: "0.1.0" }, { instructions: `Collaboration bridge MCP server. Endpoint for agent: ${agent}.`, }, ); server.registerTool( "collab_send_message", { description: "Send a message to the other agent in this collaboration session. Returns confirmation with the assigned message ID.", inputSchema: { collabId: z.string().describe("The collaboration session ID"), message: z .string() .describe("The message content to send (ASCII-safe markdown)"), }, }, async (args) => { const session = store.getSession(args.collabId); if (!session) { return { content: [ { type: "text" as const, text: `Error: session '${args.collabId}' not found.` }, ], isError: true, }; } if (session.status === "closed") { return { content: [ { type: "text" as const, text: `Error: session '${args.collabId}' is closed.` }, ], isError: true, }; } const unseenPrior = store.getLatestUnseenOutbound(args.collabId, agent); const msg = store.addMessage(args.collabId, agent, args.message); if (!msg) { return { content: [{ type: "text" as const, text: "Error: failed to add message." }], isError: true, }; } return { content: [ { type: "text" as const, text: JSON.stringify({ status: "sent", messageId: msg.id, deliveryStatus: msg.status, ...(unseenPrior ? { unseenPriorMessageId: unseenPrior.id, warning: `You sent message ${msg.id} before the other agent received message ${unseenPrior.id}. ` + "Because waits return only the latest message, send a superseding full-state message if the earlier message contained required context.", } : {}), }), }, ], }; }, ); server.registerTool( "collab_wait_for_reply", { description: "Wait for the next message from the other agent (long-polls ~5min). Pass lastSeenMessageId from your last received message (0 if none). Returns ONLY the latest message to save context.", inputSchema: { collabId: z.string().describe("The collaboration session ID"), lastSeenMessageId: z .number() .describe("ID of the last message you received (0 if none)"), }, }, async (args) => { const session = store.getSession(args.collabId); if (!session) { return { content: [ { type: "text" as const, text: `Error: session '${args.collabId}' not found.` }, ], isError: true, }; } const result = await store.waitForMessage( args.collabId, agent, args.lastSeenMessageId, 300_000, ); // Re-fetch the session so we read the freshest currentPhaseIndex (it may have // been advanced via PATCH while we were long-polling). const fresh = store.getSession(args.collabId); const phaseHint = fresh ? computePhaseHint(fresh, agent) : undefined; if (result.status === "received") { const skippedMessageIds = result.skippedMessageIds; return { content: [ { type: "text" as const, text: JSON.stringify({ status: "received", messageId: result.message.id, from: result.message.from, message: result.message.body, ...(skippedMessageIds.length > 0 ? { skippedMessageIds, warning: `Skipped older delivered message(s) ${skippedMessageIds.join(", ")} because collab_wait_for_reply returns only the latest message. ` + "Ask the sender for a superseding full-state message if any skipped context matters.", } : {}), ...(phaseHint ? { phaseHint } : {}), }), }, ], }; } if (result.status === "closed") { return { content: [ { type: "text" as const, text: JSON.stringify({ status: "closed", hint: "This collaboration session is closed. Stop waiting and do not send more messages.", ...(phaseHint ? { phaseHint } : {}), }), }, ], }; } return { content: [ { type: "text" as const, text: JSON.stringify({ status: "timeout", lastSeenMessageId: args.lastSeenMessageId, hint: "No new message within timeout. Retry with the same lastSeenMessageId -- the other agent may still be thinking.", ...(phaseHint ? { phaseHint } : {}), }), }, ], }; }, ); server.registerTool( "collab_close_session", { description: "Close this collaboration session. After closing, waits return immediately and sends are rejected.", inputSchema: { collabId: z.string().describe("The collaboration session ID"), reason: z.string().optional().describe("Optional short reason for closing the session"), }, }, async (args) => { const session = store.closeSession(args.collabId, agent, args.reason); if (!session) { return { content: [ { type: "text" as const, text: `Error: session '${args.collabId}' not found.` }, ], isError: true, }; } return { content: [ { type: "text" as const, text: JSON.stringify({ status: "closed", collabId: session.id, closedBy: session.closedBy, closedReason: session.closedReason ?? null, }), }, ], }; }, ); return server; } export async function handleMcpRequest( agent: AgentId, store: Store, req: IncomingMessage, res: ServerResponse, ): Promise { const server = createMcpServer(agent, store); const transport = new StreamableHTTPServerTransport({ sessionIdGenerator: undefined, }); res.on("close", () => { transport.close(); }); await server.connect(transport); await transport.handleRequest(req, res, (req as any).body); }