simpleSseServer.js 5.2 KB

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