Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 25 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`) |
Expand Down
42 changes: 39 additions & 3 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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}`
])
Expand Down
115 changes: 115 additions & 0 deletions src/streamable.ts
Original file line number Diff line number Diff line change
@@ -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<void> {
// Нечего запускать: входящие сообщения подаются вручную через deliver().
}

async send(message: JSONRPCMessage): Promise<void> {
this.emit(message);
}

async close(): Promise<void> {
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<StreamableResult> {
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<string>();
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<void>(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);
}
}
Loading