{"version":3,"sources":["../../../../src/adapters/mcp/serve/http_server.ts"],"names":["logger","Logger","root","child","name","runServer","server","hostname","port","app","express","use","json","transports","all","req","res","debug","method","sessionId","headers","transport","existingTransport","StreamableHTTPServerTransport","status","jsonrpc","error","code","message","id","isInitializeRequest","body","eventStore","InMemoryEventStore","sessionIdGenerator","randomUUID","onsessioninitialized","onclose","sid","connect","handleRequest","headersSent","get","info","SSEServerTransport","on","post","query","handlePostMessage","send","listen","process","exit","close"],"mappings":";;;;;;;;;;;;;;;;AAgBA,MAAMA,MAAAA,GAASC,iBAAAA,CAAOC,IAAAA,CAAKC,KAAAA,CAAM;EAC/BC,IAAAA,EAAM;AACR,CAAA,CAAA;AAEO,SAASC,SAAAA,CAAUC,MAAAA,EAAmBC,QAAAA,GAAW,WAAA,EAAaC,OAAO,GAAA,EAAI;AAE9E,EAAA,MAAMC,MAAMC,wBAAAA,EAAAA;AACZD,EAAAA,GAAAA,CAAIE,GAAAA,CAAID,wBAAAA,CAAQE,IAAAA,EAAI,CAAA;AAGpB,EAAA,MAAMC,aAAiF,EAAC;AAOxFJ,EAAAA,GAAAA,CAAIK,GAAAA,CAAI,MAAA,EAAQ,OAAOC,GAAAA,EAAcC,GAAAA,KAAAA;AACnChB,IAAAA,MAAAA,CAAOiB,KAAAA,CAAM,CAAA,SAAA,EAAYF,GAAAA,CAAIG,MAAM,CAAA,gBAAA,CAAkB,CAAA;AAErD,IAAA,IAAI;AAEF,MAAA,MAAMC,SAAAA,GAAYJ,GAAAA,CAAIK,OAAAA,CAAQ,gBAAA,CAAA;AAC9B,MAAA,IAAIC,SAAAA;AAEJ,MAAA,IAAIF,SAAAA,IAAaN,UAAAA,CAAWM,SAAAA,CAAAA,EAAY;AAEtC,QAAA,MAAMG,iBAAAA,GAAoBT,WAAWM,SAAAA,CAAAA;AACrC,QAAA,IAAIG,6BAA6BC,+CAAAA,EAA+B;AAE9DF,UAAAA,SAAAA,GAAYC,iBAAAA;QACd,CAAA,MAAO;AAELN,UAAAA,GAAAA,CAAIQ,MAAAA,CAAO,GAAA,CAAA,CAAKZ,IAAAA,CAAK;YACnBa,OAAAA,EAAS,KAAA;YACTC,KAAAA,EAAO;cACLC,IAAAA,EAAM,CAAA,IAAA;cACNC,OAAAA,EAAS;AACX,aAAA;YACAC,EAAAA,EAAI;WACN,CAAA;AACA,UAAA;AACF,QAAA;MACF,CAAA,MAAA,IAAW,CAACV,aAAaJ,GAAAA,CAAIG,MAAAA,KAAW,UAAUY,4BAAAA,CAAoBf,GAAAA,CAAIgB,IAAI,CAAA,EAAG;AAC/E,QAAA,MAAMC,UAAAA,GAAa,IAAIC,sCAAAA,EAAAA;AACvBZ,QAAAA,SAAAA,GAAY,IAAIE,+CAAAA,CAA8B;UAC5CW,kBAAAA,kBAAoB,MAAA,CAAA,MAAMC,wBAAAA,EAAN,oBAAA,CAAA;AACpBH,UAAAA,UAAAA;;AACAI,UAAAA,oBAAAA,0BAAuBjB,UAAAA,KAAAA;AAErBnB,YAAAA,MAAAA,CAAOiB,KAAAA,CAAM,CAAA,4CAAA,EAA+CE,UAAAA,CAAAA,CAAW,CAAA;AACvEN,YAAAA,UAAAA,CAAWM,UAAAA,CAAAA,GAAaE,SAAAA;UAC1B,CAAA,EAJsB,sBAAA;SAKxB,CAAA;AAGAA,QAAAA,SAAAA,CAAUgB,UAAU,MAAA;AAClB,UAAA,MAAMC,MAAMjB,SAAAA,CAAUF,SAAAA;AACtB,UAAA,IAAImB,GAAAA,IAAOzB,UAAAA,CAAWyB,GAAAA,CAAAA,EAAM;AAC1BtC,YAAAA,MAAAA,CAAOiB,KAAAA,CAAM,CAAA,6BAAA,EAAgCqB,GAAAA,CAAAA,8BAAAA,CAAmC,CAAA;AAEhF,YAAA,OAAOzB,WAAWyB,GAAAA,CAAAA;AACpB,UAAA;AACF,QAAA,CAAA;AAGA,QAAA,MAAMhC,MAAAA,CAAOiC,QAAQlB,SAAAA,CAAAA;MACvB,CAAA,MAAO;AAELL,QAAAA,GAAAA,CAAIQ,MAAAA,CAAO,GAAA,CAAA,CAAKZ,IAAAA,CAAK;UACnBa,OAAAA,EAAS,KAAA;UACTC,KAAAA,EAAO;YACLC,IAAAA,EAAM,CAAA,IAAA;YACNC,OAAAA,EAAS;AACX,WAAA;UACAC,EAAAA,EAAI;SACN,CAAA;AACA,QAAA;AACF,MAAA;AAGA,MAAA,MAAMR,SAAAA,CAAUmB,aAAAA,CAAczB,GAAAA,EAAKC,GAAAA,EAAKD,IAAIgB,IAAI,CAAA;AAClD,IAAA,CAAA,CAAA,OAASL,KAAAA,EAAO;AACd1B,MAAAA,MAAAA,CAAO0B,KAAAA,CAAM,+BAA+BA,KAAAA,CAAAA;AAC5C,MAAA,IAAI,CAACV,IAAIyB,WAAAA,EAAa;AACpBzB,QAAAA,GAAAA,CAAIQ,MAAAA,CAAO,GAAA,CAAA,CAAKZ,IAAAA,CAAK;UACnBa,OAAAA,EAAS,KAAA;UACTC,KAAAA,EAAO;YACLC,IAAAA,EAAM,MAAA;YACNC,OAAAA,EAAS;AACX,WAAA;UACAC,EAAAA,EAAI;SACN,CAAA;AACF,MAAA;AACF,IAAA;EACF,CAAA,CAAA;AAMApB,EAAAA,GAAAA,CAAIiC,GAAAA,CAAI,MAAA,EAAQ,OAAO3B,GAAAA,EAAcC,GAAAA,KAAAA;AACnChB,IAAAA,MAAAA,CAAO2C,KAAK,yDAAA,CAAA;AACZ,IAAA,MAAMtB,SAAAA,GAAY,IAAIuB,yBAAAA,CAAmB,WAAA,EAAa5B,GAAAA,CAAAA;AACtDH,IAAAA,UAAAA,CAAWQ,SAAAA,CAAUF,SAAS,CAAA,GAAIE,SAAAA;AAClCL,IAAAA,GAAAA,CAAI6B,EAAAA,CAAG,SAAS,MAAA;AAEd,MAAA,OAAOhC,UAAAA,CAAWQ,UAAUF,SAAS,CAAA;IACvC,CAAA,CAAA;AACA,IAAA,MAAMb,MAAAA,CAAOiC,QAAQlB,SAAAA,CAAAA;EACvB,CAAA,CAAA;AAEAZ,EAAAA,GAAAA,CAAIqC,IAAAA,CAAK,WAAA,EAAa,OAAO/B,GAAAA,EAAcC,GAAAA,KAAAA;AACzC,IAAA,MAAMG,SAAAA,GAAYJ,IAAIgC,KAAAA,CAAM5B,SAAAA;AAC5B,IAAA,IAAIE,SAAAA;AACJ,IAAA,MAAMC,iBAAAA,GAAoBT,WAAWM,SAAAA,CAAAA;AACrC,IAAA,IAAIG,6BAA6BsB,yBAAAA,EAAoB;AAEnDvB,MAAAA,SAAAA,GAAYC,iBAAAA;IACd,CAAA,MAAO;AAELN,MAAAA,GAAAA,CAAIQ,MAAAA,CAAO,GAAA,CAAA,CAAKZ,IAAAA,CAAK;QACnBa,OAAAA,EAAS,KAAA;QACTC,KAAAA,EAAO;UACLC,IAAAA,EAAM,KAAA;UACNC,OAAAA,EAAS;AACX,SAAA;QACAC,EAAAA,EAAI;OACN,CAAA;AACA,MAAA;AACF,IAAA;AACA,IAAA,IAAIR,SAAAA,EAAW;AACb,MAAA,MAAMA,SAAAA,CAAU2B,iBAAAA,CAAkBjC,GAAAA,EAAKC,GAAAA,EAAKD,IAAIgB,IAAI,CAAA;IACtD,CAAA,MAAO;AACLf,MAAAA,GAAAA,CAAIQ,MAAAA,CAAO,GAAA,CAAA,CAAKyB,IAAAA,CAAK,kCAAA,CAAA;AACvB,IAAA;EACF,CAAA,CAAA;AAGAxC,EAAAA,GAAAA,CAAIyC,MAAAA,CAAO1C,IAAAA,EAAMD,QAAAA,EAAU,CAACmB,KAAAA,KAAAA;AAC1B,IAAA,IAAIA,KAAAA,EAAO;AACT1B,MAAAA,MAAAA,CAAO0B,KAAAA,CAAMA,OAAO,uBAAA,CAAA;AACpByB,MAAAA,OAAAA,CAAQC,KAAK,CAAA,CAAA;AACf,IAAA;AACApD,IAAAA,MAAAA,CAAO2C,IAAAA,CAAK,CAAA,kDAAA,EAAqDpC,QAAAA,CAAAA,CAAAA,EAAYC,IAAAA,CAAAA,CAAM,CAAA;AACnFR,IAAAA,MAAAA,CAAOiB,KAAAA,CAAM;;;;;;;;;;;;;;;;;;;AAmBZ,IAAA,CAAA,CAAA;EACH,CAAA,CAAA;AAGAkC,EAAAA,OAAAA,CAAQN,EAAAA,CAAG,UAAU,YAAA;AACnB7C,IAAAA,MAAAA,CAAO2C,KAAK,yBAAA,CAAA;AAGZ,IAAA,KAAA,MAAWxB,aAAaN,UAAAA,EAAY;AAClC,MAAA,IAAI;AACFb,QAAAA,MAAAA,CAAOiB,KAAAA,CAAM,CAAA,8BAAA,EAAiCE,SAAAA,CAAAA,CAAW,CAAA;AACzD,QAAA,MAAMN,UAAAA,CAAWM,SAAAA,CAAAA,CAAWkC,KAAAA,EAAK;AAEjC,QAAA,OAAOxC,WAAWM,SAAAA,CAAAA;AACpB,MAAA,CAAA,CAAA,OAASO,KAAAA,EAAO;AACd1B,QAAAA,MAAAA,CAAO0B,KAAAA,CAAM,CAAA,oCAAA,EAAuCP,SAAAA,CAAAA,CAAAA,CAAAA,EAAcO,KAAAA,CAAAA;AACpE,MAAA;AACF,IAAA;AACA1B,IAAAA,MAAAA,CAAOiB,MAAM,0BAAA,CAAA;AACbkC,IAAAA,OAAAA,CAAQC,KAAK,CAAA,CAAA;EACf,CAAA,CAAA;AACF;AArLgB/C,MAAAA,CAAAA,SAAAA,EAAAA,WAAAA,CAAAA","file":"http_server.cjs","sourcesContent":["/**\n * Copyright 2025 © BeeAI a Series of LF Projects, LLC\n * SPDX-License-Identifier: Apache-2.0\n */\n\n// Taken from: https://github.com/modelcontextprotocol/typescript-sdk/blob/main/src/examples/server/sseAndStreamableHttpCompatibleServer.ts\n\nimport express, { Request, Response } from \"express\";\nimport { SSEServerTransport } from \"@modelcontextprotocol/sdk/server/sse.js\";\nimport { StreamableHTTPServerTransport } from \"@modelcontextprotocol/sdk/server/streamableHttp.js\";\nimport { isInitializeRequest } from \"@modelcontextprotocol/sdk/types.js\";\nimport { InMemoryEventStore } from \"./in_memory_store.js\";\nimport { randomUUID } from \"node:crypto\";\nimport { McpServer } from \"@modelcontextprotocol/sdk/server/mcp.js\";\nimport { Logger } from \"@/logger/logger.js\";\n\nconst logger = Logger.root.child({\n  name: \"MCP HTTP server\",\n});\n\nexport function runServer(server: McpServer, hostname = \"127.0.0.1\", port = 3000) {\n  // Create Express application\n  const app = express();\n  app.use(express.json());\n\n  // Store transports by session ID\n  const transports: Record<string, StreamableHTTPServerTransport | SSEServerTransport> = {};\n\n  //=============================================================================\n  // STREAMABLE HTTP TRANSPORT (PROTOCOL VERSION 2025-03-26)\n  //=============================================================================\n\n  // Handle all MCP Streamable HTTP requests (GET, POST, DELETE) on a single endpoint\n  app.all(\"/mcp\", async (req: Request, res: Response) => {\n    logger.debug(`Received ${req.method} request to /mcp`);\n\n    try {\n      // Check for existing session ID\n      const sessionId = req.headers[\"mcp-session-id\"] as string | undefined;\n      let transport: StreamableHTTPServerTransport;\n\n      if (sessionId && transports[sessionId]) {\n        // Check if the transport is of the correct type\n        const existingTransport = transports[sessionId];\n        if (existingTransport instanceof StreamableHTTPServerTransport) {\n          // Reuse existing transport\n          transport = existingTransport;\n        } else {\n          // Transport exists but is not a StreamableHTTPServerTransport (could be SSEServerTransport)\n          res.status(400).json({\n            jsonrpc: \"2.0\",\n            error: {\n              code: -32000,\n              message: \"Bad Request: Session exists but uses a different transport protocol\",\n            },\n            id: null,\n          });\n          return;\n        }\n      } else if (!sessionId && req.method === \"POST\" && isInitializeRequest(req.body)) {\n        const eventStore = new InMemoryEventStore();\n        transport = new StreamableHTTPServerTransport({\n          sessionIdGenerator: () => randomUUID(),\n          eventStore, // Enable resumability\n          onsessioninitialized: (sessionId) => {\n            // Store the transport by session ID when session is initialized\n            logger.debug(`StreamableHTTP session initialized with ID: ${sessionId}`);\n            transports[sessionId] = transport;\n          },\n        });\n\n        // Set up onclose handler to clean up transport when closed\n        transport.onclose = () => {\n          const sid = transport.sessionId;\n          if (sid && transports[sid]) {\n            logger.debug(`Transport closed for session ${sid}, removing from transports map`);\n            // eslint-disable-next-line @typescript-eslint/no-dynamic-delete\n            delete transports[sid];\n          }\n        };\n\n        // Connect the transport to the MCP server\n        await server.connect(transport);\n      } else {\n        // Invalid request - no session ID or not initialization request\n        res.status(400).json({\n          jsonrpc: \"2.0\",\n          error: {\n            code: -32000,\n            message: \"Bad Request: No valid session ID provided\",\n          },\n          id: null,\n        });\n        return;\n      }\n\n      // Handle the request with the transport\n      await transport.handleRequest(req, res, req.body);\n    } catch (error) {\n      logger.error(\"Error handling MCP request:\", error);\n      if (!res.headersSent) {\n        res.status(500).json({\n          jsonrpc: \"2.0\",\n          error: {\n            code: -32603,\n            message: \"Internal server error\",\n          },\n          id: null,\n        });\n      }\n    }\n  });\n\n  //=============================================================================\n  // DEPRECATED HTTP+SSE TRANSPORT (PROTOCOL VERSION 2024-11-05)\n  //=============================================================================\n\n  app.get(\"/sse\", async (req: Request, res: Response) => {\n    logger.info(\"Received GET request to /sse (deprecated SSE transport)\");\n    const transport = new SSEServerTransport(\"/messages\", res);\n    transports[transport.sessionId] = transport;\n    res.on(\"close\", () => {\n      // eslint-disable-next-line @typescript-eslint/no-dynamic-delete\n      delete transports[transport.sessionId];\n    });\n    await server.connect(transport);\n  });\n\n  app.post(\"/messages\", async (req: Request, res: Response) => {\n    const sessionId = req.query.sessionId as string;\n    let transport: SSEServerTransport;\n    const existingTransport = transports[sessionId];\n    if (existingTransport instanceof SSEServerTransport) {\n      // Reuse existing transport\n      transport = existingTransport;\n    } else {\n      // Transport exists but is not a SSEServerTransport (could be StreamableHTTPServerTransport)\n      res.status(400).json({\n        jsonrpc: \"2.0\",\n        error: {\n          code: -32000,\n          message: \"Bad Request: Session exists but uses a different transport protocol\",\n        },\n        id: null,\n      });\n      return;\n    }\n    if (transport) {\n      await transport.handlePostMessage(req, res, req.body);\n    } else {\n      res.status(400).send(\"No transport found for sessionId\");\n    }\n  });\n\n  // Start the server\n  app.listen(port, hostname, (error) => {\n    if (error) {\n      logger.error(error, \"Error starting server\");\n      process.exit(1);\n    }\n    logger.info(`Backwards compatible MCP server listening on port ${hostname}:${port}`);\n    logger.debug(`\n    ==============================================\n    SUPPORTED TRANSPORT OPTIONS:\n\n    1. Streamable Http(Protocol version: 2025-03-26)\n    Endpoint: /mcp\n    Methods: GET, POST, DELETE\n    Usage: \n        - Initialize with POST to /mcp\n        - Establish SSE stream with GET to /mcp\n        - Send requests with POST to /mcp\n        - Terminate session with DELETE to /mcp\n\n    2. Http + SSE (Protocol version: 2024-11-05)\n    Endpoints: /sse (GET) and /messages (POST)\n    Usage:\n        - Establish SSE stream with GET to /sse\n        - Send requests with POST to /messages?sessionId=<id>\n    ==============================================\n    `);\n  });\n\n  // Handle server shutdown\n  process.on(\"SIGINT\", async () => {\n    logger.info(\"Shutting down server...\");\n\n    // Close all active transports to properly clean up resources\n    for (const sessionId in transports) {\n      try {\n        logger.debug(`Closing transport for session ${sessionId}`);\n        await transports[sessionId].close();\n        // eslint-disable-next-line @typescript-eslint/no-dynamic-delete\n        delete transports[sessionId];\n      } catch (error) {\n        logger.error(`Error closing transport for session ${sessionId}:`, error);\n      }\n    }\n    logger.debug(\"Server shutdown complete\");\n    process.exit(0);\n  });\n}\n"]}