sseAndStreamableHttpCompatibleServer.js 8.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231
  1. import { randomUUID } from 'node:crypto';
  2. import { McpServer } from '../../server/mcp.js';
  3. import { StreamableHTTPServerTransport } from '../../server/streamableHttp.js';
  4. import { SSEServerTransport } from '../../server/sse.js';
  5. import * as z from 'zod/v4';
  6. import { isInitializeRequest } from '../../types.js';
  7. import { InMemoryEventStore } from '../shared/inMemoryEventStore.js';
  8. import { createMcpExpressApp } from '../../server/express.js';
  9. /**
  10. * This example server demonstrates backwards compatibility with both:
  11. * 1. The deprecated HTTP+SSE transport (protocol version 2024-11-05)
  12. * 2. The Streamable HTTP transport (protocol version 2025-11-25)
  13. *
  14. * It maintains a single MCP server instance but exposes two transport options:
  15. * - /mcp: The new Streamable HTTP endpoint (supports GET/POST/DELETE)
  16. * - /sse: The deprecated SSE endpoint for older clients (GET to establish stream)
  17. * - /messages: The deprecated POST endpoint for older clients (POST to send messages)
  18. */
  19. const getServer = () => {
  20. const server = new McpServer({
  21. name: 'backwards-compatible-server',
  22. version: '1.0.0'
  23. }, { capabilities: { logging: {} } });
  24. // Register a simple tool that sends notifications over time
  25. server.registerTool('start-notification-stream', {
  26. description: 'Starts sending periodic notifications for testing resumability',
  27. inputSchema: {
  28. interval: z.number().describe('Interval in milliseconds between notifications').default(100),
  29. count: z.number().describe('Number of notifications to send (0 for 100)').default(50)
  30. }
  31. }, async ({ interval, count }, extra) => {
  32. const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms));
  33. let counter = 0;
  34. while (count === 0 || counter < count) {
  35. counter++;
  36. try {
  37. await server.sendLoggingMessage({
  38. level: 'info',
  39. data: `Periodic notification #${counter} at ${new Date().toISOString()}`
  40. }, extra.sessionId);
  41. }
  42. catch (error) {
  43. console.error('Error sending notification:', error);
  44. }
  45. // Wait for the specified interval
  46. await sleep(interval);
  47. }
  48. return {
  49. content: [
  50. {
  51. type: 'text',
  52. text: `Started sending periodic notifications every ${interval}ms`
  53. }
  54. ]
  55. };
  56. });
  57. return server;
  58. };
  59. // Create Express application
  60. const app = createMcpExpressApp();
  61. // Store transports by session ID
  62. const transports = {};
  63. //=============================================================================
  64. // STREAMABLE HTTP TRANSPORT (PROTOCOL VERSION 2025-11-25)
  65. //=============================================================================
  66. // Handle all MCP Streamable HTTP requests (GET, POST, DELETE) on a single endpoint
  67. app.all('/mcp', async (req, res) => {
  68. console.log(`Received ${req.method} request to /mcp`);
  69. try {
  70. // Check for existing session ID
  71. const sessionId = req.headers['mcp-session-id'];
  72. let transport;
  73. if (sessionId && transports[sessionId]) {
  74. // Check if the transport is of the correct type
  75. const existingTransport = transports[sessionId];
  76. if (existingTransport instanceof StreamableHTTPServerTransport) {
  77. // Reuse existing transport
  78. transport = existingTransport;
  79. }
  80. else {
  81. // Transport exists but is not a StreamableHTTPServerTransport (could be SSEServerTransport)
  82. res.status(400).json({
  83. jsonrpc: '2.0',
  84. error: {
  85. code: -32000,
  86. message: 'Bad Request: Session exists but uses a different transport protocol'
  87. },
  88. id: null
  89. });
  90. return;
  91. }
  92. }
  93. else if (!sessionId && req.method === 'POST' && isInitializeRequest(req.body)) {
  94. const eventStore = new InMemoryEventStore();
  95. transport = new StreamableHTTPServerTransport({
  96. sessionIdGenerator: () => randomUUID(),
  97. eventStore, // Enable resumability
  98. onsessioninitialized: sessionId => {
  99. // Store the transport by session ID when session is initialized
  100. console.log(`StreamableHTTP session initialized with ID: ${sessionId}`);
  101. transports[sessionId] = transport;
  102. }
  103. });
  104. // Set up onclose handler to clean up transport when closed
  105. transport.onclose = () => {
  106. const sid = transport.sessionId;
  107. if (sid && transports[sid]) {
  108. console.log(`Transport closed for session ${sid}, removing from transports map`);
  109. delete transports[sid];
  110. }
  111. };
  112. // Connect the transport to the MCP server
  113. const server = getServer();
  114. await server.connect(transport);
  115. }
  116. else {
  117. // Invalid request - no session ID or not initialization request
  118. res.status(400).json({
  119. jsonrpc: '2.0',
  120. error: {
  121. code: -32000,
  122. message: 'Bad Request: No valid session ID provided'
  123. },
  124. id: null
  125. });
  126. return;
  127. }
  128. // Handle the request with the transport
  129. await transport.handleRequest(req, res, req.body);
  130. }
  131. catch (error) {
  132. console.error('Error handling MCP request:', error);
  133. if (!res.headersSent) {
  134. res.status(500).json({
  135. jsonrpc: '2.0',
  136. error: {
  137. code: -32603,
  138. message: 'Internal server error'
  139. },
  140. id: null
  141. });
  142. }
  143. }
  144. });
  145. //=============================================================================
  146. // DEPRECATED HTTP+SSE TRANSPORT (PROTOCOL VERSION 2024-11-05)
  147. //=============================================================================
  148. app.get('/sse', async (req, res) => {
  149. console.log('Received GET request to /sse (deprecated SSE transport)');
  150. const transport = new SSEServerTransport('/messages', res);
  151. transports[transport.sessionId] = transport;
  152. res.on('close', () => {
  153. delete transports[transport.sessionId];
  154. });
  155. const server = getServer();
  156. await server.connect(transport);
  157. });
  158. app.post('/messages', async (req, res) => {
  159. const sessionId = req.query.sessionId;
  160. let transport;
  161. const existingTransport = transports[sessionId];
  162. if (existingTransport instanceof SSEServerTransport) {
  163. // Reuse existing transport
  164. transport = existingTransport;
  165. }
  166. else {
  167. // Transport exists but is not a SSEServerTransport (could be StreamableHTTPServerTransport)
  168. res.status(400).json({
  169. jsonrpc: '2.0',
  170. error: {
  171. code: -32000,
  172. message: 'Bad Request: Session exists but uses a different transport protocol'
  173. },
  174. id: null
  175. });
  176. return;
  177. }
  178. if (transport) {
  179. await transport.handlePostMessage(req, res, req.body);
  180. }
  181. else {
  182. res.status(400).send('No transport found for sessionId');
  183. }
  184. });
  185. // Start the server
  186. const PORT = 3000;
  187. app.listen(PORT, error => {
  188. if (error) {
  189. console.error('Failed to start server:', error);
  190. process.exit(1);
  191. }
  192. console.log(`Backwards compatible MCP server listening on port ${PORT}`);
  193. console.log(`
  194. ==============================================
  195. SUPPORTED TRANSPORT OPTIONS:
  196. 1. Streamable Http(Protocol version: 2025-11-25)
  197. Endpoint: /mcp
  198. Methods: GET, POST, DELETE
  199. Usage:
  200. - Initialize with POST to /mcp
  201. - Establish SSE stream with GET to /mcp
  202. - Send requests with POST to /mcp
  203. - Terminate session with DELETE to /mcp
  204. 2. Http + SSE (Protocol version: 2024-11-05)
  205. Endpoints: /sse (GET) and /messages (POST)
  206. Usage:
  207. - Establish SSE stream with GET to /sse
  208. - Send requests with POST to /messages?sessionId=<id>
  209. ==============================================
  210. `);
  211. });
  212. // Handle server shutdown
  213. process.on('SIGINT', async () => {
  214. console.log('Shutting down server...');
  215. // Close all active transports to properly clean up resources
  216. for (const sessionId in transports) {
  217. try {
  218. console.log(`Closing transport for session ${sessionId}`);
  219. await transports[sessionId].close();
  220. delete transports[sessionId];
  221. }
  222. catch (error) {
  223. console.error(`Error closing transport for session ${sessionId}:`, error);
  224. }
  225. }
  226. console.log('Server shutdown complete');
  227. process.exit(0);
  228. });
  229. //# sourceMappingURL=sseAndStreamableHttpCompatibleServer.js.map