From caae3601ee6441bfa6f3c0b45664cb0793974ab5 Mon Sep 17 00:00:00 2001 From: can1357 Date: Thu, 13 Aug 2026 02:09:16 +0200 Subject: [PATCH] Revert "test(coding-agent): stabilize MCP stream resumption" This reverts commit dc283c0deb08bb4063950f9d4b3e67f00a117760. --- .../test/mcp-http-transport.test.ts | 65 +++++++++---------- 1 file changed, 32 insertions(+), 33 deletions(-) diff --git a/packages/coding-agent/test/mcp-http-transport.test.ts b/packages/coding-agent/test/mcp-http-transport.test.ts index e7244dc78..ce7c2a940 100644 --- a/packages/coding-agent/test/mcp-http-transport.test.ts +++ b/packages/coding-agent/test/mcp-http-transport.test.ts @@ -1,4 +1,4 @@ -import { afterEach, describe, expect, it, vi } from "bun:test"; +import { afterEach, describe, expect, it } from "bun:test"; import { HttpTransport } from "@oh-my-pi/pi-coding-agent/mcp/transports/http"; const encoder = new TextEncoder(); @@ -258,41 +258,41 @@ describe("MCP Streamable HTTP POST response resumption", () => { }); it("resumes after an abrupt stream drop once an event ID exists", async () => { - const observed = { posts: 0, gets: 0, lastEventId: null as string | null }; - const fetchMock: typeof globalThis.fetch = Object.assign( - async (_input: string | URL | Request, init?: RequestInit | BunFetchRequestInit): Promise => { - if (init?.method === "POST") { - observed.posts++; - let deliveredEvent = false; - return new Response( - new ReadableStream({ - pull(controller) { - if (!deliveredEvent) { - deliveredEvent = true; - controller.enqueue(encoder.encode("id: stream-1\nretry: 10\ndata:\n\n")); - return; - } - controller.error(new Error("simulated abrupt stream drop")); - }, - }), - { headers: { "Content-Type": "text/event-stream" } }, + // Bun.serve cannot produce a genuine mid-body transport failure in-process + // (stream errors surface as clean EOF client-side), so speak raw HTTP: a + // chunked response without the terminal chunk, closed mid-body, makes the + // client's body read throw. + const observed = { posts: 0, lastEventId: null as string | null }; + const sseChunk = (payload: string): string => + `HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nTransfer-Encoding: chunked\r\n\r\n${payload.length.toString(16)}\r\n${payload}\r\n`; + const listener = Bun.listen({ + hostname: "127.0.0.1", + port: 0, + socket: { + data(socket, data) { + const request = new TextDecoder().decode(data); + if (request.startsWith("POST")) { + observed.posts++; + // Priming event, then close without the terminal 0-chunk. + socket.write(sseChunk("id: stream-1\nretry: 10\ndata:\n\n")); + socket.end(); + return; + } + if (!request.startsWith("GET")) return; + observed.lastEventId = /^Last-Event-ID:\s*(.+)$/im.exec(request)?.[1]?.trim() ?? null; + socket.write( + `${sseChunk( + 'id: stream-2\ndata: {"jsonrpc":"2.0","id":1,"result":{"tools":[{"name":"resumed","inputSchema":{"type":"object"}}]}}\n\n', + )}0\r\n\r\n`, ); - } - if (init?.method !== "GET") throw new Error(`Unexpected method: ${init?.method}`); - observed.gets++; - observed.lastEventId = new Headers(init.headers).get("Last-Event-ID"); - return new Response( - 'id: stream-2\ndata: {"jsonrpc":"2.0","id":1,"result":{"tools":[{"name":"resumed","inputSchema":{"type":"object"}}]}}\n\n', - { headers: { "Content-Type": "text/event-stream" } }, - ); + socket.end(); + }, }, - { preconnect: globalThis.fetch.preconnect }, - ); - const fetchSpy = vi.spyOn(globalThis, "fetch").mockImplementation(fetchMock); + }); try { const transport = new HttpTransport({ type: "http", - url: "http://mcp.test/mcp", + url: `http://127.0.0.1:${listener.port}/mcp`, timeout: GUARD_TIMEOUT_MS, }); await transport.connect(); @@ -301,10 +301,9 @@ describe("MCP Streamable HTTP POST response resumption", () => { tools: [{ name: "resumed", inputSchema: { type: "object" } }], }); expect(observed.posts).toBe(1); - expect(observed.gets).toBe(1); expect(observed.lastEventId).toBe("stream-1"); } finally { - fetchSpy.mockRestore(); + listener.stop(true); } }); });