standaloneSseWithGetStreamableHttp.js 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124
  1. "use strict";
  2. Object.defineProperty(exports, "__esModule", { value: true });
  3. const node_crypto_1 = require("node:crypto");
  4. const mcp_js_1 = require("../../server/mcp.js");
  5. const streamableHttp_js_1 = require("../../server/streamableHttp.js");
  6. const types_js_1 = require("../../types.js");
  7. const express_js_1 = require("../../server/express.js");
  8. // Factory to create a new MCP server per session.
  9. // Each session needs its own server+transport pair to avoid cross-session contamination.
  10. const getServer = () => {
  11. const server = new mcp_js_1.McpServer({
  12. name: 'resource-list-changed-notification-server',
  13. version: '1.0.0'
  14. });
  15. const addResource = (name, content) => {
  16. const uri = `https://mcp-example.com/dynamic/${encodeURIComponent(name)}`;
  17. server.registerResource(name, uri, { mimeType: 'text/plain', description: `Dynamic resource: ${name}` }, async () => {
  18. return {
  19. contents: [{ uri, text: content }]
  20. };
  21. });
  22. };
  23. addResource('example-resource', 'Initial content for example-resource');
  24. // Periodically add new resources to demonstrate notifications
  25. const resourceChangeInterval = setInterval(() => {
  26. const name = (0, node_crypto_1.randomUUID)();
  27. addResource(name, `Content for ${name}`);
  28. }, 5000);
  29. // Clean up the interval when the server closes
  30. server.server.onclose = () => {
  31. clearInterval(resourceChangeInterval);
  32. };
  33. return server;
  34. };
  35. // Store transports by session ID to send notifications
  36. const transports = {};
  37. const app = (0, express_js_1.createMcpExpressApp)();
  38. app.post('/mcp', async (req, res) => {
  39. console.log('Received MCP request:', req.body);
  40. try {
  41. // Check for existing session ID
  42. const sessionId = req.headers['mcp-session-id'];
  43. let transport;
  44. if (sessionId && transports[sessionId]) {
  45. // Reuse existing transport
  46. transport = transports[sessionId];
  47. }
  48. else if (!sessionId && (0, types_js_1.isInitializeRequest)(req.body)) {
  49. // New initialization request
  50. transport = new streamableHttp_js_1.StreamableHTTPServerTransport({
  51. sessionIdGenerator: () => (0, node_crypto_1.randomUUID)(),
  52. onsessioninitialized: sessionId => {
  53. // Store the transport by session ID when session is initialized
  54. // This avoids race conditions where requests might come in before the session is stored
  55. console.log(`Session initialized with ID: ${sessionId}`);
  56. transports[sessionId] = transport;
  57. }
  58. });
  59. // Create a new server per session and connect it to the transport
  60. const server = getServer();
  61. await server.connect(transport);
  62. // Handle the request - the onsessioninitialized callback will store the transport
  63. await transport.handleRequest(req, res, req.body);
  64. return; // Already handled
  65. }
  66. else {
  67. // Invalid request - no session ID or not initialization request
  68. res.status(400).json({
  69. jsonrpc: '2.0',
  70. error: {
  71. code: -32000,
  72. message: 'Bad Request: No valid session ID provided'
  73. },
  74. id: null
  75. });
  76. return;
  77. }
  78. // Handle the request with existing transport
  79. await transport.handleRequest(req, res, req.body);
  80. }
  81. catch (error) {
  82. console.error('Error handling MCP request:', error);
  83. if (!res.headersSent) {
  84. res.status(500).json({
  85. jsonrpc: '2.0',
  86. error: {
  87. code: -32603,
  88. message: 'Internal server error'
  89. },
  90. id: null
  91. });
  92. }
  93. }
  94. });
  95. // Handle GET requests for SSE streams (now using built-in support from StreamableHTTP)
  96. app.get('/mcp', async (req, res) => {
  97. const sessionId = req.headers['mcp-session-id'];
  98. if (!sessionId || !transports[sessionId]) {
  99. res.status(400).send('Invalid or missing session ID');
  100. return;
  101. }
  102. console.log(`Establishing SSE stream for session ${sessionId}`);
  103. const transport = transports[sessionId];
  104. await transport.handleRequest(req, res);
  105. });
  106. // Start the server
  107. const PORT = 3000;
  108. app.listen(PORT, error => {
  109. if (error) {
  110. console.error('Failed to start server:', error);
  111. process.exit(1);
  112. }
  113. console.log(`Server listening on port ${PORT}`);
  114. });
  115. // Handle server shutdown
  116. process.on('SIGINT', async () => {
  117. console.log('Shutting down server...');
  118. for (const sessionId in transports) {
  119. await transports[sessionId].close();
  120. delete transports[sessionId];
  121. }
  122. process.exit(0);
  123. });
  124. //# sourceMappingURL=standaloneSseWithGetStreamableHttp.js.map