jsonResponseStreamableHttp.js 6.4 KB

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