simpleSseServer.js 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168
  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 mcp_js_1 = require("../../server/mcp.js");
  27. const sse_js_1 = require("../../server/sse.js");
  28. const z = __importStar(require("zod/v4"));
  29. const express_js_1 = require("../../server/express.js");
  30. /**
  31. * This example server demonstrates the deprecated HTTP+SSE transport
  32. * (protocol version 2024-11-05). It mainly used for testing backward compatible clients.
  33. *
  34. * The server exposes two endpoints:
  35. * - /mcp: For establishing the SSE stream (GET)
  36. * - /messages: For receiving client messages (POST)
  37. *
  38. */
  39. // Create an MCP server instance
  40. const getServer = () => {
  41. const server = new mcp_js_1.McpServer({
  42. name: 'simple-sse-server',
  43. version: '1.0.0'
  44. }, { capabilities: { logging: {} } });
  45. server.registerTool('start-notification-stream', {
  46. description: 'Starts sending periodic notifications',
  47. inputSchema: {
  48. interval: z.number().describe('Interval in milliseconds between notifications').default(1000),
  49. count: z.number().describe('Number of notifications to send').default(10)
  50. }
  51. }, async ({ interval, count }, extra) => {
  52. const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms));
  53. let counter = 0;
  54. // Send the initial notification
  55. await server.sendLoggingMessage({
  56. level: 'info',
  57. data: `Starting notification stream with ${count} messages every ${interval}ms`
  58. }, extra.sessionId);
  59. // Send periodic notifications
  60. while (counter < count) {
  61. counter++;
  62. await sleep(interval);
  63. try {
  64. await server.sendLoggingMessage({
  65. level: 'info',
  66. data: `Notification #${counter} at ${new Date().toISOString()}`
  67. }, extra.sessionId);
  68. }
  69. catch (error) {
  70. console.error('Error sending notification:', error);
  71. }
  72. }
  73. return {
  74. content: [
  75. {
  76. type: 'text',
  77. text: `Completed sending ${count} notifications every ${interval}ms`
  78. }
  79. ]
  80. };
  81. });
  82. return server;
  83. };
  84. const app = (0, express_js_1.createMcpExpressApp)();
  85. // Store transports by session ID
  86. const transports = {};
  87. // SSE endpoint for establishing the stream
  88. app.get('/mcp', async (req, res) => {
  89. console.log('Received GET request to /sse (establishing SSE stream)');
  90. try {
  91. // Create a new SSE transport for the client
  92. // The endpoint for POST messages is '/messages'
  93. const transport = new sse_js_1.SSEServerTransport('/messages', res);
  94. // Store the transport by session ID
  95. const sessionId = transport.sessionId;
  96. transports[sessionId] = transport;
  97. // Set up onclose handler to clean up transport when closed
  98. transport.onclose = () => {
  99. console.log(`SSE transport closed for session ${sessionId}`);
  100. delete transports[sessionId];
  101. };
  102. // Connect the transport to the MCP server
  103. const server = getServer();
  104. await server.connect(transport);
  105. console.log(`Established SSE stream with session ID: ${sessionId}`);
  106. }
  107. catch (error) {
  108. console.error('Error establishing SSE stream:', error);
  109. if (!res.headersSent) {
  110. res.status(500).send('Error establishing SSE stream');
  111. }
  112. }
  113. });
  114. // Messages endpoint for receiving client JSON-RPC requests
  115. app.post('/messages', async (req, res) => {
  116. console.log('Received POST request to /messages');
  117. // Extract session ID from URL query parameter
  118. // In the SSE protocol, this is added by the client based on the endpoint event
  119. const sessionId = req.query.sessionId;
  120. if (!sessionId) {
  121. console.error('No session ID provided in request URL');
  122. res.status(400).send('Missing sessionId parameter');
  123. return;
  124. }
  125. const transport = transports[sessionId];
  126. if (!transport) {
  127. console.error(`No active transport found for session ID: ${sessionId}`);
  128. res.status(404).send('Session not found');
  129. return;
  130. }
  131. try {
  132. // Handle the POST message with the transport
  133. await transport.handlePostMessage(req, res, req.body);
  134. }
  135. catch (error) {
  136. console.error('Error handling request:', error);
  137. if (!res.headersSent) {
  138. res.status(500).send('Error handling request');
  139. }
  140. }
  141. });
  142. // Start the server
  143. const PORT = 3000;
  144. app.listen(PORT, error => {
  145. if (error) {
  146. console.error('Failed to start server:', error);
  147. process.exit(1);
  148. }
  149. console.log(`Simple SSE Server (deprecated protocol version 2024-11-05) listening on port ${PORT}`);
  150. });
  151. // Handle server shutdown
  152. process.on('SIGINT', async () => {
  153. console.log('Shutting down server...');
  154. // Close all active transports to properly clean up resources
  155. for (const sessionId in transports) {
  156. try {
  157. console.log(`Closing transport for session ${sessionId}`);
  158. await transports[sessionId].close();
  159. delete transports[sessionId];
  160. }
  161. catch (error) {
  162. console.error(`Error closing transport for session ${sessionId}:`, error);
  163. }
  164. }
  165. console.log('Server shutdown complete');
  166. process.exit(0);
  167. });
  168. //# sourceMappingURL=simpleSseServer.js.map