/** * MUA — Mikes Unraid Agent * index.ts — MCP over HTTP (Streamable HTTP) Server * * Portiert von mcp/server.php. Nutzt Bun.serve() statt php -S — * production-grade HTTP-Server, POST-Body wird korrekt gelesen. * * Transport: MCP Streamable HTTP (POST-only, JSON-RPC 2.0). * Endpunkte: * POST /mcp — JSON-RPC Request * GET /mcp — 405 (SSE nicht implementiert, spec-konform) * DELETE /mcp — Session-End * GET /health — Health-Check */ import { MUA_SERVER_NAME, MUA_VERSION, MUA_PROTOCOL_VERSION, } from "./helpers"; import { TOOLS, toolByName, getToolRisk } from "./tools"; import { checkAuth, isToolEnabled, getEnabledTools, getAllToolsEnabled, getMaskedApiKey, generateApiKey, setEnabledTools, } from "./auth"; import { randomBytes } from "node:crypto"; import { spawn } from "node:child_process"; import { appendFile, chmodSync, statSync, writeFileSync } from "node:fs"; const PORT = Number(process.env["MUA_PORT"] ?? 3002); const HOST = process.env["MUA_HOST"] ?? "0.0.0.0"; const MAX_REQUEST_BYTES = Number(process.env["MUA_MAX_REQUEST_BYTES"] ?? 1_048_576); const RATE_LIMIT_PER_MINUTE = Number(process.env["MUA_RATE_LIMIT_PER_MINUTE"] ?? 120); const MAX_CONCURRENT_TOOL_CALLS = Number(process.env["MUA_MAX_CONCURRENT_TOOL_CALLS"] ?? 4); const CORS_ORIGIN = process.env["MUA_CORS_ORIGIN"] ?? ""; const AUDIT_LOG = process.env["MUA_AUDIT_LOG"] ?? "/var/log/plugins/mua-audit.log"; // ── JSON-RPC Helpers ──────────────────────────────────────────────────── function rpcResult(id: number | string | null, result: unknown) { return { jsonrpc: "2.0", id, result }; } function rpcError(id: number | string | null, code: number, message: string) { return { jsonrpc: "2.0", id, error: { code, message } }; } // ── Session Management ────────────────────────────────────────────────── const sessions = new Map(); const rateLimits = new Map(); let activeToolCalls = 0; function allowRequest(client: string): boolean { const now = Date.now(); const entry = rateLimits.get(client); if (!entry || now - entry.startedAt >= 60_000) { rateLimits.set(client, { startedAt: now, count: 1 }); return true; } entry.count += 1; return entry.count <= RATE_LIMIT_PER_MINUTE; } function auditTool(name: string, ok: boolean, durationMs: number): void { const line = JSON.stringify({ time: new Date().toISOString(), event: "tool_call", tool: name, risk: getToolRisk(name), ok, duration_ms: durationMs, }) + "\n"; appendFile(AUDIT_LOG, line, { mode: 0o600 }, () => {}); } function newSessionId(): string { return crypto.randomUUID(); } function touchSession(id: string) { const now = Date.now(); if (sessions.has(id)) { sessions.get(id)!.lastActivity = now; } else { sessions.set(id, { createdAt: now, lastActivity: now }); } } function deleteSession(id: string | null) { if (id) sessions.delete(id); } // ── MCP Request Handler ───────────────────────────────────────────────── async function handleMcpRequest( message: Record, sessionId: string | null, ): Promise<{ response: unknown; sessionId: string } | null> { const method = (message["method"] as string) ?? ""; const id = message["id"] as number | string | null; const params = (message["params"] as Record) ?? {}; // ── initialize ──────────────────────────────────────────────────────── if (method === "initialize") { const sid = newSessionId(); touchSession(sid); return { response: rpcResult(id, { protocolVersion: MUA_PROTOCOL_VERSION, capabilities: { tools: {} }, serverInfo: { name: MUA_SERVER_NAME, version: MUA_VERSION }, }), sessionId: sid, }; } // ── notifications/initialized (kein Response) ───────────────────────── if (method === "notifications/initialized") { if (sessionId) touchSession(sessionId); return null; } // ── tools/list ──────────────────────────────────────────────────────── if (method === "tools/list") { if (sessionId) touchSession(sessionId); // Tool-Filter: nur aktive Tools anzeigen const activeTools = TOOLS.filter((t) => isToolEnabled(t.name)); return { response: rpcResult(id, { tools: activeTools.map((t) => ({ name: t.name, description: t.description, inputSchema: t.inputSchema, })), }), sessionId: sessionId ?? "", }; } // ── tools/call ──────────────────────────────────────────────────────── if (method === "tools/call") { if (sessionId) touchSession(sessionId); const toolName = (params["name"] as string) ?? ""; const args = (params["arguments"] as Record) ?? {}; const tool = toolByName(toolName); if (!tool) { return { response: rpcResult(id, { content: [{ type: "text", text: `ERROR: Unknown tool: ${toolName}` }], isError: true, }), sessionId: sessionId ?? "", }; } // Tool-Filter: deaktivierte Tools ablehnen if (!isToolEnabled(toolName)) { return { response: rpcResult(id, { content: [{ type: "text", text: `ERROR: Tool disabled: ${toolName}` }], isError: true, }), sessionId: sessionId ?? "", }; } if (activeToolCalls >= MAX_CONCURRENT_TOOL_CALLS) { return { response: rpcResult(id, { content: [{ type: "text", text: "ERROR: Server busy; retry later" }], isError: true, }), sessionId: sessionId ?? "", }; } activeToolCalls += 1; const startedAt = Date.now(); try { const text = await tool.handler(args); auditTool(toolName, true, Date.now() - startedAt); const result: Record = { content: [{ type: "text", text }], }; // JSON-Ausgaben zusätzlich strukturiert bereitstellen. Reiner Text // bleibt unverändert im normalen MCP-Content-Feld. try { const structured = JSON.parse(text); if (structured !== null && typeof structured === "object") { result.structuredContent = structured; } } catch { // Kein JSON-Tool: Textausgabe genügt. } return { response: rpcResult(id, result), sessionId: sessionId ?? "" }; } catch (e) { auditTool(toolName, false, Date.now() - startedAt); return { response: rpcResult(id, { content: [{ type: "text", text: `ERROR: ${String(e)}` }], isError: true, }), sessionId: sessionId ?? "", }; } finally { activeToolCalls -= 1; } } // ── ping ────────────────────────────────────────────────────────────── if (method === "ping") { if (sessionId) touchSession(sessionId); return { response: rpcResult(id, {}), sessionId: sessionId ?? "" }; } // ── Unknown method ──────────────────────────────────────────────────── if (id !== null && id !== undefined) { return { response: rpcError(id, -32601, `Method not found: ${method}`), sessionId: sessionId ?? "", }; } return null; } // ── HTTP Server (Bun.serve — production-grade) ────────────────────────── const server = Bun.serve({ port: PORT, hostname: HOST, idleTimeout: 120, // Sekunden (für lange docker stats / audit) async fetch(req) { const url = new URL(req.url); const path = url.pathname; const method = req.method; const sessionHeader = req.headers.get("mcp-session-id"); // Browserzugriffe sind standardmäßig deaktiviert. Bei Bedarf kann genau // ein vertrauenswürdiger Origin per MUA_CORS_ORIGIN freigegeben werden. const origin = req.headers.get("origin") ?? ""; const corsHeaders: Record = CORS_ORIGIN && origin === CORS_ORIGIN ? { "Access-Control-Allow-Origin": CORS_ORIGIN, "Access-Control-Allow-Methods": "GET, POST, DELETE, OPTIONS", "Access-Control-Allow-Headers": "Authorization, Content-Type, Mcp-Session-Id", "Vary": "Origin", } : {}; // ── OPTIONS (CORS Preflight) ──────────────────────────────────────── if (method === "OPTIONS") { if (!CORS_ORIGIN || origin !== CORS_ORIGIN) { return Response.json({ error: "CORS origin not allowed" }, { status: 403 }); } return new Response(null, { status: 204, headers: corsHeaders }); } // ── Health-Check (ohne Auth) ──────────────────────────────────────── if (path === "/health") { return Response.json( { status: "ok", server: MUA_SERVER_NAME, version: MUA_VERSION, port: PORT, auth: "required", time: new Date().toISOString(), }, { headers: corsHeaders }, ); } // ── Auth-Check für /mcp (POST + DELETE) — Bearer-Token ───────────── if (path === "/mcp" || path === "/") { const client = server.requestIP(req)?.address ?? "unknown"; if (!allowRequest(client)) { return Response.json( rpcError(null, -32002, "Rate limit exceeded"), { status: 429, headers: { ...corsHeaders, "Retry-After": "60" } }, ); } if (!checkAuth(req)) { return Response.json( { jsonrpc: "2.0", id: null, error: { code: -32001, message: "Unauthorized: missing or invalid API key" }, }, { status: 401, headers: { ...corsHeaders, "Mcp-Session-Id": "00000000-0000-0000-0000-000000000000", }, }, ); } } // ── MCP-Endpunkt ──────────────────────────────────────────────────── if (path === "/mcp" || path === "/") { // POST: JSON-RPC Request if (method === "POST") { const contentLength = Number(req.headers.get("content-length") ?? 0); if (contentLength > MAX_REQUEST_BYTES) { return Response.json(rpcError(null, -32600, "Request too large"), { status: 413, headers: corsHeaders, }); } let body: string; try { body = await req.text(); } catch (e) { return Response.json( rpcError(null, -32700, "Parse error"), { status: 400, headers: corsHeaders }, ); } if (new TextEncoder().encode(body).byteLength > MAX_REQUEST_BYTES) { return Response.json(rpcError(null, -32600, "Request too large"), { status: 413, headers: corsHeaders, }); } let message: Record; try { message = JSON.parse(body); } catch { return Response.json( rpcError(null, -32700, "Parse error"), { status: 400, headers: corsHeaders }, ); } const result = await handleMcpRequest(message, sessionHeader); // Notification (kein Response nötig) if (result === null) { return new Response(null, { status: 202, headers: { ...corsHeaders, "Mcp-Session-Id": sessionHeader ?? "", }, }); } return Response.json(result.response, { headers: { ...corsHeaders, "Mcp-Session-Id": result.sessionId, }, }); } // GET: SSE-Stream (nicht implementiert, spec-konform 405) if (method === "GET") { return Response.json( { error: "SSE not supported. Use POST /mcp for JSON-RPC requests." }, { status: 405, headers: corsHeaders }, ); } // DELETE: Session-End if (method === "DELETE") { deleteSession(sessionHeader); return new Response(null, { status: 200, headers: corsHeaders }); } } // ── 404 ───────────────────────────────────────────────────────────── return Response.json( { error: "Not found. Use /mcp or /health" }, { status: 404, headers: corsHeaders }, ); }, }); // ── Config-Server (nur localhost:127.0.0.1:3013) ─────────────────────── // Separater Server, nur von localhost erreichbar (WebGUI auf demselben Host). // Kein Header-Spoofing, kein File-Permission-Problem. // Wenn der Port belegt ist: Warnung + weiterlaufen (MCP-Server bleibt up). const CONFIG_PORT = Number(process.env["MUA_CONFIG_PORT"] ?? 3013); const ADMIN_TOKEN_FILE = process.env["MUA_ADMIN_TOKEN_FILE"] ?? "/var/run/mua-admin.token"; const adminToken = randomBytes(32).toString("hex"); try { writeFileSync(ADMIN_TOKEN_FILE, adminToken + "\n", { mode: 0o600 }); chmodSync(ADMIN_TOKEN_FILE, 0o600); } catch (e) { console.error(`[MUA] Config-Token konnte nicht geschrieben werden: ${String(e)}`); } let configServer: ReturnType | null = null; try { configServer = Bun.serve({ port: CONFIG_PORT, hostname: "127.0.0.1", // nur localhost idleTimeout: 30, async fetch(req) { const url = new URL(req.url); const path = url.pathname; const method = req.method; if (path !== "/config" && path !== "/restart") { return Response.json({ error: "Not found" }, { status: 404 }); } if (req.headers.get("x-mua-admin-token") !== adminToken) { return Response.json({ error: "Unauthorized" }, { status: 401 }); } // POST /restart: Service neu starten (für WebGUI-Button) if (path === "/restart" && method === "POST") { const rcScript = process.env["MUA_RC_SCRIPT"] ?? "/etc/rc.d/rc.mua"; try { const child = spawn(rcScript, ["restart"], { detached: true, stdio: "ignore", }); child.unref(); return Response.json({ ok: true, message: "Restart initiated" }); } catch (e) { return Response.json({ error: String(e) }, { status: 500 }); } } // GET: Config lesen if (method === "GET") { return Response.json({ apiKeyConfigured: true, apiKeyMasked: getMaskedApiKey(), enabledTools: getEnabledTools(), allToolsEnabled: getAllToolsEnabled(), allTools: TOOLS.map((t) => ({ name: t.name, description: t.description, risk: getToolRisk(t.name), })), }); } // POST: Config schreiben if (method === "POST") { let body: Record; try { body = JSON.parse(await req.text()); } catch { return Response.json({ error: "Invalid JSON" }, { status: 400 }); } // "generate": neuen API-Key generieren let generatedApiKey: string | undefined; if (body["generate"] === true) generatedApiKey = generateApiKey(); if (Array.isArray(body["enabledTools"])) { const known = new Set(TOOLS.map((t) => t.name)); const tools = (body["enabledTools"] as unknown[]) .filter((t): t is string => typeof t === "string" && known.has(t)); setEnabledTools(tools, body["allToolsEnabled"] === true); } return Response.json({ ok: true, generatedApiKey, apiKeyMasked: getMaskedApiKey(), enabledTools: getEnabledTools(), allToolsEnabled: getAllToolsEnabled(), }); } return Response.json({ error: "Method not allowed" }, { status: 405 }); }, }); } catch (e) { console.warn( `[MUA] WARNUNG: Config-Server konnte nicht gestartet werden (Port ${CONFIG_PORT} belegt?). ` + `Einstellungen über WebGUI sind nicht verfügbar. MCP-Server läuft weiter.`, ); } console.log( `[MUA] ${MUA_SERVER_NAME} v${MUA_VERSION} listening on ${HOST}:${PORT} (MCP)` + (configServer ? ` + 127.0.0.1:${CONFIG_PORT} (config)` : " (config: NICHT VERFÜGBAR)"), ); // ── Auto-Restart bei Binary-Update ────────────────────────────────────── // Beim Plugin-Update wird das Binary auf Platte ersetzt, aber der laufende // Prozess hält das alte Binary im Speicher (Linux-Verhalten). Wir erkennen // den Change (mtime + size) und starten den Service automatisch neu. const BINARY_PATH = process.env["MUA_BINARY"] ?? "/usr/local/bin/mua"; const RC_SCRIPT = process.env["MUA_RC_SCRIPT"] ?? "/etc/rc.d/rc.mua"; let binarySnapshot: { mtimeMs: number; size: number } | null = null; try { const st = statSync(BINARY_PATH); binarySnapshot = { mtimeMs: st.mtimeMs, size: st.size }; } catch { // Binary-Pfad nicht lesbar (z.B. lokaler Test) — Auto-Restart deaktiviert } function checkBinaryUpdate() { if (!binarySnapshot) return; try { const st = statSync(BINARY_PATH); if (st.mtimeMs !== binarySnapshot.mtimeMs || st.size !== binarySnapshot.size) { console.log( `[MUA] Binary-Update erkannt (${BINARY_PATH} geändert). ` + `Starte Service neu: ${RC_SCRIPT} restart`, ); // rc.mua restart detached ausführen, dann sauber beenden const child = spawn(RC_SCRIPT, ["restart"], { detached: true, stdio: "ignore", }); child.unref(); // Kurze Pause, damit der Log-Flush durchkommt setTimeout(() => process.exit(0), 500); } } catch { // Binary temporär nicht lesbar — ignorieren } } const updateCheckInterval = setInterval( checkBinaryUpdate, Number(process.env["MUA_UPDATE_CHECK_INTERVAL"] ?? 30_000), ); updateCheckInterval.unref(); // Abgelaufene Sessions und Rate-Limit-Einträge entfernen, damit lange // Laufzeiten nicht durch beliebig viele Client-Adressen Speicher ansammeln. const housekeepingInterval = setInterval(() => { const now = Date.now(); for (const [id, session] of sessions) { if (now - session.lastActivity > 60 * 60 * 1000) sessions.delete(id); } for (const [client, entry] of rateLimits) { if (now - entry.startedAt > 2 * 60 * 1000) rateLimits.delete(client); } }, 5 * 60 * 1000); housekeepingInterval.unref(); // Graceful shutdown process.on("SIGTERM", () => { console.log("[MUA] SIGTERM received, shutting down"); server.stop(true); process.exit(0); }); process.on("SIGINT", () => { console.log("[MUA] SIGINT received, shutting down"); server.stop(true); process.exit(0); });