diff --git a/deno.lock b/deno.lock index fcb2e1c..e2aab00 100644 --- a/deno.lock +++ b/deno.lock @@ -1,6 +1,7 @@ { "version": "5", "specifiers": { + "jsr:@hono/hono@^4.10.5": "4.10.5", "jsr:@logtape/logtape@^1.2.0": "1.2.0", "jsr:@noble/ciphers@^2.0.1": "2.0.1", "jsr:@noble/curves@2.0": "2.0.1", @@ -28,6 +29,9 @@ "npm:zod@^3.25.76": "3.25.76" }, "jsr": { + "@hono/hono@4.10.5": { + "integrity": "13dbf2a528feb8189ad13394b213f0cf5f83b0ba4b2fadd0549993426db9ad2d" + }, "@logtape/logtape@1.2.0": { "integrity": "8e1d3af5c91966cc5689cfb17081a36bccfdff28ff6314769185661f5147e74d" }, @@ -53,12 +57,6 @@ "@puregarlic/randimal@1.1.1": { "integrity": "4e1fa61982cf2f610e9ad851d0fd0ff7bc3bb7b7a3c6cccae59f5ae2e68a7e47" }, - "@std/assert@1.0.14": { - "integrity": "68d0d4a43b365abc927f45a9b85c639ea18a9fab96ad92281e493e4ed84abaa4", - "dependencies": [ - "jsr:@std/internal@^1.0.10" - ] - }, "@std/assert@1.0.15": { "integrity": "d64018e951dbdfab9777335ecdb000c0b4e3df036984083be219ce5941e4703b", "dependencies": [ @@ -1724,6 +1722,7 @@ }, "packages/mcp": { "dependencies": [ + "jsr:@hono/hono@^4.10.5", "jsr:@logtape/logtape@^1.2.0", "jsr:@std/cli@^1.0.23", "npm:@modelcontextprotocol/sdk@^1.21.1", diff --git a/packages/mcp/deno.jsonc b/packages/mcp/deno.jsonc index 1ad32f0..158cf40 100644 --- a/packages/mcp/deno.jsonc +++ b/packages/mcp/deno.jsonc @@ -17,6 +17,7 @@ } }, "imports": { + "hono": "jsr:@hono/hono@^4.10.5", "@logtape/logtape": "jsr:@logtape/logtape@^1.2.0", "@modelcontextprotocol/sdk": "npm:@modelcontextprotocol/sdk@^1.21.1", "@std/cli": "jsr:@std/cli@^1.0.23", diff --git a/packages/mcp/hono.ts b/packages/mcp/hono.ts new file mode 100644 index 0000000..138166a --- /dev/null +++ b/packages/mcp/hono.ts @@ -0,0 +1,124 @@ +import { Hono } from "hono"; +import { cors } from "hono/cors"; +import { getLogger, withContext } from "@logtape/logtape"; +import { toFetchResponse, toReqRes } from "fetch-to-node"; +import { StreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/streamableHttp.js"; +import { createServer } from "./server.ts"; + +export function createApp() { + const app = new Hono(); + const logger = getLogger(["cistern", "http"]); + const sessions = new Map(); + + app.use("*", async (c, next) => { + const requestId = crypto.randomUUID(); + const startTime = Date.now(); + + await withContext({ + requestId, + method: c.req.method, + url: c.req.url, + userAgent: c.req.header("User-Agent"), + ipAddress: c.req.header("CF-Connecting-IP") || + c.req.header("X-Forwarded-For"), + }, async () => { + logger.info("{method} request started", { + method: c.req.method, + url: c.req.url, + requestId, + }); + + await next(); + + const duration = Date.now() - startTime; + + logger.info("{status} request completed in {duration}ms", { + status: c.res.status, + duration, + requestId, + }); + }); + }); + + app.onError((err, c) => { + logger.error("request error", { + error: { + name: err.name, + message: err.message, + stack: err.stack, + }, + method: c.req.method, + url: c.req.url, + }); + + return c.json({ error: "internal server error" }, 500); + }); + + app.all( + "/mcp", + cors({ + origin: "*", + allowMethods: ["GET", "POST", "DELETE", "OPTIONS"], + allowHeaders: [ + "Content-Type", + "Authorization", + "Mcp-Session-Id", + "Mcp-Protocol-Version", + ], + exposeHeaders: ["Mcp-Session-Id"], + }), + ); + + app.post("/mcp", async (ctx) => { + const sessionId = ctx.req.header("mcp-session-id") ?? crypto.randomUUID(); + let session = sessions.get(sessionId); + + if (session) { + logger.info("resuming session {sessionId}", { sessionId }); + } else { + logger.info("creating new session {sessionId}", { sessionId }); + + const server = createServer(); + + session = new StreamableHTTPServerTransport({ + sessionIdGenerator: () => sessionId, + }); + + session.onclose = () => { + logger.info("closing session {sessionId}", { sessionId }); + }; + + await server.connect(session); + + sessions.set(sessionId, session); + } + + const { req, res } = toReqRes(ctx.req.raw); + + await session.handleRequest(req, res); + + return await toFetchResponse(res); + }); + + app.on(["GET", "DELETE"], "/mcp", async (ctx) => { + const sessionId = ctx.req.header("mcp-session-id") ?? ""; + const session = sessions.get(sessionId); + + if (!session) { + logger.info("{method} invalid session {sessionId}", { + method: ctx.req.method, + sessionId, + }); + + return ctx.json({ error: "invalid or missing session" }, 401); + } + + const { req, res } = toReqRes(ctx.req.raw); + + await session.handleRequest(req, res); + + return await toFetchResponse(res); + }); + + return app; +} diff --git a/packages/mcp/index.ts b/packages/mcp/index.ts index cb4b727..3c314b4 100644 --- a/packages/mcp/index.ts +++ b/packages/mcp/index.ts @@ -1,35 +1,45 @@ import { parseArgs } from "@std/cli"; +import { AsyncLocalStorage } from "node:async_hooks"; import { configure, getConsoleSink, getLogger } from "@logtape/logtape"; import { StdioServerTransport } from "@modelcontextprotocol/sdk/server/stdio.js"; -import { StreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/streamableHttp.js"; -import { toFetchResponse, toReqRes } from "fetch-to-node"; import { createServer } from "./server.ts"; +import { createApp } from "./hono.ts"; async function main() { await configure({ sinks: { console: getConsoleSink() }, loggers: [ - { category: "cistern-mcp", lowestLevel: "trace", sinks: ["console"] }, + { + category: ["cistern", "mcp"], + lowestLevel: "trace", + sinks: ["console"], + }, + { + category: ["cistern", "http"], + lowestLevel: "info", + sinks: ["console"], + }, ], + contextLocalStorage: new AsyncLocalStorage(), }); - const logger = getLogger("cistern-mcp"); + const logger = getLogger(["cistern", "mcp"]); const args = parseArgs(Deno.args, { boolean: ["http"], }); - const server = createServer(); - if (!args.http) { logger.info("starting in stdio mode"); const transport = new StdioServerTransport(); + const server = createServer(); + await server.connect(transport); } else { logger.info("starting in streamable HTTP mode"); - const sessions: Map = new Map(); + const app = createApp(); Deno.serve( { @@ -38,75 +48,8 @@ async function main() { ...addr, }); }, - onError(error) { - logger.error( - "unexpected route error: {error}", - { error }, - ); - - return new Response(null, { status: 500 }); - }, - }, - async function handler(request: Request): Promise { - const PATH = new URLPattern({ pathname: "/mcp" }); - - if (!PATH.exec(request.url)) { - logger.info("not found", { - status: 404, - url: request.url, - }); - - return new Response(null, { status: 404 }); - } - - const sessionId = request.headers.get("mcp-session-id"); - let transport: StreamableHTTPServerTransport; - - if (sessionId && sessions.has(sessionId)) { - logger.info("{method} resuming session {sessionId}", { - sessionId, - method: request.method, - }); - - transport = sessions.get(sessionId)!; - } else if ( - request.method !== "POST" && !sessions.has(sessionId ?? "") - ) { - logger.error("{method} has invalid session {sessionId}", { - sessionId, - method: request.method, - }); - - return Response.json({ error: "invalid or missing session" }, { - status: 401, - }); - } else { - const sessionId = crypto.randomUUID(); - - logger.info("opening new session {sessionId}", { sessionId }); - - transport = new StreamableHTTPServerTransport({ - sessionIdGenerator: () => sessionId as string, - }); - - transport.onclose = () => { - logger.info("session {sessionId} closed, cleaning up", { - sessionId, - }); - sessions.delete(sessionId); - }; - - sessions.set(sessionId, transport); - - await server.connect(transport); - } - - const { req, res } = toReqRes(request); - - await transport.handleRequest(req, res); - - return await toFetchResponse(res); }, + app.fetch, ); } }