Skip to content

Commit 2dd6fe3

Browse files
committed
refactor: simplify Express implementation following mcp-echarts pattern
1 parent 55bbf84 commit 2dd6fe3

4 files changed

Lines changed: 67 additions & 311 deletions

File tree

src/services/sse.ts

Lines changed: 10 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -1,53 +1,36 @@
11
import type { Server } from "@modelcontextprotocol/sdk/server/index.js";
22
import { SSEServerTransport } from "@modelcontextprotocol/sdk/server/sse.js";
3-
import type { Request, Response } from "express";
4-
import { createExpressServer } from "../utils";
3+
import express from "express";
54

65
export const startSSEMcpServer = async (
76
server: Server,
87
endpoint = "/sse",
98
port = 1122,
109
): Promise<void> => {
11-
const activeTransports: Record<string, SSEServerTransport> = {};
12-
13-
// Custom cleanup for SSE server
14-
const cleanup = () => {
15-
// Close all active transports
16-
for (const transport of Object.values(activeTransports)) {
17-
transport.close();
18-
}
19-
server.close();
20-
};
10+
const app = express();
11+
app.use(express.json());
2112

22-
// Create Express server
23-
const { app, start } = createExpressServer({
24-
port,
25-
serverType: "SSE Server",
26-
cleanup,
27-
});
13+
const activeTransports: Record<string, SSEServerTransport> = {};
2814

2915
// Handle GET requests to the SSE endpoint
30-
app.get(endpoint, async (req: Request, res: Response) => {
16+
app.get(endpoint, async (req, res) => {
3117
const transport = new SSEServerTransport("/messages", res);
3218
activeTransports[transport.sessionId] = transport;
3319

3420
let closed = false;
3521

3622
res.on("close", async () => {
3723
closed = true;
38-
3924
try {
4025
await server.close();
4126
} catch (error) {
4227
console.error("Error closing server:", error);
4328
}
44-
4529
delete activeTransports[transport.sessionId];
4630
});
4731

4832
try {
4933
await server.connect(transport);
50-
5134
await transport.send({
5235
jsonrpc: "2.0",
5336
method: "sse/connection",
@@ -62,17 +45,15 @@ export const startSSEMcpServer = async (
6245
});
6346

6447
// Handle POST requests to the messages endpoint
65-
app.post("/messages", async (req: Request, res: Response) => {
48+
app.post("/messages", async (req, res) => {
6649
const sessionId = req.query.sessionId as string;
6750

6851
if (!sessionId) {
6952
res.status(400).send("No sessionId");
7053
return;
7154
}
7255

73-
const activeTransport: SSEServerTransport | undefined =
74-
activeTransports[sessionId];
75-
56+
const activeTransport = activeTransports[sessionId];
7657
if (!activeTransport) {
7758
res.status(400).send("No active transport");
7859
return;
@@ -81,6 +62,7 @@ export const startSSEMcpServer = async (
8162
await activeTransport.handlePostMessage(req, res);
8263
});
8364

84-
// Start the server and log endpoints
85-
start([endpoint, "/messages"]);
65+
app.listen(port, () => {
66+
console.log(`SSE Server running on http://localhost:${port}${endpoint}`);
67+
});
8668
};

src/services/streamable.ts

Lines changed: 57 additions & 151 deletions
Original file line numberDiff line numberDiff line change
@@ -5,184 +5,90 @@ import {
55
StreamableHTTPServerTransport,
66
} from "@modelcontextprotocol/sdk/server/streamableHttp.js";
77
import { isInitializeRequest } from "@modelcontextprotocol/sdk/types.js";
8-
import type { Request, Response } from "express";
9-
import { InMemoryEventStore, createExpressServer, getBody } from "../utils";
8+
import express from "express";
9+
import { InMemoryEventStore } from "../utils";
1010

1111
export const startHTTPStreamableServer = async (
1212
createServer: () => Server,
1313
endpoint = "/mcp",
1414
port = 1122,
1515
eventStore: EventStore = new InMemoryEventStore(),
1616
): Promise<void> => {
17-
const activeTransports: Record<
18-
string,
19-
{
20-
server: Server;
21-
transport: StreamableHTTPServerTransport;
22-
}
23-
> = {};
24-
25-
// Custom cleanup for streamable server
26-
const cleanup = () => {
27-
for (const { server, transport } of Object.values(activeTransports)) {
28-
transport.close();
29-
server.close();
30-
}
31-
};
32-
33-
// Create Express server
34-
const { app, start } = createExpressServer({
35-
port,
36-
serverType: "HTTP Streamable Server",
37-
cleanup,
38-
});
39-
40-
// Handle POST requests to endpoint
41-
app.post(endpoint, async (req: Request, res: Response) => {
42-
try {
43-
const sessionId = Array.isArray(req.headers["mcp-session-id"])
44-
? req.headers["mcp-session-id"][0]
45-
: req.headers["mcp-session-id"];
46-
let transport: StreamableHTTPServerTransport;
47-
let server: Server;
48-
const body = req.body;
49-
50-
/**
51-
* diagram: https://modelcontextprotocol.io/specification/2025-03-26/basic/transports#sequence-diagram.
52-
*/
53-
// 1. If the sessionId is provided and the server is already created, use the existing transport and server.
54-
if (sessionId && activeTransports[sessionId]) {
55-
transport = activeTransports[sessionId].transport;
56-
server = activeTransports[sessionId].server;
57-
58-
// 2. If the sessionId is not provided and the request is an initialize request, create a new transport for the session.
59-
} else if (!sessionId && isInitializeRequest(body)) {
60-
transport = new StreamableHTTPServerTransport({
61-
// use the event store to store the events to replay on reconnect.
62-
// more details: https://modelcontextprotocol.io/specification/2025-03-26/basic/transports#resumability-and-redelivery.
63-
eventStore: eventStore || new InMemoryEventStore(),
64-
onsessioninitialized: (_sessionId: string) => {
65-
// add only when the id Session id is generated.
66-
activeTransports[_sessionId] = {
67-
server,
68-
transport,
69-
};
70-
},
71-
sessionIdGenerator: randomUUID,
72-
});
17+
const app = express();
18+
app.use(express.json());
7319

74-
// Handle the server close event.
75-
transport.onclose = async () => {
76-
const sid = transport.sessionId;
77-
if (sid && activeTransports[sid]) {
78-
try {
79-
await server?.close();
80-
} catch (error) {
81-
console.error("Error closing server:", error);
82-
}
20+
// Store transports by session ID
21+
const transports: Record<string, StreamableHTTPServerTransport> = {};
8322

84-
// delete used transport and server to avoid memory leak.
85-
delete activeTransports[sid];
86-
}
87-
};
23+
// Handle POST requests for client-to-server communication
24+
app.post(endpoint, async (req, res) => {
25+
const sessionId = req.headers["mcp-session-id"] as string | undefined;
26+
let transport: StreamableHTTPServerTransport;
27+
28+
if (sessionId && transports[sessionId]) {
29+
// Reuse existing transport
30+
transport = transports[sessionId];
31+
} else if (!sessionId && isInitializeRequest(req.body)) {
32+
// New initialization request
33+
transport = new StreamableHTTPServerTransport({
34+
eventStore,
35+
sessionIdGenerator: () => randomUUID(),
36+
onsessioninitialized: (sessionId) => {
37+
transports[sessionId] = transport;
38+
},
39+
});
8840

89-
// Create the server
90-
try {
91-
server = createServer();
92-
} catch (error) {
93-
console.error("Error creating server:", error);
94-
res.status(500).json({
95-
error: { code: -32603, message: "Error creating server" },
96-
id: null,
97-
jsonrpc: "2.0",
98-
});
99-
return;
41+
// Clean up transport when closed
42+
transport.onclose = () => {
43+
if (transport.sessionId) {
44+
delete transports[transport.sessionId];
10045
}
46+
};
10147

102-
server.connect(transport);
103-
104-
await transport.handleRequest(req, res, body);
105-
return;
106-
} else {
107-
// Error if the server is not created but the request is not an initialize request.
108-
res.status(400).json({
109-
error: {
110-
code: -32000,
111-
message: "Bad Request: No valid session ID provided",
112-
},
113-
id: null,
114-
jsonrpc: "2.0",
115-
});
116-
return;
117-
}
118-
119-
// Handle the request if the server is already created.
120-
await transport.handleRequest(req, res, body);
121-
} catch (error) {
122-
console.error("Error handling request:", error);
123-
res.status(500).json({
124-
error: { code: -32603, message: "Internal Server Error" },
125-
id: null,
48+
const server = createServer();
49+
await server.connect(transport);
50+
} else {
51+
// Invalid request
52+
res.status(400).json({
12653
jsonrpc: "2.0",
54+
error: {
55+
code: -32000,
56+
message: "Bad Request: No valid session ID provided",
57+
},
58+
id: null,
12759
});
128-
}
129-
});
130-
131-
// Handle GET requests to endpoint
132-
app.get(endpoint, async (req: Request, res: Response) => {
133-
const sessionId = req.headers["mcp-session-id"] as string | undefined;
134-
const activeTransport:
135-
| {
136-
server: Server;
137-
transport: StreamableHTTPServerTransport;
138-
}
139-
| undefined = sessionId ? activeTransports[sessionId] : undefined;
140-
141-
if (!sessionId) {
142-
res.status(400).send("No sessionId");
143-
return;
144-
}
145-
146-
if (!activeTransport) {
147-
res.status(400).send("No active transport");
14860
return;
14961
}
15062

151-
const lastEventId = req.headers["last-event-id"] as string | undefined;
152-
if (lastEventId) {
153-
console.log(`Client reconnecting with Last-Event-ID: ${lastEventId}`);
154-
} else {
155-
console.log(`Establishing new SSE stream for session ${sessionId}`);
156-
}
157-
158-
await activeTransport.transport.handleRequest(req, res);
63+
// Handle the request
64+
await transport.handleRequest(req, res, req.body);
15965
});
16066

161-
// Handle DELETE requests to endpoint
162-
app.delete(endpoint, async (req: Request, res: Response) => {
163-
console.log("received delete request");
67+
// Handle GET requests for server-to-client notifications via SSE
68+
app.get(endpoint, async (req, res) => {
16469
const sessionId = req.headers["mcp-session-id"] as string | undefined;
165-
if (!sessionId) {
166-
res.status(400).send("Invalid or missing sessionId");
70+
if (!sessionId || !transports[sessionId]) {
71+
res.status(400).send("Invalid or missing session ID");
16772
return;
16873
}
16974

170-
console.log("received delete request for session", sessionId);
75+
const transport = transports[sessionId];
76+
await transport.handleRequest(req, res);
77+
});
17178

172-
const transport = activeTransports[sessionId]?.transport;
173-
if (!transport) {
174-
res.status(400).send("No active transport");
79+
// Handle DELETE requests for session termination
80+
app.delete(endpoint, async (req, res) => {
81+
const sessionId = req.headers["mcp-session-id"] as string | undefined;
82+
if (!sessionId || !transports[sessionId]) {
83+
res.status(400).send("Invalid or missing session ID");
17584
return;
17685
}
17786

178-
try {
179-
await transport.handleRequest(req, res);
180-
} catch (error) {
181-
console.error("Error handling delete request:", error);
182-
res.status(500).send("Error handling delete request");
183-
}
87+
const transport = transports[sessionId];
88+
await transport.handleRequest(req, res);
18489
});
18590

186-
// Start the server and log endpoints
187-
start([endpoint]);
91+
app.listen(port, () => {
92+
console.log(`Streamable HTTP Server running on http://localhost:${port}${endpoint}`);
93+
});
18894
};

0 commit comments

Comments
 (0)