diff --git a/README.md b/README.md index 9d64eaa..93997f9 100644 --- a/README.md +++ b/README.md @@ -89,12 +89,36 @@ bun run src/index.ts mode bypass --- +## 🔌 Транспорты + +Сервер отдаёт один и тот же набор тулов через два эндпоинта. + +| Эндпоинт | Транспорт | Когда использовать | +| --- | --- | --- | +| `POST /mcp` | Streamable HTTP (stateless) | Рекомендуется для всех современных клиентов и реверс-прокси | +| `GET /sse` + `POST /messages` | Legacy SSE | Только для старых клиентов без поддержки Streamable HTTP | + +Почему `/mcp` надёжнее: каждый запрос самодостаточен, сессия не хранится между вызовами, +и обрыв соединения не приводит к `404 Session not found` на следующем вызове тула — типичная проблема +долгоживущих SSE-стримов за Cloudflare-туннелем или другим реверс-прокси. + +Проверка подключения: + +```bash +curl -X POST http://127.0.0.1:3000/mcp \ + -H "Authorization: Bearer $TOKEN" \ + -H "Content-Type: application/json" \ + -d '{"jsonrpc":"2.0","id":1,"method":"tools/list"}' +``` + +--- + ## 📜 Команды CLI | Команда | Описание | | --- | --- | | `bun run setup` | Первичная настройка и параметры подключения | -| `bun run start` | Запуск MCP-сервера (SSE + Bearer) | +| `bun run start` | Запуск MCP-сервера (Streamable HTTP + legacy SSE, Bearer) | | `bun run dev` | Запуск в режиме разработки (watch mode) | | `bun run status` | Статус сервера, текущий режим и лимиты | | `bun run token` | Посмотреть или пересоздать токен (`--reset`) | diff --git a/src/index.ts b/src/index.ts index 9788b5d..e5ac77d 100644 --- a/src/index.ts +++ b/src/index.ts @@ -3,6 +3,7 @@ import { SSEServerTransport } from "@modelcontextprotocol/sdk/server/sse.js"; import type { IncomingMessage, ServerResponse } from "node:http"; import { resolve } from "node:path"; import { createServer, toolCount } from "@/mcp"; +import { handleStreamableRequest, jsonRpcError } from "@/streamable"; import { CONFIG_FILE, getProfile, @@ -211,7 +212,7 @@ function startServer(config: NotCodeConfig): void { .onRequest(({ request, set }) => { set.headers["Access-Control-Allow-Origin"] = "*"; set.headers["Access-Control-Allow-Headers"] = "*"; - set.headers["Access-Control-Allow-Methods"] = "GET, POST, OPTIONS"; + set.headers["Access-Control-Allow-Methods"] = "GET, POST, DELETE, OPTIONS"; if (request.method === "OPTIONS") return; if (new URL(request.url).pathname === "/health") return; @@ -239,6 +240,40 @@ function startServer(config: NotCodeConfig): void { watchers: watchers.list().length, workspaceRoot: getWorkspaceRoot(config) })) + // Современный транспорт (MCP Streamable HTTP), stateless-режим. + // Каждый POST самодостаточен: нет долгоживущего стрима и нет sessionId, поэтому + // обрыв соединения больше не ломает последующие вызовы тулов. + .post("/mcp", async ({ body, set }) => { + let payload: unknown = body; + + if (typeof payload === "string") { + try { + payload = JSON.parse(payload); + } catch { + const parseError = jsonRpcError(400, -32700, "Parse error: request body is not valid JSON"); + set.status = parseError.status; + return parseError.body; + } + } + + const result = await handleStreamableRequest(createServer, payload, config.limits.execTimeoutMs + 30_000); + + set.status = result.status; + if (result.body === null) return ""; + return result.body; + }) + // Серверные пуши не нужны: stateless-режим отвечает на каждый POST синхронно. + .get("/mcp", ({ set }) => { + set.status = 405; + set.headers["Allow"] = "POST, DELETE, OPTIONS"; + return jsonRpcError(405, -32000, "Method Not Allowed: this endpoint is stateless, use POST /mcp").body; + }) + // Клиент может явно закрыть сессию; закрывать нечего, но отвечаем по спеке. + .delete("/mcp", ({ set }) => { + set.status = 204; + return ""; + }) + // Legacy SSE оставлен для клиентов, которые ещё не умеют Streamable HTTP. .get("/sse", ({ request, set }) => { let streamController: ReadableStreamDefaultController; const body = new ReadableStream({ @@ -376,12 +411,13 @@ function startServer(config: NotCodeConfig): void { `🛡️ Режим: ${config.mode.toUpperCase()}`, `📁 Воркспейс: ${config.activeProfile} → ${getWorkspaceRoot(config)}`, `🧰 Тулов: ${toolCount()} (fs / terminal-сессии / git / meta)`, - `🌐 SSE: http://${config.host}:${config.port}/sse`, + `🌐 MCP: http://${config.host}:${config.port}/mcp (Streamable HTTP, рекомендуется)`, + `🕰️ SSE (legacy): http://${config.host}:${config.port}/sse`, `❤️ Health: http://${config.host}:${config.port}/health`, `🔑 Bearer Token: ${config.token}`, "-", "📌 Подключение MCP-клиента (Notion AI / Claude / Cursor):", - " 1. MCP server URL: адрес SSE выше (или HTTPS-адрес твоего реверс-прокси)", + " 1. MCP server URL: адрес /mcp выше (или HTTPS-адрес твоего реверс-прокси)", " 2. Authentication: Bearer token", ` 3. Token: ${config.token}` ]) diff --git a/src/streamable.ts b/src/streamable.ts new file mode 100644 index 0000000..38770c4 --- /dev/null +++ b/src/streamable.ts @@ -0,0 +1,115 @@ +import type { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; +import type { Transport } from "@modelcontextprotocol/sdk/shared/transport.js"; +import type { JSONRPCMessage } from "@modelcontextprotocol/sdk/types.js"; + +/** + * Транспорт без сети: сообщение приходит прямо из HTTP-хендлера, ответ уезжает в колбэк. + * + * Нужен для stateless Streamable HTTP: на каждый POST /mcp поднимается свой MCP-сервер, + * поэтому нет ни долгоживущего SSE-стрима, ни sessionId, который может протухнуть. + */ +class InlineTransport implements Transport { + onclose?: () => void; + onerror?: (error: Error) => void; + onmessage?: Transport["onmessage"]; + sessionId?: string; + + constructor(private readonly emit: (message: JSONRPCMessage) => void) {} + + async start(): Promise { + // Нечего запускать: входящие сообщения подаются вручную через deliver(). + } + + async send(message: JSONRPCMessage): Promise { + this.emit(message); + } + + async close(): Promise { + this.onclose?.(); + } + + deliver(message: JSONRPCMessage): void { + this.onmessage?.(message); + } +} + +export type StreamableResult = { + status: number; + /** null — тело не нужно (например, 202 на чистые нотификации). */ + body: unknown | null; +}; + +function messageId(message: unknown): string | null { + const id = (message as { id?: unknown } | null)?.id; + return id === undefined || id === null ? null : String(id); +} + +export function jsonRpcError(status: number, code: number, message: string): StreamableResult { + return { + status, + body: { jsonrpc: "2.0", id: null, error: { code, message } } + }; +} + +/** + * Обрабатывает один POST /mcp по спеке Streamable HTTP в stateless-режиме. + * + * На каждый запрос создаётся свежий MCP-сервер, поэтому клиенту не нужен mcp-session-id, + * а обрыв соединения больше не приводит к 404 "Session not found", как это было с legacy SSE. + */ +export async function handleStreamableRequest( + createServer: () => McpServer, + payload: unknown, + timeoutMs: number +): Promise { + const isBatch = Array.isArray(payload); + const messages = (isBatch ? payload : [payload]) as JSONRPCMessage[]; + + if (messages.length === 0 || messages.some(message => typeof message !== "object" || message === null)) { + return jsonRpcError(400, -32600, "Invalid Request: expected a JSON-RPC message or a non-empty batch"); + } + + const pending = new Set(); + for (const message of messages) { + const id = messageId(message); + if (id !== null) pending.add(id); + } + + const responses: JSONRPCMessage[] = []; + let settle: () => void = () => undefined; + const completed = new Promise(resolve => { + settle = resolve; + }); + + const transport = new InlineTransport(message => { + const id = messageId(message); + // Серверные нотификации в stateless-режиме отдавать некуда — ответом идут только результаты запросов. + if (id === null) return; + responses.push(message); + pending.delete(id); + if (pending.size === 0) settle(); + }); + + const server = createServer(); + await server.connect(transport); + + try { + const expectsResponse = pending.size > 0; + for (const message of messages) transport.deliver(message); + + // Только нотификации (например, notifications/initialized) — по спеке отвечаем 202 без тела. + if (!expectsResponse) return { status: 202, body: null }; + + const timer = setTimeout(settle, timeoutMs); + await completed; + clearTimeout(timer); + + if (responses.length === 0) { + return jsonRpcError(504, -32001, `Request timed out after ${timeoutMs} ms`); + } + + return { status: 200, body: isBatch ? responses : responses[0] }; + } finally { + await server.close().catch(() => undefined); + } +}