Skip to content
File

Blob: src/mcp-handler.ts

typescript252 lines
1import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js";
2import { StreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/streamableHttp.js";
3import { z } from "zod";
4import type { IncomingMessage, ServerResponse } from "node:http";
5import type { AgentId, Session, Store } from "./store.js";
6import { getWorkflow, getPhase } from "./workflows.js";
7 
8interface PhaseHint {
9 index: number;
10 total: number;
11 id: string;
12 name: string;
13 instructions: string;
14}
15 
16function computePhaseHint(session: Session, receiver: AgentId): PhaseHint | undefined {
17 const template = getWorkflow(session.workflow);
18 if (template.phases.length === 0) return undefined;
19 const phase = getPhase(template, session.currentPhaseIndex);
20 if (!phase) return undefined;
21 const instructions =
22 phase.reminderForReceiver[receiver] ?? phase.instructions[receiver];
23 if (!instructions) return undefined;
24 return {
25 index: session.currentPhaseIndex,
26 total: template.phases.length,
27 id: phase.id,
28 name: phase.name,
29 instructions,
30 };
31}
32 
33function createMcpServer(agent: AgentId, store: Store): McpServer {
34 const server = new McpServer(
35 { name: "collab-bridge", version: "0.1.0" },
36 {
37 instructions: `Collaboration bridge MCP server. Endpoint for agent: ${agent}.`,
38 },
39 );
40 
41 server.registerTool(
42 "collab_send_message",
43 {
44 description:
45 "Send a message to the other agent in this collaboration session. Returns confirmation with the assigned message ID.",
46 inputSchema: {
47 collabId: z.string().describe("The collaboration session ID"),
48 message: z
49 .string()
50 .describe("The message content to send (ASCII-safe markdown)"),
51 },
52 },
53 async (args) => {
54 const session = store.getSession(args.collabId);
55 if (!session) {
56 return {
57 content: [
58 { type: "text" as const, text: `Error: session '${args.collabId}' not found.` },
59 ],
60 isError: true,
61 };
62 }
63 if (session.status === "closed") {
64 return {
65 content: [
66 { type: "text" as const, text: `Error: session '${args.collabId}' is closed.` },
67 ],
68 isError: true,
69 };
70 }
71 
72 const unseenPrior = store.getLatestUnseenOutbound(args.collabId, agent);
73 const msg = store.addMessage(args.collabId, agent, args.message);
74 if (!msg) {
75 return {
76 content: [{ type: "text" as const, text: "Error: failed to add message." }],
77 isError: true,
78 };
79 }
80 
81 return {
82 content: [
83 {
84 type: "text" as const,
85 text: JSON.stringify({
86 status: "sent",
87 messageId: msg.id,
88 deliveryStatus: msg.status,
89 ...(unseenPrior
90 ? {
91 unseenPriorMessageId: unseenPrior.id,
92 warning:
93 `You sent message ${msg.id} before the other agent received message ${unseenPrior.id}. ` +
94 "Because waits return only the latest message, send a superseding full-state message if the earlier message contained required context.",
95 }
96 : {}),
97 }),
98 },
99 ],
100 };
101 },
102 );
103 
104 server.registerTool(
105 "collab_wait_for_reply",
106 {
107 description:
108 "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.",
109 inputSchema: {
110 collabId: z.string().describe("The collaboration session ID"),
111 lastSeenMessageId: z
112 .number()
113 .describe("ID of the last message you received (0 if none)"),
114 },
115 },
116 async (args) => {
117 const session = store.getSession(args.collabId);
118 if (!session) {
119 return {
120 content: [
121 { type: "text" as const, text: `Error: session '${args.collabId}' not found.` },
122 ],
123 isError: true,
124 };
125 }
126 
127 const result = await store.waitForMessage(
128 args.collabId,
129 agent,
130 args.lastSeenMessageId,
131 300_000,
132 );
133 
134 // Re-fetch the session so we read the freshest currentPhaseIndex (it may have
135 // been advanced via PATCH while we were long-polling).
136 const fresh = store.getSession(args.collabId);
137 const phaseHint = fresh ? computePhaseHint(fresh, agent) : undefined;
138 
139 if (result.status === "received") {
140 const skippedMessageIds = result.skippedMessageIds;
141 return {
142 content: [
143 {
144 type: "text" as const,
145 text: JSON.stringify({
146 status: "received",
147 messageId: result.message.id,
148 from: result.message.from,
149 message: result.message.body,
150 ...(skippedMessageIds.length > 0
151 ? {
152 skippedMessageIds,
153 warning:
154 `Skipped older delivered message(s) ${skippedMessageIds.join(", ")} because collab_wait_for_reply returns only the latest message. ` +
155 "Ask the sender for a superseding full-state message if any skipped context matters.",
156 }
157 : {}),
158 ...(phaseHint ? { phaseHint } : {}),
159 }),
160 },
161 ],
162 };
163 }
164 
165 if (result.status === "closed") {
166 return {
167 content: [
168 {
169 type: "text" as const,
170 text: JSON.stringify({
171 status: "closed",
172 hint: "This collaboration session is closed. Stop waiting and do not send more messages.",
173 ...(phaseHint ? { phaseHint } : {}),
174 }),
175 },
176 ],
177 };
178 }
179 
180 return {
181 content: [
182 {
183 type: "text" as const,
184 text: JSON.stringify({
185 status: "timeout",
186 lastSeenMessageId: args.lastSeenMessageId,
187 hint: "No new message within timeout. Retry with the same lastSeenMessageId -- the other agent may still be thinking.",
188 ...(phaseHint ? { phaseHint } : {}),
189 }),
190 },
191 ],
192 };
193 },
194 );
195 
196 server.registerTool(
197 "collab_close_session",
198 {
199 description:
200 "Close this collaboration session. After closing, waits return immediately and sends are rejected.",
201 inputSchema: {
202 collabId: z.string().describe("The collaboration session ID"),
203 reason: z.string().optional().describe("Optional short reason for closing the session"),
204 },
205 },
206 async (args) => {
207 const session = store.closeSession(args.collabId, agent, args.reason);
208 if (!session) {
209 return {
210 content: [
211 { type: "text" as const, text: `Error: session '${args.collabId}' not found.` },
212 ],
213 isError: true,
214 };
215 }
216 
217 return {
218 content: [
219 {
220 type: "text" as const,
221 text: JSON.stringify({
222 status: "closed",
223 collabId: session.id,
224 closedBy: session.closedBy,
225 closedReason: session.closedReason ?? null,
226 }),
227 },
228 ],
229 };
230 },
231 );
232 
233 return server;
234}
235 
236export async function handleMcpRequest(
237 agent: AgentId,
238 store: Store,
239 req: IncomingMessage,
240 res: ServerResponse,
241): Promise<void> {
242 const server = createMcpServer(agent, store);
243 const transport = new StreamableHTTPServerTransport({
244 sessionIdGenerator: undefined,
245 });
246 res.on("close", () => {
247 transport.close();
248 });
249 await server.connect(transport);
250 await transport.handleRequest(req, res, (req as any).body);
251}