Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/core-mcp-sse-responses.md
Original file line number Diff line number Diff line change
@@ -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.
13 changes: 13 additions & 0 deletions packages/core/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 4 additions & 0 deletions packages/core/lib/mcp-client.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
254 changes: 236 additions & 18 deletions packages/core/lib/mcp-client.js
Original file line number Diff line number Diff line change
Expand Up @@ -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") {
Expand All @@ -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) {
Expand Down Expand Up @@ -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",
Expand All @@ -214,6 +428,7 @@ async function openMcpSession({ config, fetchImpl }) {
},
}),
timeoutMs,
limits,
});

const sessionId = initialize.headers["mcp-session-id"] || initialize.headers["Mcp-Session-Id"] || "";
Expand All @@ -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({
Expand All @@ -241,6 +457,7 @@ export async function listMcpTools({ config, fetchImpl = globalThis.fetch } = {}
params: {},
}),
timeoutMs,
limits,
});
return response.body?.result ?? {};
}
Expand All @@ -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,
Expand All @@ -273,6 +490,7 @@ export async function callMcpTool({
params: toolCallParams,
}),
timeoutMs: toolCallTimeoutMs,
limits,
});
return normalizeMcpToolResult(response.body?.result);
}
Loading
Loading