| 1 | /** |
| 2 | * A small client for MCP servers over Streamable HTTP |
| 3 | * (docs.g1t.sh/guides/agent-abilities/, "MCP servers"): lists a server's |
| 4 | * tools when an owner adds it, and calls one when an agent's ability |
| 5 | * allows. JSON-RPC, one request per call, a fresh session each time; an |
| 6 | * answer may come back as JSON or as a server-sent event stream. No |
| 7 | * sign-in: servers that need a key aren't supported yet, and the docs say |
| 8 | * so. Pure apart from `fetch`, so it is tested with a fake. |
| 9 | */ |
| 10 | import type { McpTool } from "@g1t/contracts"; |
| 11 | |
| 12 | import { MAX_MCP_TOOLS } from "../../../packages/contracts/src/abilities.ts"; |
| 13 | |
| 14 | const PROTOCOL = "2025-06-18"; |
| 15 | /** How long one request to a server may take. */ |
| 16 | const TIMEOUT_MS = 15_000; |
| 17 | /** The most of a tool's answer an agent is given. */ |
| 18 | const MAX_TEXT = 40_000; |
| 19 | |
| 20 | export type Fetch = (url: string, init: RequestInit) => Promise<Response>; |
| 21 | |
| 22 | type Rpc = { id?: number | string | null; result?: unknown; error?: { code?: number; message?: string } }; |
| 23 | |
| 24 | /** One JSON-RPC exchange with the server; the body of the matching answer, or why there is none. */ |
| 25 | async function exchange(fetchFn: Fetch, url: string, session: string | null, body: Record<string, unknown>, expectAnswer: boolean): Promise<{ ok: true; answer: Rpc | null; session: string | null } | { ok: false; message: string }> { |
| 26 | let response: Response; |
| 27 | try { |
| 28 | response = await fetchFn(url, { |
| 29 | method: "POST", |
| 30 | headers: { |
| 31 | "content-type": "application/json", |
| 32 | accept: "application/json, text/event-stream", |
| 33 | "mcp-protocol-version": PROTOCOL, |
| 34 | ...(session ? { "mcp-session-id": session } : {}), |
| 35 | }, |
| 36 | body: JSON.stringify({ jsonrpc: "2.0", ...body }), |
| 37 | signal: AbortSignal.timeout(TIMEOUT_MS), |
| 38 | redirect: "error", |
| 39 | }); |
| 40 | } catch (error) { |
| 41 | return { ok: false, message: `The server couldn't be reached (${String(error).slice(0, 120)}).` }; |
| 42 | } |
| 43 | const next = response.headers.get("mcp-session-id") ?? session; |
| 44 | if (!expectAnswer) return { ok: true, answer: null, session: next }; |
| 45 | if (!response.ok) return { ok: false, message: `The server answered ${response.status}.` }; |
| 46 | const type = (response.headers.get("content-type") ?? "").toLowerCase(); |
| 47 | let text: string; |
| 48 | try { |
| 49 | text = await response.text(); |
| 50 | } catch { |
| 51 | return { ok: false, message: "The server's answer couldn't be read." }; |
| 52 | } |
| 53 | const wanted = body.id; |
| 54 | const candidates: Rpc[] = []; |
| 55 | if (type.includes("text/event-stream")) { |
| 56 | for (const event of text.split(/\n\n+/)) { |
| 57 | const data = event |
| 58 | .split("\n") |
| 59 | .filter((line) => line.startsWith("data:")) |
| 60 | .map((line) => line.slice(5).trim()) |
| 61 | .join("\n"); |
| 62 | if (!data) continue; |
| 63 | try { |
| 64 | const parsed = JSON.parse(data) as Rpc | Rpc[]; |
| 65 | candidates.push(...(Array.isArray(parsed) ? parsed : [parsed])); |
| 66 | } catch { |
| 67 | // Not JSON: a keep-alive or a comment. |
| 68 | } |
| 69 | } |
| 70 | } else { |
| 71 | try { |
| 72 | const parsed = JSON.parse(text) as Rpc | Rpc[]; |
| 73 | candidates.push(...(Array.isArray(parsed) ? parsed : [parsed])); |
| 74 | } catch { |
| 75 | return { ok: false, message: "The server didn't answer with JSON-RPC." }; |
| 76 | } |
| 77 | } |
| 78 | const answer = candidates.find((c) => c && typeof c === "object" && c.id === wanted) ?? null; |
| 79 | if (!answer) return { ok: false, message: "The server didn't answer the request." }; |
| 80 | return { ok: true, answer, session: next }; |
| 81 | } |
| 82 | |
| 83 | /** Opens a session: initialize, then the initialized notification. The session id, if the server gave one. */ |
| 84 | async function handshake(fetchFn: Fetch, url: string): Promise<{ ok: true; session: string | null } | { ok: false; message: string }> { |
| 85 | const opened = await exchange( |
| 86 | fetchFn, |
| 87 | url, |
| 88 | null, |
| 89 | { id: 1, method: "initialize", params: { protocolVersion: PROTOCOL, capabilities: {}, clientInfo: { name: "g1t-agents", version: "1" } } }, |
| 90 | true, |
| 91 | ); |
| 92 | if (!opened.ok) return opened; |
| 93 | if (opened.answer?.error) return { ok: false, message: `The server refused to start: ${opened.answer.error.message ?? "no reason given"}.` }; |
| 94 | const told = await exchange(fetchFn, url, opened.session, { method: "notifications/initialized" }, false); |
| 95 | return { ok: true, session: told.ok ? told.session : opened.session }; |
| 96 | } |
| 97 | |
| 98 | /** A tool as the server lists it, as kept: its kind from `readOnlyHint`, a write when the server doesn't say. */ |
| 99 | function toolOf(raw: unknown): McpTool | null { |
| 100 | if (!raw || typeof raw !== "object") return null; |
| 101 | const t = raw as { name?: unknown; description?: unknown; inputSchema?: unknown; annotations?: { readOnlyHint?: unknown } }; |
| 102 | if (typeof t.name !== "string" || !/^[A-Za-z0-9_.-]{1,64}$/.test(t.name)) return null; |
| 103 | const schema = t.inputSchema && typeof t.inputSchema === "object" && !Array.isArray(t.inputSchema) ? (t.inputSchema as Record<string, unknown>) : { type: "object", properties: {} }; |
| 104 | return { |
| 105 | name: t.name, |
| 106 | description: typeof t.description === "string" ? t.description.trim().slice(0, 500) : "", |
| 107 | kind: t.annotations?.readOnlyHint === true ? "read" : "write", |
| 108 | input_schema: JSON.stringify(schema).length > 20_000 ? { type: "object", properties: {} } : schema, |
| 109 | }; |
| 110 | } |
| 111 | |
| 112 | /** The server's tools, at most `MAX_MCP_TOOLS`, or why they couldn't be listed. */ |
| 113 | export async function listMcpTools(url: string, fetchFn: Fetch = (u, init) => fetch(u, init)): Promise<{ ok: true; tools: McpTool[] } | { ok: false; message: string }> { |
| 114 | const opened = await handshake(fetchFn, url); |
| 115 | if (!opened.ok) return opened; |
| 116 | const listed = await exchange(fetchFn, url, opened.session, { id: 2, method: "tools/list", params: {} }, true); |
| 117 | if (!listed.ok) return listed; |
| 118 | if (listed.answer?.error) return { ok: false, message: `The server couldn't list its tools: ${listed.answer.error.message ?? "no reason given"}.` }; |
| 119 | const raw = (listed.answer?.result as { tools?: unknown } | undefined)?.tools; |
| 120 | if (!Array.isArray(raw)) return { ok: false, message: "The server listed no tools." }; |
| 121 | const tools = raw.map(toolOf).filter((tool): tool is McpTool => tool !== null); |
| 122 | const names = new Set<string>(); |
| 123 | return { ok: true, tools: tools.filter((tool) => (names.has(tool.name) ? false : (names.add(tool.name), true))).slice(0, MAX_MCP_TOOLS) }; |
| 124 | } |
| 125 | |
| 126 | /** Calls one tool; the text it returned (its text parts joined), or why it didn't work. */ |
| 127 | export async function callMcpTool(url: string, name: string, args: Record<string, unknown>, fetchFn: Fetch = (u, init) => fetch(u, init)): Promise<{ ok: true; text: string } | { ok: false; message: string }> { |
| 128 | const opened = await handshake(fetchFn, url); |
| 129 | if (!opened.ok) return opened; |
| 130 | const called = await exchange(fetchFn, url, opened.session, { id: 3, method: "tools/call", params: { name, arguments: args } }, true); |
| 131 | if (!called.ok) return called; |
| 132 | if (called.answer?.error) return { ok: false, message: called.answer.error.message ?? "the server refused" }; |
| 133 | const result = called.answer?.result as { content?: unknown; isError?: unknown; structuredContent?: unknown } | undefined; |
| 134 | const parts = Array.isArray(result?.content) ? (result!.content as { type?: string; text?: string }[]) : []; |
| 135 | let text = parts |
| 136 | .map((part) => (part?.type === "text" && typeof part.text === "string" ? part.text : part?.type ? `[${part.type} content]` : "")) |
| 137 | .filter(Boolean) |
| 138 | .join("\n"); |
| 139 | if (!text && result?.structuredContent !== undefined) text = JSON.stringify(result.structuredContent).slice(0, MAX_TEXT); |
| 140 | if (result?.isError === true) return { ok: false, message: text.slice(0, 500) || "the tool reported an error" }; |
| 141 | return { ok: true, text: (text || "(no content)").slice(0, MAX_TEXT) }; |
| 142 | } |