diff --git a/.changeset/core-mcp-sse-responses.md b/.changeset/core-mcp-sse-responses.md new file mode 100644 index 0000000..aa59784 --- /dev/null +++ b/.changeset/core-mcp-sse-responses.md @@ -0,0 +1,5 @@ +--- +"@call-e/core": patch +--- + +Decode `text/event-stream` MCP responses instead of returning an empty result. The client reads the stream incrementally, returns the JSON-RPC response that matches the request id and stops reading as soon as it arrives. JSON-RPC errors in the stream are surfaced as errors. Streams that are empty, truncated, malformed, larger than 8 MiB or longer than 10,000 events are rejected with an `mcp_protocol_error`. The limits can be overridden with `maxSseResponseBytes` and `maxSseEvents` in the client config. diff --git a/packages/core/README.md b/packages/core/README.md index 7524ef2..1205725 100644 --- a/packages/core/README.md +++ b/packages/core/README.md @@ -60,6 +60,19 @@ The override applies only to the tool call. MCP session initialization keeps using `config.timeoutSeconds`; when the override is omitted, the tool call uses that configured timeout too. +## Streamed Responses + +The MCP client accepts both `application/json` and `text/event-stream` +responses. For event streams it returns the JSON-RPC response whose `id` +matches the request and stops reading once it arrives. It skips interleaved +server notifications. A JSON-RPC `error` in the stream is thrown as an +`McpHttpError` with `code: "mcp_error"`, the same as for JSON responses. + +A stream that ends without a matching response, contains malformed JSON, +exceeds 8 MiB, or carries more than 10,000 events is rejected with +`code: "mcp_protocol_error"`. `config.maxSseResponseBytes` and +`config.maxSseEvents` override those limits. + ## Tool Result Payloads `callMcpTool` returns the MCP `CallToolResult` envelope and preserves its raw diff --git a/packages/core/lib/mcp-client.d.ts b/packages/core/lib/mcp-client.d.ts index e910935..1ed5a76 100644 --- a/packages/core/lib/mcp-client.d.ts +++ b/packages/core/lib/mcp-client.d.ts @@ -9,6 +9,10 @@ export interface McpClientConfig { mcpClientName?: string; mcpClientVersion?: string; cliVersion?: string; + /** Byte cap for one text/event-stream response. Defaults to 8 MiB. */ + maxSseResponseBytes?: number; + /** Event cap for one text/event-stream response. Defaults to 10,000. */ + maxSseEvents?: number; } export interface McpHttpErrorOptions { diff --git a/packages/core/lib/mcp-client.js b/packages/core/lib/mcp-client.js index 2c412a7..1213410 100644 --- a/packages/core/lib/mcp-client.js +++ b/packages/core/lib/mcp-client.js @@ -50,7 +50,215 @@ function parseResponseBody(text) { return JSON.parse(text); } -async function requestJsonRpc(fetchImpl, url, { headers, payload, timeoutMs }) { +// Upper bounds for a single text/event-stream response. The stream is rejected +// with an mcp_protocol_error once either limit is exceeded. +const DEFAULT_MAX_SSE_RESPONSE_BYTES = 8 * 1024 * 1024; +const DEFAULT_MAX_SSE_EVENTS = 10_000; + +function positiveLimit(value, fallback) { + const limit = Number(value); + return Number.isFinite(limit) && limit > 0 ? Math.floor(limit) : fallback; +} + +function sseLimits(config) { + return { + maxBytes: positiveLimit(config.maxSseResponseBytes, DEFAULT_MAX_SSE_RESPONSE_BYTES), + maxEvents: positiveLimit(config.maxSseEvents, DEFAULT_MAX_SSE_EVENTS), + }; +} + +function isEventStreamResponse(headers) { + const contentType = String(headers["content-type"] || ""); + return contentType.split(";")[0].trim().toLowerCase() === "text/event-stream"; +} + +function createSseParser() { + let buffer = ""; + let started = false; + let eventType = ""; + let dataLines = []; + + function processLine(line) { + if (line === "") { + const event = dataLines.length > 0 + ? { event: eventType || "message", data: dataLines.join("\n") } + : null; + eventType = ""; + dataLines = []; + return event; + } + if (line.startsWith(":")) { + return null; + } + const colonIndex = line.indexOf(":"); + const field = colonIndex === -1 ? line : line.slice(0, colonIndex); + let value = colonIndex === -1 ? "" : line.slice(colonIndex + 1); + if (value.startsWith(" ")) { + value = value.slice(1); + } + if (field === "data") { + dataLines.push(value); + } else if (field === "event") { + eventType = value; + } + return null; + } + + function drain({ final }) { + const events = []; + const lineBreak = /\r\n|\r|\n/gu; + let lineStart = 0; + let match; + while ((match = lineBreak.exec(buffer)) !== null) { + if (!final && match[0] === "\r" && match.index === buffer.length - 1) { + // A CR at the end of a chunk may be the first half of a CRLF. + break; + } + const event = processLine(buffer.slice(lineStart, match.index)); + lineStart = match.index + match[0].length; + if (event) { + events.push(event); + } + } + // Anything left is an unterminated line. At the end of the stream it is + // dropped together with any event that was not closed by a blank line. + buffer = final ? "" : buffer.slice(lineStart); + return events; + } + + return { + push(chunk) { + buffer += started ? chunk : chunk.replace(/^\uFEFF/u, ""); + started = started || buffer.length > 0; + return drain({ final: false }); + }, + finish() { + return drain({ final: true }); + }, + }; +} + +function sseProtocolError(message, { responseText, headers }) { + return new McpHttpError(message, { + responseText, + headers, + code: "mcp_protocol_error", + }); +} + +async function* sseTextChunks(response, { maxBytes, onLimit }) { + const reader = typeof response.body?.getReader === "function" ? response.body.getReader() : null; + if (!reader) { + // Some fetch implementations (and test doubles) only expose text(). + const text = await response.text(); + if (Buffer.byteLength(text, "utf8") > maxBytes) { + onLimit(text); + } + yield text; + return; + } + + const decoder = new TextDecoder("utf-8", { ignoreBOM: true }); + let bytesRead = 0; + let finished = false; + try { + while (true) { + const { done, value } = await reader.read(); + if (done) { + finished = true; + const tail = decoder.decode(); + if (tail) { + yield tail; + } + return; + } + bytesRead += value.byteLength; + if (bytesRead > maxBytes) { + onLimit(decoder.decode(value, { stream: true })); + } + yield decoder.decode(value, { stream: true }); + } + } finally { + if (!finished) { + // Stop the transfer when a limit is hit, an error is thrown, or the + // matching response has already been found. + reader.cancel().catch(() => {}); + } + } +} + +async function readJsonRpcEventStream(response, { payload, headers, limits }) { + const expectsResponse = payload.id !== undefined; + const parser = createSseParser(); + let responseText = ""; + let eventCount = 0; + let orphanError = null; + + const protocolError = (message) => sseProtocolError(`${message} for ${payload.method}`, { responseText, headers }); + + // Returns the matching JSON-RPC response, or null to keep reading. + const handleEvents = (events) => { + for (const event of events) { + eventCount += 1; + if (eventCount > limits.maxEvents) { + throw protocolError(`MCP event stream exceeded ${limits.maxEvents} events`); + } + if (event.event !== "message") { + continue; + } + let message; + try { + message = JSON.parse(event.data); + } catch { + throw protocolError("MCP event stream contained malformed JSON"); + } + for (const candidate of Array.isArray(message) ? message : [message]) { + const record = objectRecord(candidate); + if (!record || record.method !== undefined || !("result" in record || "error" in record)) { + // Server-initiated requests and notifications can be interleaved before the response. + continue; + } + if (expectsResponse && record.id === payload.id) { + return record; + } + if (record.id === null && record.error && !orphanError) { + orphanError = record; + } + } + } + return null; + }; + + const chunks = sseTextChunks(response, { + maxBytes: limits.maxBytes, + onLimit(chunk) { + responseText += chunk; + throw protocolError(`MCP event stream exceeded ${limits.maxBytes} bytes`); + }, + }); + for await (const chunk of chunks) { + responseText += chunk; + const matched = handleEvents(parser.push(chunk)); + if (matched) { + // Leaving the loop cancels the rest of the stream. + return matched; + } + } + const matched = handleEvents(parser.finish()); + if (matched) { + return matched; + } + + if (orphanError) { + return orphanError; + } + if (!expectsResponse) { + return null; + } + throw protocolError("MCP event stream ended without a JSON-RPC response"); +} + +async function requestJsonRpc(fetchImpl, url, { headers, payload, timeoutMs, limits }) { const controller = new AbortController(); const timeout = setTimeout(() => controller.abort(), timeoutMs); if (typeof timeout.unref === "function") { @@ -64,22 +272,27 @@ async function requestJsonRpc(fetchImpl, url, { headers, payload, timeoutMs }) { body: JSON.stringify(payload), signal: controller.signal, }); - const text = await response.text(); - let body = null; - try { - body = parseResponseBody(text); - } catch { - body = null; - } const responseHeaders = Object.fromEntries(response.headers.entries()); + let body = null; - if (!response.ok) { - throw new McpHttpError(`MCP HTTP ${response.status} for ${payload.method}`, { - statusCode: response.status, - responseText: text, - payload: body, - headers: responseHeaders, - }); + if (response.ok && isEventStreamResponse(responseHeaders)) { + body = await readJsonRpcEventStream(response, { payload, headers: responseHeaders, limits }); + } else { + const text = await response.text(); + try { + body = parseResponseBody(text); + } catch { + body = null; + } + + if (!response.ok) { + throw new McpHttpError(`MCP HTTP ${response.status} for ${payload.method}`, { + statusCode: response.status, + responseText: text, + payload: body, + headers: responseHeaders, + }); + } } if (body?.error) { @@ -194,6 +407,7 @@ async function openMcpSession({ config, fetchImpl }) { requireFetch(fetchImpl); const accessToken = accessTokenFromCache(config); const timeoutMs = Math.max(Math.ceil(Number(config.timeoutSeconds || 15) * 1000), 1000); + const limits = sseLimits(config); const commonHeaders = { Accept: "application/json, text/event-stream", "Content-Type": "application/json", @@ -214,6 +428,7 @@ async function openMcpSession({ config, fetchImpl }) { }, }), timeoutMs, + limits, }); const sessionId = initialize.headers["mcp-session-id"] || initialize.headers["Mcp-Session-Id"] || ""; @@ -226,13 +441,14 @@ async function openMcpSession({ config, fetchImpl }) { params: {}, }), timeoutMs, + limits, }); - return { rpcHeaders, timeoutMs }; + return { rpcHeaders, timeoutMs, limits }; } export async function listMcpTools({ config, fetchImpl = globalThis.fetch } = {}) { - const { rpcHeaders, timeoutMs } = await openMcpSession({ config, fetchImpl }); + const { rpcHeaders, timeoutMs, limits } = await openMcpSession({ config, fetchImpl }); const response = await requestJsonRpc(fetchImpl, config.serverUrl, { headers: rpcHeaders, payload: buildJsonRpcPayload({ @@ -241,6 +457,7 @@ export async function listMcpTools({ config, fetchImpl = globalThis.fetch } = {} params: {}, }), timeoutMs, + limits, }); return response.body?.result ?? {}; } @@ -253,7 +470,7 @@ export async function callMcpTool({ timeoutSeconds = null, fetchImpl = globalThis.fetch, } = {}) { - const { rpcHeaders, timeoutMs } = await openMcpSession({ config, fetchImpl }); + const { rpcHeaders, timeoutMs, limits } = await openMcpSession({ config, fetchImpl }); const toolCallParams = { name: toolName, arguments: toolArguments, @@ -273,6 +490,7 @@ export async function callMcpTool({ params: toolCallParams, }), timeoutMs: toolCallTimeoutMs, + limits, }); return normalizeMcpToolResult(response.body?.result); } diff --git a/packages/core/test/core.test.js b/packages/core/test/core.test.js index 6fccdb9..d52fe2c 100644 --- a/packages/core/test/core.test.js +++ b/packages/core/test/core.test.js @@ -54,6 +54,52 @@ function jsonResponse(body, { status = 200, statusText = "OK", headers = {} } = }; } +function sseResponse(text, { status = 200, statusText = "OK", headers = {} } = {}) { + return { + ok: status >= 200 && status < 300, + status, + statusText, + headers: new Headers({ "content-type": "text/event-stream", ...headers }), + async text() { + return text; + }, + }; +} + +function sseMessage(message, { event = "message" } = {}) { + return `${event ? `event: ${event}\n` : ""}data: ${JSON.stringify(message)}\n\n`; +} + +function sseSessionFetch(respond) { + return async (_url, init) => { + const payload = JSON.parse(init.body); + if (payload.method === "initialize") { + return sseResponse(sseMessage({ jsonrpc: "2.0", id: payload.id, result: {} }), { + headers: { "mcp-session-id": "mcp-session-sse" }, + }); + } + if (payload.method === "notifications/initialized") { + assert.equal(init.headers["mcp-session-id"], "mcp-session-sse"); + return { ok: true, status: 202, statusText: "Accepted", headers: new Headers(), async text() { return ""; } }; + } + return respond(payload); + }; +} + +async function assertSseProtocolError(config, streamText, pattern) { + const fetchImpl = sseSessionFetch(() => sseResponse(streamText)); + await assert.rejects( + () => callMcpTool({ config, toolName: "plan_call", fetchImpl }), + (error) => { + assert.ok(error instanceof McpHttpError); + assert.equal(error.code, "mcp_protocol_error"); + assert.match(error.message, pattern); + assert.equal(error.responseText, streamText); + return true; + }, + ); +} + function mcpConfig(cacheRoot) { const serverUrl = "https://example.test/mcp/openagent_oauth"; writePrivateJson(tokenCachePath(cacheRoot, serverUrl), { @@ -553,3 +599,284 @@ test("MCP client reports request timeouts", async () => { }, ); }); + +test("MCP client decodes tool lists from SSE responses like JSON responses", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-tools")); + const tools = { tools: [{ name: "plan_call", inputSchema: { type: "object" } }] }; + const fromJson = await listMcpTools({ + config, + fetchImpl: async (_url, init) => { + const payload = JSON.parse(init.body); + return jsonResponse(payload.method === "tools/list" ? { jsonrpc: "2.0", id: payload.id, result: tools } : { result: {} }); + }, + }); + const fromSse = await listMcpTools({ + config, + fetchImpl: sseSessionFetch((payload) => { + assert.equal(payload.method, "tools/list"); + return sseResponse( + `: keep-alive\nid: 1\n${sseMessage({ jsonrpc: "2.0", id: payload.id, result: tools })}`, + { headers: { "content-type": "text/event-stream; charset=utf-8" } }, + ); + }), + }); + + assert.deepEqual(fromJson, tools); + assert.deepEqual(fromSse, tools); +}); + +test("MCP client joins multi-line SSE data and accepts CRLF framing", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-multiline")); + const fetchImpl = sseSessionFetch((payload) => { + const json = JSON.stringify( + { jsonrpc: "2.0", id: payload.id, result: { content: [{ type: "text", text: "ok" }], isError: false } }, + null, + 2, + ); + const data = json.split("\n").map((line) => `data: ${line}`).join("\r\n"); + return sseResponse(`event: message\r\n${data}\r\n\r\n`); + }); + + const result = await callMcpTool({ config, toolName: "plan_call", fetchImpl }); + assert.deepEqual(result, { content: [{ type: "text", text: "ok" }], isError: false }); +}); + +test("MCP client skips interleaved SSE notifications and mismatched ids", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-interleaved")); + const fetchImpl = sseSessionFetch((payload) => sseResponse([ + sseMessage({ jsonrpc: "2.0", method: "notifications/progress", params: { progress: 1 } }), + sseMessage({ jsonrpc: "2.0", id: "other-request", result: { content: [{ type: "text", text: "wrong" }] } }), + sseMessage({ jsonrpc: "2.0", id: payload.id, result: { content: [{ type: "text", text: "right" }] } }, { event: "" }), + ].join(""))); + + const result = await callMcpTool({ config, toolName: "plan_call", fetchImpl }); + assert.deepEqual(result, { content: [{ type: "text", text: "right" }] }); +}); + +test("MCP client preserves isError tool results from SSE responses", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-iserror")); + const toolResult = { content: [{ type: "text", text: "lookup failed" }], isError: true }; + const fetchImpl = sseSessionFetch((payload) => sseResponse(sseMessage({ jsonrpc: "2.0", id: payload.id, result: toolResult }))); + + const result = await callMcpTool({ config, toolName: "get_call_run", fetchImpl }); + assert.deepEqual(result, toolResult); +}); + +test("MCP client surfaces JSON-RPC errors delivered over SSE", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-error")); + const fetchImpl = sseSessionFetch((payload) => sseResponse(sseMessage({ + jsonrpc: "2.0", + id: payload.id, + error: { code: -32602, message: "invalid params" }, + }))); + + await assert.rejects( + () => callMcpTool({ config, toolName: "plan_call", fetchImpl }), + (error) => { + assert.ok(error instanceof McpHttpError); + assert.equal(error.code, "mcp_error"); + assert.deepEqual(error.payload, { code: -32602, message: "invalid params" }); + return true; + }, + ); +}); + +test("MCP client rejects SSE responses without a matching JSON-RPC response", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-missing")); + const missing = /ended without a JSON-RPC response for tools\/call/; + + await assertSseProtocolError(config, "", missing); + await assertSseProtocolError(config, ": keep-alive\n\n", missing); + await assertSseProtocolError( + config, + sseMessage({ jsonrpc: "2.0", id: "other-request", result: { content: [] } }), + missing, + ); + await assertSseProtocolError( + config, + sseMessage({ jsonrpc: "2.0", id: "calle-plan_call", result: {} }, { event: "ping" }), + missing, + ); + // A final event that is not terminated by a blank line is truncated and must not count. + await assertSseProtocolError( + config, + `data: ${JSON.stringify({ jsonrpc: "2.0", id: "calle-plan_call", result: {} })}\n`, + missing, + ); +}); + +test("MCP client rejects malformed JSON in SSE responses", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-malformed")); + + await assertSseProtocolError(config, 'data: {"jsonrpc":"2.0","id":\n\n', /malformed JSON for tools\/call/); +}); + +function streamedSseResponse(init, { onCancel = () => {}, pull, signal } = {}) { + const encoder = new TextEncoder(); + const body = new ReadableStream({ + start(controller) { + signal?.addEventListener("abort", () => { + const abortError = new Error("aborted"); + abortError.name = "AbortError"; + controller.error(abortError); + }); + for (const chunk of init) { + controller.enqueue(typeof chunk === "string" ? encoder.encode(chunk) : chunk); + } + }, + pull: pull ? (controller) => pull(controller, encoder) : undefined, + cancel: onCancel, + }); + return { + ok: true, + status: 200, + statusText: "OK", + headers: new Headers({ "content-type": "text/event-stream" }), + body, + async text() { + throw new Error("streamed SSE responses must be read through body"); + }, + }; +} + +function sseStreamFetch(respond) { + return async (url, init) => { + const payload = JSON.parse(init.body); + if (payload.method === "initialize" || payload.method === "notifications/initialized") { + return sseSessionFetch(() => null)(url, init); + } + return respond(payload, init.signal); + }; +} + +async function assertMcpProtocolError(promise, pattern) { + await assert.rejects(promise, (error) => { + assert.ok(error instanceof McpHttpError); + assert.equal(error.code, "mcp_protocol_error"); + assert.match(error.message, pattern); + return true; + }); +} + +const tick = () => new Promise((resolve) => setImmediate(resolve)); + +test("MCP client decodes streamed SSE split across arbitrary chunk boundaries", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-chunks")); + const toolResult = { content: [{ type: "text", text: "café ☎ confirmed" }] }; + const fetchImpl = sseStreamFetch((payload) => { + const text = `\uFEFF: hello\r\nevent: message\r\ndata: ${JSON.stringify({ jsonrpc: "2.0", id: payload.id, result: toolResult })}\r\n\r\n`; + const bytes = new TextEncoder().encode(text); + return streamedSseResponse([...bytes].map((byte) => Uint8Array.of(byte))); + }); + + const result = await callMcpTool({ config, toolName: "plan_call", fetchImpl }); + assert.deepEqual(result, toolResult); +}); + +test("MCP client stops reading the SSE stream once the matching response arrives", async () => { + const config = mcpConfig(makeTempRoot("calle-core-mcp-sse-early")); + let cancelled = false; + const fetchImpl = sseStreamFetch((payload, signal) => streamedSseResponse( + [ + sseMessage({ jsonrpc: "2.0", method: "notifications/progress", params: { progress: 1 } }), + sseMessage({ jsonrpc: "2.0", id: payload.id, result: { content: [{ type: "text", text: "done" }] } }), + ], + // The stream never closes, so the call only resolves if the client stops reading. + { signal, onCancel: () => { cancelled = true; } }, + )); + + const result = await callMcpTool({ config, toolName: "plan_call", fetchImpl }); + await tick(); + assert.deepEqual(result, { content: [{ type: "text", text: "done" }] }); + assert.equal(cancelled, true); +}); + +test("MCP client rejects SSE responses over the byte limit", async () => { + const config = { ...mcpConfig(makeTempRoot("calle-core-mcp-sse-bytes")), maxSseResponseBytes: 256 }; + const bigResult = { content: [{ type: "text", text: "x".repeat(512) }] }; + + // Fallback when the response only exposes text(). + await assertMcpProtocolError( + callMcpTool({ + config, + toolName: "plan_call", + fetchImpl: sseSessionFetch((payload) => sseResponse(sseMessage({ jsonrpc: "2.0", id: payload.id, result: bigResult }))), + }), + /exceeded 256 bytes for tools\/call/, + ); + + // Streaming body that never ends: the reader is cancelled once the cap is crossed. + let cancelled = false; + let chunksPulled = 0; + await assertMcpProtocolError( + callMcpTool({ + config, + toolName: "plan_call", + fetchImpl: sseStreamFetch(() => streamedSseResponse([], { + pull(controller, encoder) { + chunksPulled += 1; + controller.enqueue(encoder.encode(": padding padding padding padding\n")); + }, + onCancel: () => { cancelled = true; }, + })), + }), + /exceeded 256 bytes for tools\/call/, + ); + await tick(); + assert.equal(cancelled, true); + assert.ok(chunksPulled < 20, `expected reading to stop near the cap, pulled ${chunksPulled} chunks`); +}); + +test("MCP client rejects SSE responses with too many events", async () => { + const config = { ...mcpConfig(makeTempRoot("calle-core-mcp-sse-events")), maxSseEvents: 3 }; + const progress = sseMessage({ jsonrpc: "2.0", method: "notifications/progress", params: { progress: 1 } }); + + await assertMcpProtocolError( + callMcpTool({ + config, + toolName: "plan_call", + fetchImpl: sseSessionFetch((payload) => sseResponse( + progress.repeat(3) + sseMessage({ jsonrpc: "2.0", id: payload.id, result: {} }), + )), + }), + /exceeded 3 events for tools\/call/, + ); + + let cancelled = false; + await assertMcpProtocolError( + callMcpTool({ + config, + toolName: "plan_call", + fetchImpl: sseStreamFetch(() => streamedSseResponse([], { + pull(controller, encoder) { + controller.enqueue(encoder.encode(progress)); + }, + onCancel: () => { cancelled = true; }, + })), + }), + /exceeded 3 events for tools\/call/, + ); + await tick(); + assert.equal(cancelled, true); +}); + +test("MCP client reports timeouts while waiting on an SSE stream", async () => { + const config = { ...mcpConfig(makeTempRoot("calle-core-mcp-sse-timeout")), timeoutSeconds: 1 }; + const fetchImpl = sseStreamFetch((_payload, signal) => streamedSseResponse([": waiting\n\n"], { signal })); + // The client's timer is unref'd; a real socket would keep the process alive meanwhile. + const keepAlive = setTimeout(() => {}, 10_000); + + try { + await assert.rejects( + () => callMcpTool({ config, toolName: "plan_call", fetchImpl }), + (error) => { + assert.ok(error instanceof McpHttpError); + assert.equal(error.code, "http_error"); + assert.match(error.message, /timed out/i); + return true; + }, + ); + } finally { + clearTimeout(keepAlive); + } +});