standaloneSseWithGetStreamableHttp.js 4.7 KB

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