/** * Lightweight MCP HTTP/SSE client. * Supports tools/list and tools/call via: * 1. Streamable HTTP (POST, response is JSON or SSE stream) * 2. Legacy SSE session transport (GET /sse → session URL, then POST) */ export interface McpTool { name: string; description?: string; inputSchema?: Record; } export interface McpCallResult { content: Array<{ type: string; text?: string }>; isError?: boolean; } interface McpRequest { jsonrpc: '2.0'; id: number; method: string; params?: unknown; } interface McpResponse { jsonrpc: '2.0'; id: number; result?: unknown; error?: { code: number; message: string }; } /** Parse a single SSE event stream looking for the first data: line with valid JSON-RPC */ function parseSseResponse(text: string): McpResponse | null { for (const line of text.split('\n')) { if (line.startsWith('data:')) { const json = line.slice(5).trim(); try { return JSON.parse(json) as McpResponse; } catch { // continue looking } } } return null; } export class McpClient { private url: string; private apiKey?: string; private timeout: number; /** Resolved session endpoint for legacy SSE transport, or null if using streamable HTTP */ private sessionUrl: string | null = null; private transportDetected = false; constructor(url: string, apiKey?: string, timeoutMs = 15000) { this.url = url.replace(/\/$/, ''); this.apiKey = apiKey; this.timeout = timeoutMs; } private headers(extra?: Record): Record { const h: Record = { 'Content-Type': 'application/json', Accept: 'application/json, text/event-stream' }; if (this.apiKey) h['Authorization'] = `Bearer ${this.apiKey}`; return { ...h, ...extra }; } /** * Try legacy SSE transport: GET /sse → read `endpoint` event → set sessionUrl. * Returns true if the server supports the legacy SSE handshake. */ private async tryLegacySseHandshake(): Promise { const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), Math.min(this.timeout, 5000)); try { const sseUrl = `${this.url}/sse`; const res = await fetch(sseUrl, { headers: { Accept: 'text/event-stream', ...(this.apiKey ? { Authorization: `Bearer ${this.apiKey}` } : {}) }, signal: controller.signal }); if (!res.ok || !res.headers.get('content-type')?.includes('event-stream')) return false; // Read only first few KB to find endpoint event const reader = res.body?.getReader(); if (!reader) return false; let accumulated = ''; let found = false; while (accumulated.length < 4096) { const { done, value } = await reader.read(); if (done) break; accumulated += new TextDecoder().decode(value); // Look for "event: endpoint\ndata: " const match = accumulated.match(/event:\s*endpoint\s*\ndata:\s*(\S+)/); if (match) { const endpoint = match[1].trim(); this.sessionUrl = endpoint.startsWith('http') ? endpoint : `${this.url}${endpoint}`; found = true; reader.cancel().catch(() => {}); break; } } return found; } catch { return false; } finally { clearTimeout(timer); } } /** * Send MCP initialize + initialized handshake. * Per MCP spec, server expects initialize request before any other method. * We ignore errors here (many simple servers skip this step). */ private async sendInitialize(endpoint: string): Promise { const initReq: McpRequest = { jsonrpc: '2.0', id: Date.now(), method: 'initialize', params: { protocolVersion: '2024-11-05', capabilities: { tools: {} }, clientInfo: { name: 'iistwin-ai-bot', version: '1.0.0' }, }, }; try { await this.rpcPost(endpoint, initReq); // Send initialized notification (no response expected) const notif: { jsonrpc: '2.0'; method: string } = { jsonrpc: '2.0', method: 'notifications/initialized' }; const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), 3000); try { await fetch(endpoint, { method: 'POST', headers: this.headers(), body: JSON.stringify(notif), signal: controller.signal }); } catch { /* notification; ignore errors */ } finally { clearTimeout(timer); } } catch { /* many servers don't require initialization; ignore */ } } private async rpc(method: string, params?: unknown): Promise { const body: McpRequest = { jsonrpc: '2.0', id: Date.now(), method, params }; // Detect transport on first call if (!this.transportDetected) { this.transportDetected = true; // First try streamable HTTP (primary MCP 2024 spec) try { await this.sendInitialize(this.url); const result = await this.rpcPost(this.url, body); return result; } catch (primaryErr) { // Fall back to legacy SSE transport const legacyOk = await this.tryLegacySseHandshake(); if (legacyOk && this.sessionUrl) { await this.sendInitialize(this.sessionUrl); return this.rpcPost(this.sessionUrl, body); } throw primaryErr; } } const endpoint = this.sessionUrl ?? this.url; return this.rpcPost(endpoint, body); } private async rpcPost(endpoint: string, body: McpRequest): Promise { const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), this.timeout); try { const res = await fetch(endpoint, { method: 'POST', headers: this.headers(), body: JSON.stringify(body), signal: controller.signal, }); if (!res.ok) { throw new Error(`MCP server responded with ${res.status}: ${await res.text().catch(() => '')}`); } const contentType = res.headers.get('content-type') ?? ''; let data: McpResponse; if (contentType.includes('event-stream')) { // SSE response for streamable HTTP transport const text = await res.text(); const parsed = parseSseResponse(text); if (!parsed) throw new Error('MCP SSE response contained no valid JSON-RPC data'); data = parsed; } else { data = await res.json() as McpResponse; } if (data.error) { throw new Error(`MCP error ${data.error.code}: ${data.error.message}`); } return data.result; } finally { clearTimeout(timer); } } async listTools(): Promise { const result = await this.rpc('tools/list') as { tools?: McpTool[] }; return result?.tools ?? []; } async callTool(name: string, args: Record): Promise { const result = await this.rpc('tools/call', { name, arguments: args }) as McpCallResult; return result ?? { content: [] }; } async ping(): Promise<{ ok: boolean; error?: string }> { try { await this.rpc('tools/list'); return { ok: true }; } catch (err: unknown) { return { ok: false, error: err instanceof Error ? err.message : String(err) }; } } }