From a3fed47e08d2aef0aff96b0acd05cd0adbddc279 Mon Sep 17 00:00:00 2001 From: hyperpolymath <6759885+hyperpolymath@users.noreply.github.com> Date: Wed, 20 May 2026 08:57:01 +0100 Subject: [PATCH] feat(mcp-bridge): Streamable HTTP transport (epic #87 item 14, PR1) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements the bridge half of ADR-0013: HTTP+SSE transport alongside stdio. PR1 of 2 — Workers / Durable-Objects shim and mTLS/OIDC auth are owed in PR2. What's in: - BOJ_TRANSPORT=stdio|http|both mode selection in main.js - lib/http-transport.js — POST /mcp + GET /mcp (SSE) + DELETE /mcp + GET /healthz; per-session UUID issuance via Mcp-Session-Id; 30-min idle expiry; bearer / none auth; loopback-refuse on auth=none + non-loopback bind; zero new deps (Deno.serve / node:http) - lib/dispatcher.js — transport-neutral dispatch core consumed by both stdio (main.js) and HTTP, same hardeningGate, same tool surface - boj://capabilities/deployment resource flagging the 5 host-local-only cartridges (browser-mcp, container-mcp, local-coord-mcp, sandbox-mcp, ffmpeg-mcp) so clients can detect Worker-incompatible capabilities - glama.json declares BOJ_TRANSPORT, BOJ_HTTP_PORT (7780), BOJ_HTTP_BIND, BOJ_HTTP_AUTH, BOJ_HTTP_AUTH_TOKENS - tests/http_transport_test.js — 26 tests covering session manager, auth, parse / size / unknown-session paths, end-to-end POST flow, bearer reject/accept, DELETE teardown, healthz, deployment resource Test plan: - node --test mcp-bridge/tests/dispatch_test.js mcp-bridge/tests/http_transport_test.js → 41/41 pass (was 15/15 baseline; +26 new) - deno test --allow-net --allow-env --allow-read mcp-bridge/tests/ → 41/41 pass - stdio smoke (Deno + Node): unchanged behavior - http smoke (Deno + Node): Mcp-Session-Id issued, dispatch parity - bearer reject (no token / bad token) → 401; valid token → 200 - BOJ_HTTP_AUTH=none + BOJ_HTTP_BIND=0.0.0.0 refuses to start, exit 1 What's owed in PR2 (tracked): - mTLS + OIDC auth modes - Cloudflare Workers deployment guide + wrangler.toml example - Durable-Object wrapper around SessionManager - Static cartridge-manifest bundling for cold-start - Per-cartridge `requires_local` flag wired from cartridge.json (deployment resource list is static for PR1) Refs #87 (item 14), ADR-0013. Co-Authored-By: Claude Opus 4.7 (1M context) --- CHANGELOG.md | 15 + glama.json | 24 + mcp-bridge/lib/dispatcher.js | 358 ++++++++++++++ mcp-bridge/lib/http-transport.js | 570 ++++++++++++++++++++++ mcp-bridge/lib/resources.js | 48 ++ mcp-bridge/main.js | 608 +++--------------------- mcp-bridge/tests/http_transport_test.js | 333 +++++++++++++ 7 files changed, 1426 insertions(+), 530 deletions(-) create mode 100644 mcp-bridge/lib/dispatcher.js create mode 100644 mcp-bridge/lib/http-transport.js create mode 100644 mcp-bridge/tests/http_transport_test.js diff --git a/CHANGELOG.md b/CHANGELOG.md index 5790a6c1..53994259 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,21 @@ All notable changes to Bundle of Joy Server are documented here. ### Added +- **Streamable HTTP transport (ADR-0013, PR1 of 2)** — MCP bridge now selects + between stdio (default), `http`, and `both` via `BOJ_TRANSPORT`. HTTP + endpoints: `POST /mcp` for JSON-RPC, `GET /mcp` for the server-initiated + SSE notifications stream, `DELETE /mcp` for explicit session teardown, + `GET /healthz` for liveness. Sessions are server-issued UUIDs in the + `Mcp-Session-Id` header; the manager expires idle sessions after 30 min + and fans events out across attached SSE streams. Auth: `none` (loopback + only — refuses non-loopback binds) or `bearer` (token list via + `BOJ_HTTP_AUTH_TOKENS`). The same `hardeningGate` runs on every request. + Zero new deps — built on `Deno.serve` and `node:http`. mTLS / OIDC auth + and the Cloudflare Workers / Durable-Objects shim are owed in PR2. +- **`boj://capabilities/deployment` resource** — reports per-deployment + cartridge availability so clients can avoid invoking host-local-only + cartridges (browser-mcp, container-mcp, local-coord-mcp, sandbox-mcp, + ffmpeg-mcp) against a Worker / remote-HTTP deployment. - **k9iser-mcp cartridge** — reference implementation of the `-iser` regeneration-cartridge pattern (central K9 contract regeneration), mirroring ssg-mcp: `cartridge.json`, `mod.js`, Idris2 ABI, Zig FFI, panels. diff --git a/glama.json b/glama.json index 717af72b..0ea98c82 100644 --- a/glama.json +++ b/glama.json @@ -76,6 +76,30 @@ "OTEL_EXPORTER_OTLP_HEADERS": { "type": "string", "description": "Optional comma-separated key=value headers attached to OTLP export POSTs (e.g. 'authorization=Bearer xyz,x-honeycomb-team=abc'). Used for hosted collectors that require auth." + }, + "BOJ_TRANSPORT": { + "type": "string", + "description": "MCP transport selection (ADR-0013). 'stdio' (default) reads JSON-RPC from stdin and writes to stdout — the protocol clients like Claude Code launch the bridge as a subprocess over. 'http' starts an HTTP+SSE listener on BOJ_HTTP_PORT for remote / Workers / browser deployments. 'both' runs both simultaneously.", + "default": "stdio" + }, + "BOJ_HTTP_PORT": { + "type": "string", + "description": "TCP port for the HTTP transport (used only when BOJ_TRANSPORT=http or BOJ_TRANSPORT=both).", + "default": "7780" + }, + "BOJ_HTTP_BIND": { + "type": "string", + "description": "Bind address for the HTTP transport. Defaults to 127.0.0.1 (loopback only). Set to 0.0.0.0 for remote access — BOJ_HTTP_AUTH=none is refused on non-loopback binds.", + "default": "127.0.0.1" + }, + "BOJ_HTTP_AUTH": { + "type": "string", + "description": "Authentication mode for the HTTP transport. 'none' is permitted only on loopback (refused otherwise). 'bearer' requires Authorization: Bearer against the BOJ_HTTP_AUTH_TOKENS list. 'mtls' and 'oidc' are owed in PR2 of ADR-0013 — not yet implemented.", + "default": "none" + }, + "BOJ_HTTP_AUTH_TOKENS": { + "type": "string", + "description": "Comma-separated list of accepted bearer tokens when BOJ_HTTP_AUTH=bearer. Whitespace around each token is trimmed; empty entries dropped. Required when BOJ_HTTP_AUTH=bearer." } } } diff --git a/mcp-bridge/lib/dispatcher.js b/mcp-bridge/lib/dispatcher.js new file mode 100644 index 00000000..bd452cdf --- /dev/null +++ b/mcp-bridge/lib/dispatcher.js @@ -0,0 +1,358 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +// +// BoJ Server — MCP JSON-RPC dispatch (transport-agnostic) +// +// Shared dispatch core consumed by both stdio (main.js) and HTTP +// (http-transport.js) transports. Returns a JSON-RPC response object +// or null for notifications. Does not touch stdout / sockets. + +import { env } from "./runtime.js"; +import { + RATE_LIMIT, + isInputSizeOk, + isValidToolName, + rateLimitAllow, + sanitizeErrorMessage, + scanObjectForInjection, + validateRequiredStrings, +} from "./security.js"; +import { + fetchCartridgeInfo, + fetchCartridges, + fetchHealth, + fetchMenu, + handleGitHubTool, + handleGitLabTool, + invokeCartridge, +} from "./api-clients.js"; +import { buildToolList } from "./tools.js"; +import { listResources, readResource } from "./resources.js"; +import { listPrompts, getPrompt } from "./prompts.js"; +import { + initValidator, + tryParseEnvelope, + validateEnvelope, +} from "./nickel-validator.js"; +import { info, warn, error as logError, setLevel as setLogLevel } from "./logger.js"; +import * as otel from "./otel.js"; + +const SERVER_NAME = "boj-server"; +const SERVER_VERSION = "0.4.7"; + +const LOCAL_COORD_URL = env.get("COORD_BACKEND_URL") ?? "http://127.0.0.1:7745"; +const ENVELOPE_CARRYING_TOOLS = new Set(["coord_send", "coord_send_gated"]); + +initValidator(); + +function result(id, value) { + return { jsonrpc: "2.0", id, result: value }; +} + +function rpcError(id, code, message) { + return { jsonrpc: "2.0", id, error: { code, message: sanitizeErrorMessage(message) } }; +} + +function hardeningGate(toolName, args) { + if (!rateLimitAllow()) { + return { code: -32000, message: "Rate limit exceeded. Max " + RATE_LIMIT + " tool calls per minute." }; + } + if (!isValidToolName(toolName)) { + return { code: -32602, message: "Invalid tool name" }; + } + if (!isInputSizeOk(args)) { + return { code: -32600, message: "Tool arguments exceed maximum size (1 MB)" }; + } + const injectionLevel = scanObjectForInjection(args); + if (injectionLevel === "critical" || injectionLevel === "high") { + logError("Injection blocked", { tool: toolName, confidence: injectionLevel }); + return { code: -32600, message: "Request rejected: suspicious content detected" }; + } + if (injectionLevel === "medium") { + warn("Injection warning", { tool: toolName, confidence: injectionLevel }); + } + + let validationError = null; + if (toolName === "boj_cartridge_info" || toolName === "boj_cartridge_invoke") { + validationError = validateRequiredStrings(args, ["name"]); + } else if (toolName === "boj_browser_navigate") { + validationError = validateRequiredStrings(args, ["url"]); + } else if (toolName === "boj_browser_click") { + validationError = validateRequiredStrings(args, ["selector"]); + } else if (toolName === "boj_browser_type") { + validationError = validateRequiredStrings(args, ["selector", "text"]); + } else if (toolName === "boj_browser_execute_js") { + validationError = validateRequiredStrings(args, ["script"]); + } else if (toolName.startsWith("boj_github_") && toolName !== "boj_github_list_repos") { + if (toolName === "boj_github_graphql" || toolName === "boj_github_search_code" || toolName === "boj_github_search_issues") { + validationError = validateRequiredStrings(args, ["query"]); + } else { + validationError = validateRequiredStrings(args, ["owner", "repo"]); + } + } else if (toolName.startsWith("boj_gitlab_") && toolName !== "boj_gitlab_list_projects") { + validationError = validateRequiredStrings(args, ["project_id"]); + } else if (toolName.startsWith("boj_cloud_") || toolName.startsWith("boj_comms_") || toolName === "boj_ml_huggingface" || toolName === "boj_research" || toolName === "boj_codeseeker" || toolName === "boj_search") { + validationError = validateRequiredStrings(args, ["operation"]); + } else if (toolName === "boj_browser_tabs") { + validationError = validateRequiredStrings(args, ["operation"]); + } + + if (validationError) { + return { code: -32602, message: validationError }; + } + return null; +} + +async function dispatchTool(toolName, args) { + switch (toolName) { + case "boj_health": + return fetchHealth(); + case "boj_menu": + return fetchMenu(); + case "boj_cartridges": + return fetchCartridges(); + case "boj_cartridge_info": + return fetchCartridgeInfo(args.name); + case "boj_cartridge_invoke": + return invokeCartridge(args.name, args.params); + + case "boj_cloud_verpex": + case "boj_cloud_cloudflare": + case "boj_cloud_vercel": + return invokeCartridge("cloud-mcp", { provider: toolName.replace("boj_cloud_", ""), ...args }); + + case "boj_comms_gmail": + case "boj_comms_calendar": + return invokeCartridge("comms-mcp", { provider: toolName.replace("boj_comms_", ""), ...args }); + + case "boj_ml_huggingface": + return invokeCartridge("ml-mcp", { provider: "huggingface", ...args }); + + case "boj_browser_navigate": + case "boj_browser_click": + case "boj_browser_type": + case "boj_browser_read_page": + case "boj_browser_screenshot": + case "boj_browser_tabs": + case "boj_browser_execute_js": + return invokeCartridge("browser-mcp", { action: toolName.replace("boj_browser_", ""), ...args }); + + case "boj_github_list_repos": + case "boj_github_get_repo": + case "boj_github_create_issue": + case "boj_github_list_issues": + case "boj_github_get_issue": + case "boj_github_comment_issue": + case "boj_github_create_pr": + case "boj_github_list_prs": + case "boj_github_get_pr": + case "boj_github_merge_pr": + case "boj_github_search_code": + case "boj_github_search_issues": + case "boj_github_get_file": + case "boj_github_graphql": + return handleGitHubTool(toolName, args); + + case "boj_gitlab_list_projects": + case "boj_gitlab_get_project": + case "boj_gitlab_create_issue": + case "boj_gitlab_list_issues": + case "boj_gitlab_create_mr": + case "boj_gitlab_list_mrs": + case "boj_gitlab_list_pipelines": + case "boj_gitlab_setup_mirror": + return handleGitLabTool(toolName, args); + + case "boj_codeseeker": + return invokeCartridge("codeseeker-mcp", args); + + case "boj_research": + return invokeCartridge("research-mcp", args); + + case "boj_search": + return invokeCartridge("search-mcp", args); + + case "coord_register": + case "coord_list_peers": + case "coord_send": + case "coord_receive": + case "coord_claim_task": + case "coord_status": + case "coord_promote_to_master": + case "coord_promote_to_supervisor": + case "coord_send_gated": + case "coord_review": + case "coord_review_entry": + case "coord_approve": + case "coord_reject": + case "coord_report_outcome": + case "coord_get_affinities": + case "coord_set_declared_affinities": + case "coord_scan_suggestions": + case "coord_transfer_master": + case "coord_set_variant": + case "coord_set_capabilities": + case "coord_get_peer_capabilities": + case "coord_health": + case "coord_progress": + case "coord_sweep_watchdog": + return dispatchLocalCoord(toolName, args); + + default: + return null; + } +} + +async function dispatchLocalCoord(toolName, args) { + if (ENVELOPE_CARRYING_TOOLS.has(toolName) && args && typeof args.message === "string") { + const envelope = tryParseEnvelope(args.message); + if (envelope && typeof envelope === "object") { + const senderRole = envelope._meta?.sender_role || args.sender_role; + const r = validateEnvelope(envelope, senderRole); + if (!r.ok) { + return { + success: false, + error: `envelope validation failed: ${r.error}`, + hint: "See cartridges/local-coord-mcp/schemas/coord-messages-contracts.ncl for the active contracts", + }; + } + } + } + try { + const res = await fetch(`${LOCAL_COORD_URL}/tools/${toolName}`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(args || {}), + }); + try { + return await res.json(); + } catch { + return { success: false, error: "local-coord-mcp backend returned non-JSON" }; + } + } catch (e) { + return { + success: false, + error: `local-coord-mcp backend unavailable at ${LOCAL_COORD_URL}: ${e.message}`, + hint: "Start the adapter: cd cartridges/local-coord-mcp/adapter && zig build run", + }; + } +} + +/** + * Dispatch a parsed JSON-RPC message. Returns a response object or + * null for notifications. Transport-agnostic — the caller is responsible + * for delivery (stdout write, HTTP body, SSE event, etc.). + * + * @param {object} msg parsed JSON-RPC message + * @param {object} [ctx] optional context — { transport: "stdio"|"http", sessionId } + * @returns {Promise} + */ +export async function dispatchMcpMessage(msg, ctx = {}) { + const { id, method, params } = msg; + + switch (method) { + case "initialize": { + info("MCP initialize", { client: params?.clientInfo?.name, transport: ctx.transport ?? "stdio" }); + return result(id, { + protocolVersion: "2024-11-05", + capabilities: { + tools: { listChanged: true }, + resources: { subscribe: false }, + prompts: { listChanged: false }, + logging: { levels: ["debug", "info", "warn", "error"] }, + }, + serverInfo: { name: SERVER_NAME, version: SERVER_VERSION }, + }); + } + + case "notifications/initialized": + return null; + + case "logging/setLevel": { + const level = params?.level; + const applied = setLogLevel(level); + if (!applied) { + return rpcError(id, -32602, `Unknown log level: ${level}. Expected debug|info|warn|error|silent.`); + } + info("MCP logging/setLevel", { level }); + return result(id, {}); + } + + case "tools/list": + return result(id, { tools: buildToolList() }); + + case "resources/list": + return result(id, { resources: listResources() }); + + case "resources/read": { + const uri = params?.uri; + if (typeof uri !== "string" || !uri.startsWith("boj://")) { + return rpcError(id, -32602, "resources/read requires a boj:// URI"); + } + try { + const r = await readResource(uri); + if (r === null) return rpcError(id, -32602, `Unknown resource: ${uri}`); + return result(id, r); + } catch (e) { + return rpcError(id, -32603, e?.message ?? String(e)); + } + } + + case "prompts/list": + return result(id, { prompts: listPrompts() }); + + case "prompts/get": { + const name = params?.name; + const args = params?.arguments ?? {}; + if (typeof name !== "string" || name.length === 0) { + return rpcError(id, -32602, "prompts/get requires a 'name'"); + } + const { result: promptResult, error: promptError } = getPrompt(name, args); + if (promptError) return rpcError(id, promptError.code, promptError.message); + return result(id, promptResult); + } + + case "tools/call": { + const toolName = params?.name; + const args = params?.arguments || {}; + const span = otel.startSpan("mcp.tools.call", { + "mcp.tool.name": toolName, + "mcp.tool.arg_count": Object.keys(args).length, + "mcp.transport": ctx.transport ?? "stdio", + }); + const rejection = hardeningGate(toolName, args); + if (rejection) { + otel.endSpan(span, { + status: "error", + error: "gate_rejected", + attributes: { "mcp.rejection.code": rejection.code }, + }); + return rpcError(id, rejection.code, rejection.message); + } + try { + const toolResult = await dispatchTool(toolName, args); + if (toolResult === null) { + otel.endSpan(span, { status: "error", error: "unknown_tool" }); + return rpcError(id, -32601, "Unknown tool"); + } + otel.endSpan(span, { status: "ok" }); + return result(id, { + content: [{ type: "text", text: JSON.stringify(toolResult, null, 2) }], + }); + } catch (e) { + otel.endSpan(span, { + status: "error", + error: sanitizeErrorMessage(e?.message ?? String(e)), + }); + return rpcError(id, -32603, e?.message ?? String(e)); + } + } + + case "ping": + return result(id, {}); + + default: + if (id !== undefined) return rpcError(id, -32601, "Method not found"); + return null; + } +} diff --git a/mcp-bridge/lib/http-transport.js b/mcp-bridge/lib/http-transport.js new file mode 100644 index 00000000..79331447 --- /dev/null +++ b/mcp-bridge/lib/http-transport.js @@ -0,0 +1,570 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +// +// BoJ Server — MCP Streamable HTTP transport (per ADR-0013, PR1 of 2) +// +// Adds an HTTP+SSE transport alongside stdio. Same `dispatchMcpMessage`, +// same `hardeningGate`, same tool surface — only the I/O layer differs. +// +// Endpoints: +// POST /mcp — submit a JSON-RPC request, get a JSON response +// GET /mcp — open an SSE stream for server-initiated notifications +// GET /healthz — liveness probe (200 OK, no auth) +// +// Headers: +// Mcp-Session-Id — server-issued on `initialize`; client sends it +// on every subsequent request +// Authorization — `Bearer ` when BOJ_HTTP_AUTH=bearer +// +// Zero deps. Uses Deno.serve / node:http via runtime detection. + +import { isDeno, env } from "./runtime.js"; +import { dispatchMcpMessage } from "./dispatcher.js"; +import { info, warn, error as logError } from "./logger.js"; + +const DEFAULT_PORT = 7780; +const DEFAULT_BIND = "127.0.0.1"; +const SESSION_TIMEOUT_MS = 30 * 60 * 1000; // 30 minutes +const SSE_KEEPALIVE_MS = 25 * 1000; +const MAX_BODY_BYTES = 2 * 1024 * 1024; // 2 MB — same as stdio buffer cap + +const LOOPBACK_HOSTS = new Set(["127.0.0.1", "localhost", "::1", "0:0:0:0:0:0:0:1"]); + +// globalThis.crypto.randomUUID is stable on Deno (all versions) and Node +// (>=18.17). package.json declares engines.node >=18.0.0 — older Node 18 +// builds without globalThis.crypto.randomUUID are out of support. +function makeSessionId() { + return globalThis.crypto.randomUUID(); +} + +/** + * Parse the BOJ_HTTP_AUTH_TOKENS CSV into a Set. Empty tokens are + * dropped. Whitespace around each token is trimmed. + */ +function parseTokens(raw) { + if (!raw) return new Set(); + return new Set( + raw.split(",") + .map((t) => t.trim()) + .filter((t) => t.length > 0), + ); +} + +function isLoopback(host) { + if (!host) return false; + return LOOPBACK_HOSTS.has(host); +} + +/** + * Read configuration from environment. Returns a normalised opts object + * used by createHttpServer. Throws on configurations that would expose + * the bridge unsafely (auth=none + non-loopback bind). + */ +export function configFromEnv() { + const port = parseInt(env.get("BOJ_HTTP_PORT") ?? `${DEFAULT_PORT}`, 10) || DEFAULT_PORT; + const bind = env.get("BOJ_HTTP_BIND") ?? DEFAULT_BIND; + const authMode = (env.get("BOJ_HTTP_AUTH") ?? (isLoopback(bind) ? "none" : "bearer")).toLowerCase(); + const tokens = parseTokens(env.get("BOJ_HTTP_AUTH_TOKENS")); + + if (authMode === "none" && !isLoopback(bind)) { + throw new Error( + `BOJ_HTTP_AUTH=none refuses to serve on non-loopback bind '${bind}'. ` + + `Set BOJ_HTTP_AUTH=bearer (with BOJ_HTTP_AUTH_TOKENS) or bind to 127.0.0.1.`, + ); + } + if (authMode === "bearer" && tokens.size === 0) { + throw new Error( + "BOJ_HTTP_AUTH=bearer requires BOJ_HTTP_AUTH_TOKENS (CSV of accepted tokens).", + ); + } + if (authMode !== "none" && authMode !== "bearer") { + throw new Error(`Unsupported BOJ_HTTP_AUTH='${authMode}'. PR1 supports 'none' and 'bearer'; mTLS/OIDC owed in PR2.`); + } + + return { port, bind, authMode, tokens }; +} + +// ===================================================================== +// Session manager +// ===================================================================== + +class SessionManager { + constructor({ timeoutMs = SESSION_TIMEOUT_MS } = {}) { + this.timeoutMs = timeoutMs; + this.sessions = new Map(); + } + + create() { + const id = makeSessionId(); + this.sessions.set(id, { + id, + createdMs: Date.now(), + lastSeenMs: Date.now(), + sseStreams: new Set(), + }); + return id; + } + + touch(id) { + const s = this.sessions.get(id); + if (!s) return false; + s.lastSeenMs = Date.now(); + return true; + } + + get(id) { + return this.sessions.get(id); + } + + delete(id) { + const s = this.sessions.get(id); + if (!s) return; + for (const stream of s.sseStreams) { + try { stream.close(); } catch { /* ignore */ } + } + this.sessions.delete(id); + } + + /** Drop sessions idle beyond timeoutMs. Returns the count removed. */ + expireIdle(now = Date.now()) { + let removed = 0; + for (const [id, s] of this.sessions) { + if (now - s.lastSeenMs > this.timeoutMs) { + this.delete(id); + removed += 1; + } + } + return removed; + } + + attachStream(id, stream) { + const s = this.sessions.get(id); + if (!s) return false; + s.sseStreams.add(stream); + return true; + } + + detachStream(id, stream) { + const s = this.sessions.get(id); + if (!s) return; + s.sseStreams.delete(stream); + } + + /** + * Fan out an event to every SSE stream for the given session. Used by + * the ADR-0011 notifications path (wiring is owed; this is the seam). + */ + emit(id, eventName, data) { + const s = this.sessions.get(id); + if (!s) return 0; + let sent = 0; + for (const stream of s.sseStreams) { + try { + stream.send(eventName, data); + sent += 1; + } catch { + s.sseStreams.delete(stream); + } + } + return sent; + } +} + +// ===================================================================== +// Authentication +// ===================================================================== + +function checkAuth({ authMode, tokens, headerValue }) { + if (authMode === "none") return { ok: true }; + if (authMode === "bearer") { + if (!headerValue || typeof headerValue !== "string") { + return { ok: false, code: 401, body: { error: "missing Authorization header" } }; + } + const m = headerValue.match(/^Bearer\s+(.+)$/i); + if (!m) return { ok: false, code: 401, body: { error: "expected 'Bearer ' Authorization header" } }; + if (!tokens.has(m[1].trim())) return { ok: false, code: 401, body: { error: "invalid bearer token" } }; + return { ok: true }; + } + return { ok: false, code: 500, body: { error: `unsupported auth mode ${authMode}` } }; +} + +// ===================================================================== +// Request handler (runtime-neutral) +// ===================================================================== + +const JSON_HEADERS = { "Content-Type": "application/json; charset=utf-8" }; + +/** + * Build the JSON-RPC response object for a POST /mcp call. Returns + * `{ status, headers, body }` where body is a JS object (caller + * stringifies). Pure-ish: only side effect is session-table mutation. + */ +async function handleMcpPost({ rawBody, headers, sessions, auth }) { + const sessionHeader = headers["mcp-session-id"]; + + const authResult = checkAuth({ + authMode: auth.authMode, + tokens: auth.tokens, + headerValue: headers["authorization"], + }); + if (!authResult.ok) { + return { status: authResult.code, headers: JSON_HEADERS, body: authResult.body }; + } + + if (rawBody.length > MAX_BODY_BYTES) { + return { + status: 413, + headers: JSON_HEADERS, + body: { jsonrpc: "2.0", error: { code: -32600, message: "Message too large" } }, + }; + } + + let msg; + try { + msg = JSON.parse(rawBody); + } catch { + return { + status: 400, + headers: JSON_HEADERS, + body: { jsonrpc: "2.0", id: null, error: { code: -32700, message: "Parse error" } }, + }; + } + + let sessionId = sessionHeader; + if (msg.method === "initialize") { + // initialize mints a new session — ignore any client-supplied id + sessionId = sessions.create(); + } else if (sessionId) { + if (!sessions.touch(sessionId)) { + return { + status: 404, + headers: JSON_HEADERS, + body: { jsonrpc: "2.0", id: msg.id ?? null, error: { code: -32001, message: "Unknown or expired Mcp-Session-Id" } }, + }; + } + } + + const response = await dispatchMcpMessage(msg, { transport: "http", sessionId }); + + // notifications/* and any other no-response method + if (response === null) { + const respHeaders = { ...JSON_HEADERS }; + if (sessionId) respHeaders["Mcp-Session-Id"] = sessionId; + return { status: 202, headers: respHeaders, body: "" }; + } + + const respHeaders = { ...JSON_HEADERS }; + if (sessionId) respHeaders["Mcp-Session-Id"] = sessionId; + return { status: 200, headers: respHeaders, body: response }; +} + +// ===================================================================== +// SSE stream wrapper (runtime-neutral handle) +// ===================================================================== + +function makeSseDenoStream(controller) { + const encoder = new TextEncoder(); + return { + send(event, data) { + const payload = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`; + controller.enqueue(encoder.encode(payload)); + }, + ping() { + controller.enqueue(encoder.encode(`: keepalive\n\n`)); + }, + close() { + try { controller.close(); } catch { /* idempotent */ } + }, + }; +} + +function makeSseNodeStream(res) { + return { + send(event, data) { + res.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`); + }, + ping() { + res.write(`: keepalive\n\n`); + }, + close() { + try { res.end(); } catch { /* idempotent */ } + }, + }; +} + +// ===================================================================== +// Server lifecycle — Deno path +// ===================================================================== + +async function startDeno({ port, bind, sessions, auth }) { + const handler = async (req) => { + const url = new URL(req.url); + const headers = {}; + for (const [k, v] of req.headers) headers[k.toLowerCase()] = v; + + if (url.pathname === "/healthz" && req.method === "GET") { + return new Response("ok\n", { status: 200, headers: { "Content-Type": "text/plain" } }); + } + + if (url.pathname !== "/mcp") { + return new Response(JSON.stringify({ error: "not found" }), { status: 404, headers: JSON_HEADERS }); + } + + if (req.method === "POST") { + const rawBody = await req.text(); + const out = await handleMcpPost({ rawBody, headers, sessions, auth }); + return new Response( + typeof out.body === "string" ? out.body : JSON.stringify(out.body), + { status: out.status, headers: out.headers }, + ); + } + + if (req.method === "GET") { + const authResult = checkAuth({ + authMode: auth.authMode, + tokens: auth.tokens, + headerValue: headers["authorization"], + }); + if (!authResult.ok) { + return new Response(JSON.stringify(authResult.body), { status: authResult.code, headers: JSON_HEADERS }); + } + const sessionId = headers["mcp-session-id"]; + if (!sessionId || !sessions.touch(sessionId)) { + return new Response(JSON.stringify({ error: "missing or invalid Mcp-Session-Id" }), { status: 400, headers: JSON_HEADERS }); + } + let stream; + const body = new ReadableStream({ + start(controller) { + stream = makeSseDenoStream(controller); + sessions.attachStream(sessionId, stream); + stream.send("ready", { sessionId }); + const keepalive = setInterval(() => { + try { stream.ping(); } catch { clearInterval(keepalive); } + }, SSE_KEEPALIVE_MS); + stream._keepalive = keepalive; + }, + cancel() { + if (stream?._keepalive) clearInterval(stream._keepalive); + sessions.detachStream(sessionId, stream); + }, + }); + return new Response(body, { + status: 200, + headers: { + "Content-Type": "text/event-stream", + "Cache-Control": "no-cache", + "Connection": "keep-alive", + "Mcp-Session-Id": sessionId, + }, + }); + } + + if (req.method === "DELETE") { + const sessionId = headers["mcp-session-id"]; + if (sessionId) sessions.delete(sessionId); + // 204 responses MUST NOT carry a body per RFC 9110; Deno's + // Response constructor enforces null-body status codes. + return new Response(null, { status: 204 }); + } + + return new Response(JSON.stringify({ error: "method not allowed" }), { status: 405, headers: JSON_HEADERS }); + }; + + const server = Deno.serve({ port, hostname: bind, onListen: () => {} }, handler); + info("MCP HTTP transport listening", { bind, port, authMode: auth.authMode }); + return { + address: { host: bind, port }, + async stop() { + try { await server.shutdown(); } catch { /* ignore */ } + }, + }; +} + +// ===================================================================== +// Server lifecycle — Node path +// ===================================================================== + +async function startNode({ port, bind, sessions, auth }) { + const { createServer } = await import("node:http"); + + const server = createServer(async (req, res) => { + const url = new URL(req.url, `http://${bind}`); + const headers = {}; + for (const [k, v] of Object.entries(req.headers)) { + headers[k.toLowerCase()] = Array.isArray(v) ? v.join(",") : v; + } + + if (url.pathname === "/healthz" && req.method === "GET") { + res.writeHead(200, { "Content-Type": "text/plain" }); + res.end("ok\n"); + return; + } + + if (url.pathname !== "/mcp") { + res.writeHead(404, JSON_HEADERS); + res.end(JSON.stringify({ error: "not found" })); + return; + } + + if (req.method === "POST") { + const chunks = []; + let total = 0; + let oversized = false; + for await (const chunk of req) { + total += chunk.length; + if (total > MAX_BODY_BYTES) { oversized = true; break; } + chunks.push(chunk); + } + if (oversized) { + res.writeHead(413, JSON_HEADERS); + res.end(JSON.stringify({ jsonrpc: "2.0", error: { code: -32600, message: "Message too large" } })); + return; + } + const rawBody = Buffer.concat(chunks).toString("utf8"); + const out = await handleMcpPost({ rawBody, headers, sessions, auth }); + res.writeHead(out.status, out.headers); + res.end(typeof out.body === "string" ? out.body : JSON.stringify(out.body)); + return; + } + + if (req.method === "GET") { + const authResult = checkAuth({ + authMode: auth.authMode, + tokens: auth.tokens, + headerValue: headers["authorization"], + }); + if (!authResult.ok) { + res.writeHead(authResult.code, JSON_HEADERS); + res.end(JSON.stringify(authResult.body)); + return; + } + const sessionId = headers["mcp-session-id"]; + if (!sessionId || !sessions.touch(sessionId)) { + res.writeHead(400, JSON_HEADERS); + res.end(JSON.stringify({ error: "missing or invalid Mcp-Session-Id" })); + return; + } + res.writeHead(200, { + "Content-Type": "text/event-stream", + "Cache-Control": "no-cache", + "Connection": "keep-alive", + "Mcp-Session-Id": sessionId, + }); + const stream = makeSseNodeStream(res); + sessions.attachStream(sessionId, stream); + stream.send("ready", { sessionId }); + const keepalive = setInterval(() => { + try { stream.ping(); } catch { clearInterval(keepalive); } + }, SSE_KEEPALIVE_MS); + req.on("close", () => { + clearInterval(keepalive); + sessions.detachStream(sessionId, stream); + }); + return; + } + + if (req.method === "DELETE") { + const sessionId = headers["mcp-session-id"]; + if (sessionId) sessions.delete(sessionId); + res.writeHead(204); + res.end(); + return; + } + + res.writeHead(405, JSON_HEADERS); + res.end(JSON.stringify({ error: "method not allowed" })); + }); + + await new Promise((resolve, reject) => { + const onError = (e) => { server.off("listening", onListen); reject(e); }; + const onListen = () => { server.off("error", onError); resolve(); }; + server.once("error", onError); + server.once("listening", onListen); + server.listen(port, bind); + }); + + info("MCP HTTP transport listening", { bind, port, authMode: auth.authMode }); + + return { + address: server.address(), + async stop() { + await new Promise((resolve) => server.close(() => resolve())); + }, + }; +} + +// ===================================================================== +// Public entry point +// ===================================================================== + +/** + * Start the HTTP transport. Returns a handle with `.stop()` for tests + * and graceful shutdown. Errors during startup propagate. + * + * @param {object} [opts] override env-derived config; useful in tests + * @param {number} [opts.port] + * @param {string} [opts.bind] + * @param {"none"|"bearer"} [opts.authMode] + * @param {Iterable} [opts.tokens] + * @param {number} [opts.sessionTimeoutMs] + * @returns {Promise<{ address: object, stop: () => Promise, sessions: SessionManager }>} + */ +export async function startHttpTransport(opts = {}) { + const envConfig = (() => { + try { return configFromEnv(); } catch (e) { + // tests pass opts directly; only fail if env was the source of config + if (opts.port === undefined && opts.bind === undefined && opts.authMode === undefined) { + throw e; + } + return null; + } + })(); + + const port = opts.port ?? envConfig?.port ?? DEFAULT_PORT; + const bind = opts.bind ?? envConfig?.bind ?? DEFAULT_BIND; + const authMode = opts.authMode ?? envConfig?.authMode ?? (isLoopback(bind) ? "none" : "bearer"); + const tokens = opts.tokens ? new Set(opts.tokens) : (envConfig?.tokens ?? new Set()); + + if (authMode === "none" && !isLoopback(bind)) { + throw new Error(`Refusing to start: BOJ_HTTP_AUTH=none on non-loopback bind '${bind}'.`); + } + if (authMode === "bearer" && tokens.size === 0) { + throw new Error("BOJ_HTTP_AUTH=bearer requires at least one token."); + } + + const sessions = new SessionManager({ timeoutMs: opts.sessionTimeoutMs ?? SESSION_TIMEOUT_MS }); + + const reaper = setInterval(() => { + const removed = sessions.expireIdle(); + if (removed > 0) info("HTTP sessions expired", { count: removed }); + }, 60 * 1000); + if (typeof reaper.unref === "function") reaper.unref(); + + const auth = { authMode, tokens }; + + let handle; + if (isDeno) { + handle = await startDeno({ port, bind, sessions, auth }); + } else { + handle = await startNode({ port, bind, sessions, auth }); + } + + return { + address: handle.address, + sessions, + async stop() { + clearInterval(reaper); + await handle.stop(); + }, + }; +} + +// Exposed for tests +export const _internals = { + SessionManager, + checkAuth, + parseTokens, + isLoopback, + handleMcpPost, + configFromEnv, +}; diff --git a/mcp-bridge/lib/resources.js b/mcp-bridge/lib/resources.js index b51a5c41..32e42903 100644 --- a/mcp-bridge/lib/resources.js +++ b/mcp-bridge/lib/resources.js @@ -13,6 +13,20 @@ import { OFFLINE_MENU } from "./offline-menu.js"; import { fetchMenu, fetchCartridges, fetchCartridgeInfo } from "./api-clients.js"; +import { env } from "./runtime.js"; + +// Cartridges that need host-local resources (subprocess, loopback bus, +// browser, codec binaries). Under transport=http on a remote deployment +// (e.g. a Cloudflare Worker) they don't function. ADR-0013 §"Cartridge +// compatibility under HTTP". Kept static for PR1; PR2 derives from +// per-cartridge cartridge.json `requires_local: true` flags. +const LOCAL_ONLY_CARTRIDGES = [ + { name: "browser-mcp", reason: "Marionette session against a host Firefox profile" }, + { name: "container-mcp", reason: "shells out to podman / kubectl on the host" }, + { name: "local-coord-mcp", reason: "loopback bus on 127.0.0.1:7745" }, + { name: "sandbox-mcp", reason: "local backend uses host process isolation (SaaS backends work over HTTP)" }, + { name: "ffmpeg-mcp", reason: "invokes the host ffmpeg binary" }, +]; const STATIC_RESOURCES = [ { @@ -45,6 +59,12 @@ const STATIC_RESOURCES = [ description: "Server name, version, runtime, supported protocol versions, advertised MCP capabilities. JSON.", mimeType: "application/json", }, + { + uri: "boj://capabilities/deployment", + name: "Deployment capabilities", + description: "Per-deployment cartridge availability. Flags cartridges that need host-local resources (browser, podman, ffmpeg, loopback bus) and so cannot serve traffic when the bridge runs HTTP-only on a remote target (e.g. Cloudflare Worker). MCP clients should consult this before invoking flagged tools. JSON. (ADR-0013)", + mimeType: "application/json", + }, { uri: "boj://docs/architecture", name: "Architecture overview", @@ -223,6 +243,34 @@ async function readResource(uri) { }; } + if (uri === "boj://capabilities/deployment") { + const transport = (env.get("BOJ_TRANSPORT") ?? "stdio").toLowerCase(); + const bind = env.get("BOJ_HTTP_BIND") ?? "127.0.0.1"; + const remoteHttp = transport !== "stdio" && bind !== "127.0.0.1" && bind !== "localhost" && bind !== "::1"; + const payload = { + transport, + bind, + remote_http: remoteHttp, + // PR1 list is static; PR2 will read `requires_local` from per-cartridge manifests. + local_only_cartridges: LOCAL_ONLY_CARTRIDGES.map((c) => ({ + name: c.name, + requires_local: true, + available: !remoteHttp, + reason: c.reason, + })), + notes: [ + "Cartridges flagged requires_local=true need host-local resources and do not function on remote HTTP deployments.", + "Cartridges absent from the list are HTTP-API-based and work on any deployment.", + "PR2 (epic #87 item 14 follow-up) wires per-cartridge `requires_local` flags from cartridge.json.", + ], + }; + return { + contents: [ + { uri, mimeType: "application/json", text: JSON.stringify(payload, null, 2) }, + ], + }; + } + if (uri === "boj://server/info") { return { contents: [ diff --git a/mcp-bridge/main.js b/mcp-bridge/main.js index 75bbf278..1f8008bb 100755 --- a/mcp-bridge/main.js +++ b/mcp-bridge/main.js @@ -2,329 +2,48 @@ // SPDX-License-Identifier: MPL-2.0 // Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) // -// BoJ Server — MCP stdio transport bridge +// BoJ Server — MCP transport bridge (stdio + HTTP per ADR-0013) // // Bridges the running BoJ REST API (port 7700) to the MCP JSON-RPC -// stdio protocol so that Claude Code, Glama, and other MCP clients -// can discover and invoke BoJ cartridge tools. +// protocol so that Claude Code, Glama, and other MCP clients can +// discover and invoke BoJ cartridge tools. Two transports: // -// Usage: deno run --allow-net --allow-env --allow-read main.js +// BOJ_TRANSPORT=stdio default — read JSON-RPC from stdin, write to stdout +// BOJ_TRANSPORT=http listen on BOJ_HTTP_PORT (default 7780) +// BOJ_TRANSPORT=both run both simultaneously +// +// Usage: +// deno run --allow-net --allow-env --allow-read mcp-bridge/main.js import { env, stdout } from "./lib/runtime.js"; -import { - RATE_LIMIT, - isInputSizeOk, - isValidToolName, - rateLimitAllow, - sanitizeErrorMessage, - scanObjectForInjection, - validateRequiredStrings, -} from "./lib/security.js"; -import { - fetchCartridgeInfo, - fetchCartridges, - fetchHealth, - fetchMenu, - handleGitHubTool, - handleGitLabTool, - invokeCartridge, -} from "./lib/api-clients.js"; -import { buildToolList } from "./lib/tools.js"; -import { listResources, readResource } from "./lib/resources.js"; -import { listPrompts, getPrompt } from "./lib/prompts.js"; -import { - initValidator, - tryParseEnvelope, - validateEnvelope, -} from "./lib/nickel-validator.js"; -import { info, warn, error as logError, setLevel as setLogLevel } from "./lib/logger.js"; +import { sanitizeErrorMessage } from "./lib/security.js"; +import { dispatchMcpMessage } from "./lib/dispatcher.js"; +import { info, error as logError } from "./lib/logger.js"; import * as otel from "./lib/otel.js"; +import { startHttpTransport } from "./lib/http-transport.js"; -const BOJ_BASE = env.get("BOJ_URL") ?? "http://localhost:7700"; -const SERVER_NAME = "boj-server"; -const SERVER_VERSION = "0.4.0"; +const TRANSPORT = (env.get("BOJ_TRANSPORT") ?? "stdio").toLowerCase(); -// Initialise OTel batch-flush + shutdown hooks (no-op unless -// OTEL_EXPORTER_OTLP_ENDPOINT is set). otel.init(); // =================================================================== -// JSON-RPC stdio transport +// stdio transport // =================================================================== const decoder = new TextDecoder(); - let buffer = ""; const MAX_BUFFER_BYTES = 2 * 1_048_576; // 2 MB - const pendingMessages = []; function send(obj) { stdout.writeSync(JSON.stringify(obj) + "\n"); } -function sendResult(id, result) { - send({ jsonrpc: "2.0", id, result }); -} - function sendError(id, code, message) { send({ jsonrpc: "2.0", id, error: { code, message: sanitizeErrorMessage(message) } }); } -// =================================================================== -// Hardening gate — validates every tool call before dispatch -// =================================================================== - -/** - * Run all security checks on a tool call. - * Returns an error object {code, message} if rejected, or null if OK. - * @param {string} toolName - * @param {Record} args - * @returns {{code: number, message: string}|null} - */ -function hardeningGate(toolName, args) { - // 1. Rate limiting - if (!rateLimitAllow()) { - return { code: -32000, message: "Rate limit exceeded. Max " + RATE_LIMIT + " tool calls per minute." }; - } - - // 2. Tool name validation - if (!isValidToolName(toolName)) { - return { code: -32602, message: "Invalid tool name" }; - } - - // 3. Input size check - if (!isInputSizeOk(args)) { - return { code: -32600, message: "Tool arguments exceed maximum size (1 MB)" }; - } - - // 4. Prompt injection detection - const injectionLevel = scanObjectForInjection(args); - if (injectionLevel === "critical" || injectionLevel === "high") { - logError("Injection blocked", { tool: toolName, confidence: injectionLevel }); - return { code: -32600, message: "Request rejected: suspicious content detected" }; - } - if (injectionLevel === "medium") { - warn("Injection warning", { tool: toolName, confidence: injectionLevel }); - } - - // 5. Required field validation - let validationError = null; - if (toolName === "boj_cartridge_info" || toolName === "boj_cartridge_invoke") { - validationError = validateRequiredStrings(args, ["name"]); - } else if (toolName === "boj_browser_navigate") { - validationError = validateRequiredStrings(args, ["url"]); - } else if (toolName === "boj_browser_click") { - validationError = validateRequiredStrings(args, ["selector"]); - } else if (toolName === "boj_browser_type") { - validationError = validateRequiredStrings(args, ["selector", "text"]); - } else if (toolName === "boj_browser_execute_js") { - validationError = validateRequiredStrings(args, ["script"]); - } else if (toolName.startsWith("boj_github_") && toolName !== "boj_github_list_repos") { - if (toolName === "boj_github_graphql" || toolName === "boj_github_search_code" || toolName === "boj_github_search_issues") { - validationError = validateRequiredStrings(args, ["query"]); - } else { - validationError = validateRequiredStrings(args, ["owner", "repo"]); - } - } else if (toolName.startsWith("boj_gitlab_") && toolName !== "boj_gitlab_list_projects") { - validationError = validateRequiredStrings(args, ["project_id"]); - } else if (toolName.startsWith("boj_cloud_") || toolName.startsWith("boj_comms_") || toolName === "boj_ml_huggingface" || toolName === "boj_research" || toolName === "boj_codeseeker" || toolName === "boj_search" || toolName === "boj_vector" || toolName === "boj_multimodal") { - validationError = validateRequiredStrings(args, ["operation"]); - } else if (toolName === "boj_browser_tabs") { - validationError = validateRequiredStrings(args, ["operation"]); - } - - if (validationError) { - return { code: -32602, message: validationError }; - } - - return null; -} - -// =================================================================== -// Tool dispatch -// =================================================================== - -/** - * Dispatch a validated tool call to the appropriate handler. - * @param {string} toolName - * @param {Record} args - * @returns {Promise} - */ -async function dispatchTool(toolName, args) { - switch (toolName) { - case "boj_health": - return fetchHealth(); - case "boj_menu": - return fetchMenu(); - case "boj_cartridges": - return fetchCartridges(); - case "boj_cartridge_info": - return fetchCartridgeInfo(args.name); - case "boj_cartridge_invoke": - return invokeCartridge(args.name, args.params); - - case "boj_cloud_verpex": - case "boj_cloud_cloudflare": - case "boj_cloud_vercel": - return invokeCartridge("cloud-mcp", { provider: toolName.replace("boj_cloud_", ""), ...args }); - - case "boj_comms_gmail": - case "boj_comms_calendar": - return invokeCartridge("comms-mcp", { provider: toolName.replace("boj_comms_", ""), ...args }); - - case "boj_ml_huggingface": - return invokeCartridge("ml-mcp", { provider: "huggingface", ...args }); - - case "boj_browser_navigate": - case "boj_browser_click": - case "boj_browser_type": - case "boj_browser_read_page": - case "boj_browser_screenshot": - case "boj_browser_tabs": - case "boj_browser_execute_js": - return invokeCartridge("browser-mcp", { action: toolName.replace("boj_browser_", ""), ...args }); - - case "boj_github_list_repos": - case "boj_github_get_repo": - case "boj_github_create_issue": - case "boj_github_list_issues": - case "boj_github_get_issue": - case "boj_github_comment_issue": - case "boj_github_create_pr": - case "boj_github_list_prs": - case "boj_github_get_pr": - case "boj_github_merge_pr": - case "boj_github_search_code": - case "boj_github_search_issues": - case "boj_github_get_file": - case "boj_github_graphql": - return handleGitHubTool(toolName, args); - - case "boj_gitlab_list_projects": - case "boj_gitlab_get_project": - case "boj_gitlab_create_issue": - case "boj_gitlab_list_issues": - case "boj_gitlab_create_mr": - case "boj_gitlab_list_mrs": - case "boj_gitlab_list_pipelines": - case "boj_gitlab_setup_mirror": - return handleGitLabTool(toolName, args); - - case "boj_codeseeker": - return invokeCartridge("codeseeker-mcp", args); - - case "boj_research": - return invokeCartridge("research-mcp", args); - - case "boj_search": - return invokeCartridge("search-mcp", args); - - case "boj_vector": { - const providerToCartridge = { pinecone: "pinecone-mcp", weaviate: "weaviate-mcp", qdrant: "qdrant-mcp", chromadb: "chromadb-mcp" }; - const cart = providerToCartridge[args.provider]; - if (!cart) return { error: "unknown provider", hint: "boj_vector requires provider: pinecone | weaviate | qdrant | chromadb" }; - return invokeCartridge(cart, args); - } - - case "boj_multimodal": { - const providerToCartridge = { whisper: "whisper-mcp", elevenlabs: "elevenlabs-mcp", replicate: "replicate-mcp", ffmpeg: "ffmpeg-mcp" }; - const cart = providerToCartridge[args.provider]; - if (!cart) return { error: "unknown provider", hint: "boj_multimodal requires provider: whisper | elevenlabs | replicate | ffmpeg" }; - return invokeCartridge(cart, args); - } - - // Local coordination — direct to loopback backend on port 7745 - case "coord_register": - case "coord_list_peers": - case "coord_send": - case "coord_receive": - case "coord_claim_task": - case "coord_status": - case "coord_promote_to_master": - case "coord_promote_to_supervisor": // legacy alias — DD-32 rename; accepted for one release - case "coord_send_gated": - case "coord_review": - case "coord_review_entry": - case "coord_approve": - case "coord_reject": - case "coord_report_outcome": - case "coord_get_affinities": - case "coord_set_declared_affinities": - case "coord_scan_suggestions": - case "coord_transfer_master": - case "coord_set_variant": - case "coord_set_capabilities": - case "coord_get_peer_capabilities": - case "coord_health": - case "coord_progress": - case "coord_sweep_watchdog": - return dispatchLocalCoord(toolName, args); - - default: - return null; // unknown tool - } -} - -// =================================================================== -// local-coord-mcp direct dispatch (loopback only, port 7745) -// =================================================================== - -const LOCAL_COORD_URL = env.get("COORD_BACKEND_URL") ?? "http://127.0.0.1:7745"; - -// Nickel contracts run on coord_send / coord_send_gated only — those -// are the two tools whose `message` argument carries an A2ML envelope. -// Other coord_* calls are RPC-shaped (register/list/review/approve/...) -// and bypass contract validation. Expansion to more tools is a roadmap -// item — see Appendix K of COORD-MCP-DESIGN-LOG.md (Task #17 extension). -const ENVELOPE_CARRYING_TOOLS = new Set(["coord_send", "coord_send_gated"]); - -initValidator(); - -async function dispatchLocalCoord(toolName, args) { - // Runtime envelope validation (Task #17) — BEFORE the HTTP forward. - if (ENVELOPE_CARRYING_TOOLS.has(toolName) && args && typeof args.message === "string") { - const env = tryParseEnvelope(args.message); - if (env && typeof env === "object") { - const senderRole = env._meta?.sender_role || args.sender_role; - const result = validateEnvelope(env, senderRole); - if (!result.ok) { - return { - success: false, - error: `envelope validation failed: ${result.error}`, - hint: "See cartridges/local-coord-mcp/schemas/coord-messages-contracts.ncl for the active contracts", - }; - } - } - // Plain-string messages (non-JSON) skip validation — they're not - // A2ML envelopes. The Zig adapter still enforces shape + gating. - } - - try { - const res = await fetch(`${LOCAL_COORD_URL}/tools/${toolName}`, { - method: "POST", - headers: { "Content-Type": "application/json" }, - body: JSON.stringify(args || {}), - }); - try { - return await res.json(); - } catch { - return { success: false, error: "local-coord-mcp backend returned non-JSON" }; - } - } catch (e) { - return { - success: false, - error: `local-coord-mcp backend unavailable at ${LOCAL_COORD_URL}: ${e.message}`, - hint: "Start the adapter: cd cartridges/local-coord-mcp/adapter && zig build run", - }; - } -} - -// =================================================================== -// MCP message handler -// =================================================================== - -async function handleMessage(line) { +async function handleStdioLine(line) { let msg; try { msg = JSON.parse(line); @@ -332,232 +51,10 @@ async function handleMessage(line) { sendError(null, -32700, "Parse error"); return; } - - const { id, method, params } = msg; - - switch (method) { - case "initialize": { - info("MCP initialize", { client: params?.clientInfo?.name }); - sendResult(id, { - protocolVersion: "2024-11-05", - capabilities: { - tools: { listChanged: true }, - resources: { subscribe: false }, - prompts: { listChanged: false }, - logging: { - levels: ["debug", "info", "warn", "error"] - }, - }, - serverInfo: { name: SERVER_NAME, version: SERVER_VERSION }, - }); - break; - } - - case "notifications/initialized": - break; - - case "logging/setLevel": { - const level = params?.level; - const applied = setLogLevel(level); - if (!applied) { - sendError(id, -32602, `Unknown log level: ${level}. Expected debug|info|warn|error|silent.`); - } else { - info("MCP logging/setLevel", { level }); - sendResult(id, {}); - } - break; - } - - case "tools/list": { - // buildToolList() returns whole tool objects (name, description, - // inputSchema, annotations, outputSchema) and reads BOJ_TOOL_SCOPE - // internally for the Teranga scoped-surface lever. We pass the - // objects through verbatim so annotations + outputSchema reach the - // wire unmodified — no field whitelisting here by design. - const tools = buildToolList(); - sendResult(id, { tools }); - break; - } - - case "resources/list": { - sendResult(id, { resources: listResources() }); - break; - } - - case "resources/read": { - const uri = params?.uri; - if (typeof uri !== "string" || !uri.startsWith("boj://")) { - sendError(id, -32602, "resources/read requires a boj:// URI"); - break; - } - try { - const result = await readResource(uri); - if (result === null) { - sendError(id, -32602, `Unknown resource: ${uri}`); - } else { - sendResult(id, result); - } - } catch (e) { - sendError(id, -32603, sanitizeErrorMessage(e?.message ?? String(e))); - } - break; - } - - case "prompts/list": { - sendResult(id, { prompts: listPrompts() }); - break; - } - - case "prompts/get": { - const name = params?.name; - const args = params?.arguments ?? {}; - if (typeof name !== "string" || name.length === 0) { - sendError(id, -32602, "prompts/get requires a 'name'"); - break; - } - const { result, error: promptError } = getPrompt(name, args); - if (promptError) { - sendError(id, promptError.code, promptError.message); - } else { - sendResult(id, result); - } - break; - } - - case "tools/call": { - const toolName = params?.name; - const args = params?.arguments || {}; - const token = params?.token; - - const span = otel.startSpan("mcp.tools.call", { - "mcp.tool.name": toolName, - "mcp.tool.arg_count": Object.keys(args).length, - }); - - const rejection = hardeningGate(toolName, args, token); - if (rejection) { - otel.endSpan(span, { - status: "error", - error: "gate_rejected", - attributes: { "mcp.rejection.code": rejection.code }, - }); - sendError(id, rejection.code, rejection.message); - break; - } - - try { - const result = await dispatchTool(toolName, args); - if (result === null) { - otel.endSpan(span, { status: "error", error: "unknown_tool" }); - sendError(id, -32601, "Unknown tool"); - } else { - otel.endSpan(span, { status: "ok" }); - sendResult(id, { - content: [{ type: "text", text: JSON.stringify(result, null, 2) }], - }); - } - } catch (e) { - otel.endSpan(span, { - status: "error", - error: sanitizeErrorMessage(e?.message ?? String(e)), - }); - throw e; - } - break; - } - - case "groups/list": { - const groups = await fetchGroups(); - sendResult(id, { groups }); - break; - } - - case "groups/call": { - const groupId = params?.groupId; - const toolName = params?.name; - const args = params?.arguments || {}; - - if (!groupId) { - sendError(id, -32602, "Group ID is required"); - break; - } - - const rejection = hardeningGate(toolName, args); - if (rejection) { - sendError(id, rejection.code, rejection.message); - break; - } - - const result = await dispatchGroupTool(groupId, toolName, args); - if (result === null) { - sendError(id, -32601, "Unknown tool or group"); - } else { - sendResult(id, { - content: [{ type: "text", text: JSON.stringify(result, null, 2) }], - }); - } - break; - } - - case "prompts/list": { - const prompts = await fetchPrompts(); - sendResult(id, { prompts }); - break; - } - - case "prompts/get": { - const promptId = params?.promptId; - if (!promptId) { - sendError(id, -32602, "Prompt ID is required"); - break; - } - - const prompt = await fetchPrompt(promptId); - if (prompt === null) { - sendError(id, -32601, "Unknown prompt"); - } else { - sendResult(id, { prompt }); - } - break; - } - - case "resources/list": { - const resources = await fetchResources(); - sendResult(id, { resources }); - break; - } - - case "resources/get": { - const resourceId = params?.resourceId; - if (!resourceId) { - sendError(id, -32602, "Resource ID is required"); - break; - } - - const resource = await fetchResource(resourceId); - if (resource === null) { - sendError(id, -32601, "Unknown resource"); - } else { - sendResult(id, { resource }); - } - break; - } - - case "ping": - sendResult(id, {}); - break; - - default: - if (id !== undefined) { - sendError(id, -32601, "Method not found"); - } - } + const response = await dispatchMcpMessage(msg, { transport: "stdio" }); + if (response !== null) send(response); } -// =================================================================== -// Main I/O loop — runtime-agnostic -// =================================================================== - function processChunk(chunk) { buffer += decoder.decode(chunk); if (buffer.length > MAX_BUFFER_BYTES) { @@ -570,23 +67,74 @@ function processChunk(chunk) { const line = buffer.slice(0, boundary).trim(); buffer = buffer.slice(boundary + 1); if (line.length > 0) { - const p = handleMessage(line).catch(() => {}); + const p = handleStdioLine(line).catch(() => {}); pendingMessages.push(p); } } } -if (typeof Deno !== "undefined") { - for await (const chunk of Deno.stdin.readable) { - processChunk(chunk); +async function runStdio() { + if (typeof Deno !== "undefined") { + for await (const chunk of Deno.stdin.readable) processChunk(chunk); + } else { + // @ts-ignore: process is global in Node + for await (const chunk of process.stdin) processChunk(chunk); } -} else { - // @ts-ignore: process is global in Node - for await (const chunk of process.stdin) { - processChunk(chunk); + await Promise.allSettled(pendingMessages); +} + +// =================================================================== +// HTTP transport +// =================================================================== + +async function runHttp() { + const handle = await startHttpTransport(); + await new Promise((resolve) => { + const stop = async () => { + info("MCP HTTP transport shutting down"); + await handle.stop(); + resolve(); + }; + if (typeof Deno !== "undefined") { + Deno.addSignalListener("SIGINT", stop); + try { Deno.addSignalListener("SIGTERM", stop); } catch { /* not supported on all platforms */ } + } else { + // @ts-ignore: process is global in Node + process.on("SIGINT", stop); + // @ts-ignore + process.on("SIGTERM", stop); + } + }); +} + +// =================================================================== +// Main +// =================================================================== + +async function main() { + if (TRANSPORT === "stdio") { + await runStdio(); + } else if (TRANSPORT === "http") { + await runHttp(); + } else if (TRANSPORT === "both") { + await Promise.all([runStdio(), runHttp()]); + } else { + logError("Unknown BOJ_TRANSPORT", { value: TRANSPORT, supported: ["stdio", "http", "both"] }); + if (typeof Deno !== "undefined") Deno.exit(2); + // @ts-ignore: process is global in Node + else process.exit(2); } } -await Promise.allSettled(pendingMessages); +try { + await main(); +} catch (e) { + logError("MCP bridge fatal", { error: e?.message ?? String(e) }); + if (typeof Deno !== "undefined") Deno.exit(1); + // @ts-ignore: process is global in Node + else process.exit(1); +} + if (typeof Deno !== "undefined") Deno.exit(0); +// @ts-ignore: process is global in Node else process.exit(0); diff --git a/mcp-bridge/tests/http_transport_test.js b/mcp-bridge/tests/http_transport_test.js new file mode 100644 index 00000000..b6c66df5 --- /dev/null +++ b/mcp-bridge/tests/http_transport_test.js @@ -0,0 +1,333 @@ +// SPDX-License-Identifier: MPL-2.0 +// Copyright (c) 2026 Jonathan D.A. Jewell (hyperpolymath) +// +// BoJ Server — Streamable HTTP transport tests (ADR-0013, PR1). +// +// Covers dispatch-parity with stdio, auth modes, session lifecycle, and +// the loopback-refuse safety gate. Mirrors `dispatch_test.js` patterns: +// `node --test`, runtime-neutral assertions, no fixtures of network +// services beyond loopback. +// +// Run: node --test mcp-bridge/tests/http_transport_test.js + +import { test } from "node:test"; +import assert from "node:assert/strict"; + +import { startHttpTransport, _internals } from "../lib/http-transport.js"; + +const { SessionManager, checkAuth, parseTokens, isLoopback, handleMcpPost } = _internals; + +// Reserve an ephemeral port via the kernel rather than guessing. +async function pickEphemeralPort() { + const { createServer } = await import("node:net"); + return new Promise((resolve, reject) => { + const srv = createServer(); + srv.unref(); + srv.on("error", reject); + srv.listen(0, "127.0.0.1", () => { + const port = srv.address().port; + srv.close(() => resolve(port)); + }); + }); +} + +async function withServer(opts, fn) { + const port = opts.port ?? (await pickEphemeralPort()); + const handle = await startHttpTransport({ ...opts, port }); + try { + return await fn({ handle, port }); + } finally { + await handle.stop(); + } +} + +// ----------------------------------------------------------------- +// 1. SessionManager — create/touch/delete/expire semantics +// ----------------------------------------------------------------- +test("SessionManager: create returns unique IDs", () => { + const sm = new SessionManager(); + const a = sm.create(); + const b = sm.create(); + assert.notEqual(a, b); + assert.ok(sm.get(a)); + assert.ok(sm.get(b)); +}); + +test("SessionManager: touch fails for unknown sessions", () => { + const sm = new SessionManager(); + assert.equal(sm.touch("not-a-real-uuid"), false); +}); + +test("SessionManager: expireIdle drops sessions older than the timeout", () => { + const sm = new SessionManager({ timeoutMs: 1000 }); + const id = sm.create(); + // simulate the session having been idle for two hours + sm.get(id).lastSeenMs = Date.now() - 2 * 60 * 60 * 1000; + const removed = sm.expireIdle(); + assert.equal(removed, 1); + assert.equal(sm.get(id), undefined); +}); + +// ----------------------------------------------------------------- +// 2. checkAuth — bearer reject / accept; none permissive +// ----------------------------------------------------------------- +test("checkAuth: bearer accepts a known token", () => { + const result = checkAuth({ + authMode: "bearer", + tokens: new Set(["alpha", "beta"]), + headerValue: "Bearer alpha", + }); + assert.equal(result.ok, true); +}); + +test("checkAuth: bearer rejects unknown tokens with 401", () => { + const result = checkAuth({ + authMode: "bearer", + tokens: new Set(["alpha"]), + headerValue: "Bearer wrong", + }); + assert.equal(result.ok, false); + assert.equal(result.code, 401); +}); + +test("checkAuth: bearer rejects missing Authorization header", () => { + const result = checkAuth({ authMode: "bearer", tokens: new Set(["x"]), headerValue: undefined }); + assert.equal(result.ok, false); + assert.equal(result.code, 401); +}); + +test("checkAuth: bearer rejects malformed Authorization header", () => { + const result = checkAuth({ authMode: "bearer", tokens: new Set(["x"]), headerValue: "Basic xyz" }); + assert.equal(result.ok, false); + assert.equal(result.code, 401); +}); + +test("checkAuth: none mode always passes", () => { + const result = checkAuth({ authMode: "none", tokens: new Set(), headerValue: undefined }); + assert.equal(result.ok, true); +}); + +// ----------------------------------------------------------------- +// 3. parseTokens — CSV normalisation +// ----------------------------------------------------------------- +test("parseTokens: trims whitespace and drops empty entries", () => { + const got = parseTokens(" a, b ,, c "); + assert.deepEqual([...got].sort(), ["a", "b", "c"]); +}); + +test("parseTokens: empty input is an empty Set", () => { + assert.equal(parseTokens("").size, 0); + assert.equal(parseTokens(undefined).size, 0); +}); + +// ----------------------------------------------------------------- +// 4. isLoopback — host classification +// ----------------------------------------------------------------- +test("isLoopback: recognises loopback addresses", () => { + for (const h of ["127.0.0.1", "localhost", "::1"]) assert.equal(isLoopback(h), true); +}); + +test("isLoopback: rejects 0.0.0.0 and public addresses", () => { + for (const h of ["0.0.0.0", "10.0.0.1", "203.0.113.7"]) assert.equal(isLoopback(h), false); +}); + +// ----------------------------------------------------------------- +// 5. startup-safety — auth=none + non-loopback refuses to start +// ----------------------------------------------------------------- +test("startHttpTransport: refuses auth=none on a non-loopback bind", async () => { + await assert.rejects( + () => startHttpTransport({ port: 0, bind: "0.0.0.0", authMode: "none", tokens: [] }), + /Refusing to start.*non-loopback/i, + ); +}); + +test("startHttpTransport: refuses bearer mode with no tokens", async () => { + await assert.rejects( + () => startHttpTransport({ port: 0, bind: "127.0.0.1", authMode: "bearer", tokens: [] }), + /requires at least one token/i, + ); +}); + +// ----------------------------------------------------------------- +// 6. handleMcpPost — pure-ish dispatch (no live socket) +// ----------------------------------------------------------------- +test("handleMcpPost: initialize mints a session id and returns the MCP serverInfo", async () => { + const sessions = new SessionManager(); + const out = await handleMcpPost({ + rawBody: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize", params: { clientInfo: { name: "t" } } }), + headers: {}, + sessions, + auth: { authMode: "none", tokens: new Set() }, + }); + assert.equal(out.status, 200); + assert.ok(out.headers["Mcp-Session-Id"], "session id header must be present"); + assert.equal(out.body.result.serverInfo.name, "boj-server"); + // session is recorded + assert.ok(sessions.get(out.headers["Mcp-Session-Id"])); +}); + +test("handleMcpPost: rejects unknown session id with -32001", async () => { + const sessions = new SessionManager(); + const out = await handleMcpPost({ + rawBody: JSON.stringify({ jsonrpc: "2.0", id: 2, method: "tools/list" }), + headers: { "mcp-session-id": "not-a-known-session" }, + sessions, + auth: { authMode: "none", tokens: new Set() }, + }); + assert.equal(out.status, 404); + assert.equal(out.body.error.code, -32001); +}); + +test("handleMcpPost: bearer auth rejects unauthenticated requests", async () => { + const sessions = new SessionManager(); + const out = await handleMcpPost({ + rawBody: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize" }), + headers: {}, + sessions, + auth: { authMode: "bearer", tokens: new Set(["secret"]) }, + }); + assert.equal(out.status, 401); +}); + +test("handleMcpPost: bearer auth accepts a valid token", async () => { + const sessions = new SessionManager(); + const out = await handleMcpPost({ + rawBody: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize" }), + headers: { authorization: "Bearer secret" }, + sessions, + auth: { authMode: "bearer", tokens: new Set(["secret"]) }, + }); + assert.equal(out.status, 200); +}); + +test("handleMcpPost: malformed JSON returns -32700 Parse error", async () => { + const sessions = new SessionManager(); + const out = await handleMcpPost({ + rawBody: "{not json", + headers: {}, + sessions, + auth: { authMode: "none", tokens: new Set() }, + }); + assert.equal(out.status, 400); + assert.equal(out.body.error.code, -32700); +}); + +test("handleMcpPost: notifications/initialized returns 202 with empty body", async () => { + const sessions = new SessionManager(); + // First, mint a session via initialize so the notification carries a valid id. + const init = await handleMcpPost({ + rawBody: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize" }), + headers: {}, + sessions, + auth: { authMode: "none", tokens: new Set() }, + }); + const sessionId = init.headers["Mcp-Session-Id"]; + const out = await handleMcpPost({ + rawBody: JSON.stringify({ jsonrpc: "2.0", method: "notifications/initialized" }), + headers: { "mcp-session-id": sessionId }, + sessions, + auth: { authMode: "none", tokens: new Set() }, + }); + assert.equal(out.status, 202); + assert.equal(out.body, ""); +}); + +// ----------------------------------------------------------------- +// 7. End-to-end through the real listener — initialize + tools/list +// ----------------------------------------------------------------- +test("HTTP listener: end-to-end initialize then tools/list with session header", async () => { + await withServer({ bind: "127.0.0.1", authMode: "none", tokens: [] }, async ({ port }) => { + const initRes = await fetch(`http://127.0.0.1:${port}/mcp`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize", params: { clientInfo: { name: "e2e" } } }), + }); + assert.equal(initRes.status, 200); + const sessionId = initRes.headers.get("mcp-session-id"); + assert.ok(sessionId, "Mcp-Session-Id must be issued on initialize"); + const initBody = await initRes.json(); + assert.equal(initBody.result.serverInfo.name, "boj-server"); + + const toolsRes = await fetch(`http://127.0.0.1:${port}/mcp`, { + method: "POST", + headers: { "Content-Type": "application/json", "Mcp-Session-Id": sessionId }, + body: JSON.stringify({ jsonrpc: "2.0", id: 2, method: "tools/list" }), + }); + assert.equal(toolsRes.status, 200); + const toolsBody = await toolsRes.json(); + assert.ok(Array.isArray(toolsBody.result.tools)); + assert.ok(toolsBody.result.tools.length > 0); + }); +}); + +test("HTTP listener: bearer mode rejects request without token, accepts with token", async () => { + await withServer({ bind: "127.0.0.1", authMode: "bearer", tokens: ["s3cret"] }, async ({ port }) => { + const noAuth = await fetch(`http://127.0.0.1:${port}/mcp`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize" }), + }); + assert.equal(noAuth.status, 401); + + const withAuth = await fetch(`http://127.0.0.1:${port}/mcp`, { + method: "POST", + headers: { "Content-Type": "application/json", "Authorization": "Bearer s3cret" }, + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize" }), + }); + assert.equal(withAuth.status, 200); + }); +}); + +test("HTTP listener: GET /healthz returns 200 without auth", async () => { + await withServer({ bind: "127.0.0.1", authMode: "bearer", tokens: ["t"] }, async ({ port }) => { + const res = await fetch(`http://127.0.0.1:${port}/healthz`); + assert.equal(res.status, 200); + }); +}); + +test("HTTP listener: unknown path 404s", async () => { + await withServer({ bind: "127.0.0.1", authMode: "none", tokens: [] }, async ({ port }) => { + const res = await fetch(`http://127.0.0.1:${port}/nope`); + assert.equal(res.status, 404); + }); +}); + +test("HTTP listener: DELETE /mcp tears down a session", async () => { + await withServer({ bind: "127.0.0.1", authMode: "none", tokens: [] }, async ({ port, handle }) => { + const initRes = await fetch(`http://127.0.0.1:${port}/mcp`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "initialize" }), + }); + const sessionId = initRes.headers.get("mcp-session-id"); + assert.ok(handle.sessions.get(sessionId)); + const del = await fetch(`http://127.0.0.1:${port}/mcp`, { + method: "DELETE", + headers: { "Mcp-Session-Id": sessionId }, + }); + assert.equal(del.status, 204); + assert.equal(handle.sessions.get(sessionId), undefined); + }); +}); + +// ----------------------------------------------------------------- +// 8. boj_capabilities deployment resource +// ----------------------------------------------------------------- +test("resources: boj://capabilities/deployment lists the 5 local-only cartridges", async () => { + const { readResource, listResources } = await import("../lib/resources.js"); + const uris = listResources().map((r) => r.uri); + assert.ok(uris.includes("boj://capabilities/deployment"), "must advertise the deployment-capabilities resource"); + const r = await readResource("boj://capabilities/deployment"); + assert.ok(r && r.contents[0].text); + const payload = JSON.parse(r.contents[0].text); + const names = payload.local_only_cartridges.map((c) => c.name).sort(); + assert.deepEqual( + names, + ["browser-mcp", "container-mcp", "ffmpeg-mcp", "local-coord-mcp", "sandbox-mcp"], + "ADR-0013 PR1 names 5 local-only cartridges", + ); + for (const c of payload.local_only_cartridges) { + assert.equal(c.requires_local, true); + assert.ok(typeof c.reason === "string" && c.reason.length > 0); + } +});