jsonResponseStreamableHttp.js 5.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146
  1. import { randomUUID } from 'node:crypto';
  2. import { McpServer } from '../../server/mcp.js';
  3. import { StreamableHTTPServerTransport } from '../../server/streamableHttp.js';
  4. import * as z from 'zod/v4';
  5. import { isInitializeRequest } from '../../types.js';
  6. import { createMcpExpressApp } from '../../server/express.js';
  7. // Create an MCP server with implementation details
  8. const getServer = () => {
  9. const server = new McpServer({
  10. name: 'json-response-streamable-http-server',
  11. version: '1.0.0'
  12. }, {
  13. capabilities: {
  14. logging: {}
  15. }
  16. });
  17. // Register a simple tool that returns a greeting
  18. server.registerTool('greet', {
  19. description: 'A simple greeting tool',
  20. inputSchema: {
  21. name: z.string().describe('Name to greet')
  22. }
  23. }, async ({ name }) => {
  24. return {
  25. content: [
  26. {
  27. type: 'text',
  28. text: `Hello, ${name}!`
  29. }
  30. ]
  31. };
  32. });
  33. // Register a tool that sends multiple greetings with notifications
  34. server.registerTool('multi-greet', {
  35. description: 'A tool that sends different greetings with delays between them',
  36. inputSchema: {
  37. name: z.string().describe('Name to greet')
  38. }
  39. }, async ({ name }, extra) => {
  40. const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms));
  41. await server.sendLoggingMessage({
  42. level: 'debug',
  43. data: `Starting multi-greet for ${name}`
  44. }, extra.sessionId);
  45. await sleep(1000); // Wait 1 second before first greeting
  46. await server.sendLoggingMessage({
  47. level: 'info',
  48. data: `Sending first greeting to ${name}`
  49. }, extra.sessionId);
  50. await sleep(1000); // Wait another second before second greeting
  51. await server.sendLoggingMessage({
  52. level: 'info',
  53. data: `Sending second greeting to ${name}`
  54. }, extra.sessionId);
  55. return {
  56. content: [
  57. {
  58. type: 'text',
  59. text: `Good morning, ${name}!`
  60. }
  61. ]
  62. };
  63. });
  64. return server;
  65. };
  66. const app = createMcpExpressApp();
  67. // Map to store transports by session ID
  68. const transports = {};
  69. app.post('/mcp', async (req, res) => {
  70. console.log('Received MCP request:', req.body);
  71. try {
  72. // Check for existing session ID
  73. const sessionId = req.headers['mcp-session-id'];
  74. let transport;
  75. if (sessionId && transports[sessionId]) {
  76. // Reuse existing transport
  77. transport = transports[sessionId];
  78. }
  79. else if (!sessionId && isInitializeRequest(req.body)) {
  80. // New initialization request - use JSON response mode
  81. transport = new StreamableHTTPServerTransport({
  82. sessionIdGenerator: () => randomUUID(),
  83. enableJsonResponse: true, // Enable JSON response mode
  84. onsessioninitialized: sessionId => {
  85. // Store the transport by session ID when session is initialized
  86. // This avoids race conditions where requests might come in before the session is stored
  87. console.log(`Session initialized with ID: ${sessionId}`);
  88. transports[sessionId] = transport;
  89. }
  90. });
  91. // Connect the transport to the MCP server BEFORE handling the request
  92. const server = getServer();
  93. await server.connect(transport);
  94. await transport.handleRequest(req, res, req.body);
  95. return; // Already handled
  96. }
  97. else {
  98. // Invalid request - no session ID or not initialization request
  99. res.status(400).json({
  100. jsonrpc: '2.0',
  101. error: {
  102. code: -32000,
  103. message: 'Bad Request: No valid session ID provided'
  104. },
  105. id: null
  106. });
  107. return;
  108. }
  109. // Handle the request with existing transport - no need to reconnect
  110. await transport.handleRequest(req, res, req.body);
  111. }
  112. catch (error) {
  113. console.error('Error handling MCP request:', error);
  114. if (!res.headersSent) {
  115. res.status(500).json({
  116. jsonrpc: '2.0',
  117. error: {
  118. code: -32603,
  119. message: 'Internal server error'
  120. },
  121. id: null
  122. });
  123. }
  124. }
  125. });
  126. // Handle GET requests for SSE streams according to spec
  127. app.get('/mcp', async (req, res) => {
  128. // Since this is a very simple example, we don't support GET requests for this server
  129. // The spec requires returning 405 Method Not Allowed in this case
  130. res.status(405).set('Allow', 'POST').send('Method Not Allowed');
  131. });
  132. // Start the server
  133. const PORT = 3000;
  134. app.listen(PORT, error => {
  135. if (error) {
  136. console.error('Failed to start server:', error);
  137. process.exit(1);
  138. }
  139. console.log(`MCP Streamable HTTP Server listening on port ${PORT}`);
  140. });
  141. // Handle server shutdown
  142. process.on('SIGINT', async () => {
  143. console.log('Shutting down server...');
  144. process.exit(0);
  145. });
  146. //# sourceMappingURL=jsonResponseStreamableHttp.js.map