sseAndStreamableHttpCompatibleServer.js 10 KB

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