Skip to content

Commit 87ba30f

Browse files
committed
fix(edge): answer JSON-mode POSTs whose requests are never answered
1 parent 628e364 commit 87ba30f

2 files changed

Lines changed: 228 additions & 7 deletions

File tree

‎src/edge/WebStreamableHTTPServerTransport.test.ts‎

Lines changed: 185 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import { WebStreamableHTTPServerTransport } from "./WebStreamableHTTPServerTrans
1111
type JsonResponse = any;
1212

1313
type TransportInternals = {
14+
_jsonResponseCollectors: Map<string, unknown>;
1415
_requestToStreamMapping: Map<number | string, string>;
1516
_streamMapping: Map<string, WritableStreamDefaultWriter<Uint8Array>>;
1617
};
@@ -27,6 +28,12 @@ const createGetRequest = (lastEventId?: string) =>
2728
},
2829
});
2930

31+
const createDeleteRequest = () =>
32+
new Request("http://localhost/mcp", {
33+
headers: { "mcp-session-id": "test-session" },
34+
method: "DELETE",
35+
});
36+
3037
const createPostRequest = (body: unknown, accept: string, sessionId?: string) =>
3138
new Request("http://localhost/mcp", {
3239
body: JSON.stringify(body),
@@ -72,6 +79,32 @@ const createSlowToolServer = (delayMs: number) => {
7279
return server;
7380
};
7481

82+
/**
83+
* Server whose only tool never answers on its own: the request stays open until
84+
* something aborts it, which is what a client cancellation or a terminated
85+
* session does.
86+
*/
87+
const createHangingToolServer = () => {
88+
const server = new Server(
89+
{ name: "TestServer", version: "1.0.0" },
90+
{ capabilities: { tools: {} } },
91+
);
92+
server.setRequestHandler(CallToolRequestSchema, (_request, extra) => {
93+
return new Promise<never>((_resolve, reject) => {
94+
extra.signal.addEventListener("abort", () =>
95+
reject(extra.signal.reason ?? new Error("aborted")),
96+
);
97+
});
98+
});
99+
return server;
100+
};
101+
102+
const createCancelledNotification = (requestId: number) => ({
103+
jsonrpc: "2.0",
104+
method: "notifications/cancelled",
105+
params: { reason: "client gave up", requestId },
106+
});
107+
75108
const createToolCall = (id: number, name: string) => ({
76109
id,
77110
jsonrpc: "2.0",
@@ -690,6 +723,158 @@ describe("WebStreamableHTTPServerTransport", () => {
690723
expect(body[0].result.content[0].text).toBe("done:one");
691724
});
692725

726+
it("should answer a JSON-mode POST with an error when the session ends first", async () => {
727+
const transport = new WebStreamableHTTPServerTransport({
728+
enableJsonResponse: true,
729+
sessionIdGenerator: () => "test-session",
730+
});
731+
await createHangingToolServer().connect(transport);
732+
await transport.handleRequest(createInitializeRequest("application/json"));
733+
734+
const pending = transport.handleRequest(
735+
createPostRequest(
736+
createToolCall(2, "never"),
737+
"application/json",
738+
"test-session",
739+
),
740+
);
741+
await flushMicrotasks();
742+
743+
const deleted = await transport.handleRequest(createDeleteRequest());
744+
expect(deleted.status).toBe(204);
745+
746+
// The tool call can no longer be answered, but the POST still owes the
747+
// client a JSON-RPC message: an empty body under a 200 with a JSON content
748+
// type is unparseable, and `response.json()` throws on it.
749+
const response = await pending;
750+
expect(response.status).toBe(200);
751+
const body: JsonResponse = await response.json();
752+
expect(body.id).toBe(2);
753+
expect(body.error.code).toBe(-32000);
754+
});
755+
756+
it("should keep the array envelope when a batch loses its session", async () => {
757+
const transport = new WebStreamableHTTPServerTransport({
758+
enableJsonResponse: true,
759+
sessionIdGenerator: () => "test-session",
760+
});
761+
await createHangingToolServer().connect(transport);
762+
await transport.handleRequest(createInitializeRequest("application/json"));
763+
764+
const pending = transport.handleRequest(
765+
createPostRequest(
766+
[createToolCall(2, "never")],
767+
"application/json",
768+
"test-session",
769+
),
770+
);
771+
await flushMicrotasks();
772+
await transport.handleRequest(createDeleteRequest());
773+
774+
// JSON-RPC forbids answering a batch with an empty array, which is what
775+
// dropping the unanswered response leaves behind.
776+
const body: JsonResponse = await (await pending).json();
777+
expect(Array.isArray(body)).toBe(true);
778+
expect(body).toHaveLength(1);
779+
expect(body[0].id).toBe(2);
780+
expect(body[0].error.code).toBe(-32000);
781+
});
782+
783+
it("should release a JSON-mode POST whose request the client cancels", async () => {
784+
const transport = new WebStreamableHTTPServerTransport({
785+
enableJsonResponse: true,
786+
sessionIdGenerator: () => "test-session",
787+
});
788+
await createHangingToolServer().connect(transport);
789+
await transport.handleRequest(createInitializeRequest("application/json"));
790+
791+
const internals = transport as unknown as TransportInternals;
792+
const pending = transport.handleRequest(
793+
createPostRequest(
794+
createToolCall(2, "never"),
795+
"application/json",
796+
"test-session",
797+
),
798+
);
799+
await flushMicrotasks();
800+
expect(internals._jsonResponseCollectors.size).toBe(1);
801+
802+
const cancelled = await transport.handleRequest(
803+
createPostRequest(
804+
createCancelledNotification(2),
805+
"application/json",
806+
"test-session",
807+
),
808+
);
809+
expect(cancelled.status).toBe(202);
810+
811+
// A cancelled request is never answered, so this POST carries nothing back
812+
// — but it must still end, and leave no routing state behind. Without the
813+
// release it waits on a response that is never coming.
814+
const response = await Promise.race([
815+
pending,
816+
new Promise<null>((resolve) => setTimeout(() => resolve(null), 1000)),
817+
]);
818+
expect(response).not.toBeNull();
819+
expect(response?.status).toBe(202);
820+
expect(await response?.text()).toBe("");
821+
expect(internals._jsonResponseCollectors.size).toBe(0);
822+
expect(internals._requestToStreamMapping.size).toBe(0);
823+
});
824+
825+
it("should answer the rest of a batch when one of its requests is cancelled", async () => {
826+
const transport = new WebStreamableHTTPServerTransport({
827+
enableJsonResponse: true,
828+
sessionIdGenerator: () => "test-session",
829+
});
830+
const server = new Server(
831+
{ name: "TestServer", version: "1.0.0" },
832+
{ capabilities: { tools: {} } },
833+
);
834+
server.setRequestHandler(CallToolRequestSchema, (request, extra) => {
835+
if (request.params.name === "never") {
836+
return new Promise<never>((_resolve, reject) => {
837+
extra.signal.addEventListener("abort", () =>
838+
reject(extra.signal.reason ?? new Error("aborted")),
839+
);
840+
});
841+
}
842+
return Promise.resolve({
843+
content: [{ text: `done:${request.params.name}`, type: "text" }],
844+
});
845+
});
846+
await server.connect(transport);
847+
await transport.handleRequest(createInitializeRequest("application/json"));
848+
849+
const pending = transport.handleRequest(
850+
createPostRequest(
851+
[createToolCall(2, "never"), createToolCall(3, "one")],
852+
"application/json",
853+
"test-session",
854+
),
855+
);
856+
await flushMicrotasks();
857+
858+
await transport.handleRequest(
859+
createPostRequest(
860+
createCancelledNotification(2),
861+
"application/json",
862+
"test-session",
863+
),
864+
);
865+
866+
// Dropping the cancelled id must not take the batch's other answer with it.
867+
const response = await Promise.race([
868+
pending,
869+
new Promise<null>((resolve) => setTimeout(() => resolve(null), 1000)),
870+
]);
871+
expect(response?.status).toBe(200);
872+
const body: JsonResponse = await response?.json();
873+
expect(body).toHaveLength(1);
874+
expect(body[0].id).toBe(3);
875+
expect(body[0].result.content[0].text).toBe("done:one");
876+
});
877+
693878
it("should reject an empty JSON-RPC batch with one Invalid Request object", async () => {
694879
const transport = new WebStreamableHTTPServerTransport({
695880
enableJsonResponse: true,

‎src/edge/WebStreamableHTTPServerTransport.ts‎

Lines changed: 43 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
import { Transport } from "@modelcontextprotocol/sdk/shared/transport.js";
1010
import {
1111
CancelledNotificationSchema,
12+
ErrorCode,
1213
isInitializeRequest,
1314
isJSONRPCError,
1415
isJSONRPCRequest,
@@ -579,6 +580,16 @@ export class WebStreamableHTTPServerTransport implements Transport {
579580
// Wait for the handlers to actually answer, however long they take
580581
const responses = await jsonResponses;
581582

583+
// Every request this POST carried was cancelled, so there is nothing to
584+
// put in the body — and an empty body is not a JSON-RPC message. Answer
585+
// the way a POST that carried no requests at all does.
586+
if (responses.length === 0) {
587+
return new Response(null, {
588+
headers: this.getResponseHeaders(),
589+
status: 202,
590+
});
591+
}
592+
582593
const responseBody = JSON.stringify(isBatch ? responses : responses[0]);
583594

584595
return new Response(responseBody, {
@@ -655,13 +666,24 @@ export class WebStreamableHTTPServerTransport implements Transport {
655666
* therefore never produce a response: the same drain `send()` performs once
656667
* a response has been written, minus the response. Idempotent, so it is a
657668
* no-op if the response won the race against the cancellation.
658-
*
659-
* JSON-mode POSTs are left alone: `_jsonResponseCollectors` owns their
660-
* settlement.
661669
*/
662670
private async releaseCancelledRequest(requestId: RequestId): Promise<void> {
663671
const streamId = this._requestToStreamMapping.get(requestId);
664-
if (streamId === undefined || this._jsonResponseCollectors.has(streamId)) {
672+
if (streamId === undefined) {
673+
return;
674+
}
675+
676+
// A JSON-mode POST has no stream to close: it is parked on its collector,
677+
// which would keep waiting for a response that is never coming. Stop
678+
// expecting that one and answer with whatever the POST's other requests
679+
// produce.
680+
const collector = this._jsonResponseCollectors.get(streamId);
681+
if (collector) {
682+
collector.expectedIds = collector.expectedIds.filter(
683+
(id) => id !== requestId,
684+
);
685+
this._requestToStreamMapping.delete(requestId);
686+
this.settleJsonResponse(streamId, collector);
665687
return;
666688
}
667689

@@ -696,9 +718,23 @@ export class WebStreamableHTTPServerTransport implements Transport {
696718
private releasePendingRequests(): void {
697719
for (const collector of this._jsonResponseCollectors.values()) {
698720
collector.resolve(
699-
collector.expectedIds
700-
.map((id) => collector.responses.get(id))
701-
.filter((response) => response !== undefined),
721+
collector.expectedIds.map(
722+
(id) =>
723+
collector.responses.get(id) ??
724+
// Nobody can answer this one now. Dropping it silently would leave
725+
// the POST with an empty body for a single request and an empty
726+
// array for a batch, neither of which a client can read as a
727+
// result.
728+
({
729+
error: {
730+
code: ErrorCode.ConnectionClosed,
731+
message:
732+
"Connection closed: the session ended before this request was answered",
733+
},
734+
id,
735+
jsonrpc: "2.0",
736+
} satisfies JSONRPCMessage),
737+
),
702738
);
703739
}
704740
this._jsonResponseCollectors.clear();

0 commit comments

Comments
 (0)