From 83f77fd90fbb0c39eae080555f32e6b9c295682d Mon Sep 17 00:00:00 2001 From: Christopher Hertel Date: Wed, 7 Oct 2026 22:17:11 +0200 Subject: [PATCH 1/3] [Client] Listen on the standalone GET stream over HTTP A 2025-era server sends requests and notifications that belong to none of the client's requests on the GET stream, which the transport never opened, so they never arrived. The TypeScript SDK's reference server asks for roots that way, and the tool waiting for them timed out. With the new listen option the transport opens the stream after the handshake and reads it alongside the response of the request in flight. Both are read without blocking, since the request that response waits for can be the one arriving on the other stream. A PSR-18 client that buffers whole bodies cannot do that, so the option is off by default. The interop suite's HTTP run now listens, so its roots scenario runs instead of being skipped. --- CHANGELOG.md | 1 + docs/client/transports.md | 20 ++ src/Client/Transport/HttpTransport.php | 196 ++++++++++++++-- .../Client/EverythingServerTestCase.php | 35 +-- .../Client/HttpEverythingServerTest.php | 9 +- .../Transport/HttpTransportListenTest.php | 211 ++++++++++++++++++ 6 files changed, 420 insertions(+), 52 deletions(-) create mode 100644 tests/Unit/Client/Transport/HttpTransportListenTest.php diff --git a/CHANGELOG.md b/CHANGELOG.md index 1bf3d225..ea6d7076 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,6 +24,7 @@ All notable changes to `mcp/sdk` will be documented in this file. * Add `Tool::$execution` (`ToolExecution` with a `TaskSupport` enum) and `ServerCapabilities::$tasks`, which were dropped when parsing a 2025-11-25 server's `tools/list` and `initialize` results. * Add `Client::getServerCapabilities()`, returning what the server declared in `initialize` or, from `2026-07-28` on, in `server/discover`. * [BC Break] `ClientStateInterface` declares `setServerCapabilities()` and `getServerCapabilities()`, which a custom implementation has to add. +* Add a `listen` option to the client's `HttpTransport`, opening the standalone GET stream on which a 2025-era server sends requests and notifications outside of a client request, like `roots/list`. Needs a PSR-18 client that streams response bodies, such as `symfony/http-client`. 0.8.0 ----- diff --git a/docs/client/transports.md b/docs/client/transports.md index 6d0b6a5d..8969b6ee 100644 --- a/docs/client/transports.md +++ b/docs/client/transports.md @@ -48,6 +48,7 @@ $transport = new HttpTransport( - `streamFactory` (StreamFactoryInterface|null): PSR-17 stream factory (auto-discovered) - `logger` (LoggerInterface|null): Optional PSR-3 logger - `maxSseBufferBytes` (int): Maximum buffered bytes for a streamed SSE response +- `listen` (bool): Open the server's listening stream after the handshake (see below) **PSR-18 Auto-Discovery:** @@ -63,6 +64,25 @@ The transport automatically discovers PSR-18 HTTP clients from: composer require php-http/guzzle7-adapter ``` +**Listening for server messages:** + +Most of what a server sends belongs to one of the client's requests and arrives on that request's +response. Up to the 2025 revisions, a server may also send requests and notifications that belong to +none of them, like asking for the client's roots, on a standalone GET stream. Pass `listen: true` to +open it: + +```php +$transport = new HttpTransport('http://localhost:8000', listen: true); +``` + +- The stream is read while one of the client's requests is in flight. A message that arrives while + the client is idle waits for its next request. +- It needs a PSR-18 client that returns before the response body has ended and whose body can be read + without blocking, such as `symfony/http-client`. A client that buffers the whole body never + returns from a stream the server keeps open. +- A server without a listening stream answers `405`, and the connection carries on without one. +- On `2026-07-28`, which has no standalone stream, the option does nothing. + ## Cancellation and deadlines `callTool()` accepts optional `cancellation: ?CancellationTokenInterface` and diff --git a/src/Client/Transport/HttpTransport.php b/src/Client/Transport/HttpTransport.php index c0e455f9..0dde6885 100644 --- a/src/Client/Transport/HttpTransport.php +++ b/src/Client/Transport/HttpTransport.php @@ -61,6 +61,12 @@ class HttpTransport extends BaseTransport implements HeaderAwareTransportInterfa /** @var string Buffer for incomplete SSE data */ private string $sseBuffer = ''; + /** @var StreamInterface|null The standalone GET stream the server may send on unprompted */ + private ?StreamInterface $listenStream = null; + + /** @var string Buffer for incomplete SSE data on the listening stream */ + private string $listenBuffer = ''; + /** * Default cap on the bytes buffered while waiting for a complete SSE event. */ @@ -80,6 +86,16 @@ class HttpTransport extends BaseTransport implements HeaderAwareTransportInterfa * and exhaust client memory; reaching the cap aborts the * stream instead. Raise it for servers that legitimately * emit single events larger than the default. + * @param bool $listen Open the standalone GET stream after the handshake, on + * which a server sends requests and notifications that + * belong to none of the client's requests (2025 revisions + * only). It is read while a request of the client is in + * flight, so a message arriving while the client is idle + * waits for its next request. Needs a PSR-18 client that + * returns before the response body has ended and whose body + * can be read without blocking, such as the Psr18Client of + * symfony/http-client: one that buffers the whole body never + * returns from a stream the server keeps open. */ public function __construct( private readonly string $endpoint, @@ -89,6 +105,7 @@ public function __construct( ?StreamFactoryInterface $streamFactory = null, ?LoggerInterface $logger = null, int $maxSseBufferBytes = self::DEFAULT_MAX_SSE_BUFFER_BYTES, + private readonly bool $listen = false, ) { parent::__construct($logger); @@ -120,6 +137,10 @@ public function connect(): void } $this->logger->info('HTTP client connected and initialized', ['endpoint' => $this->endpoint]); + + if ($this->listen) { + $this->openListenStream(); + } } public function onHeaders(callable $callback): void @@ -190,7 +211,11 @@ public function send(string $data): void $contentType = strtolower($response->getHeaderLine('Content-Type')); if (str_contains($contentType, 'text/event-stream')) { - $this->activeStream = $response->getBody(); + // While listening, a request on the GET stream can be what this + // response waits for, so neither stream may block the other. + $this->activeStream = $this->listen + ? $this->nonBlocking($response->getBody()) + : $response->getBody(); $this->sseBuffer = ''; } elseif (str_contains($contentType, 'application/json')) { $body = $response->getBody()->getContents(); @@ -256,6 +281,9 @@ public function close(): void $this->sessionId = null; $this->activeStream = null; + $this->listenStream?->close(); + $this->listenStream = null; + $this->listenBuffer = ''; $this->handleClose('Transport closed'); } @@ -280,10 +308,103 @@ private function protocolHeaders(string $payload): array } } + /** + * Open the standalone GET stream, if the server offers one. + * + * Failing to is not an error: the stream is optional on both sides, and + * everything that belongs to a request still arrives on its response. + */ + private function openListenStream(): void + { + $version = $this->state?->getProtocolVersion(); + if (null !== $version && $version->isModern()) { + // 2026-07-28 has no standalone stream: a server asks for input + // within the result of the request that needs it. + return; + } + + $request = $this->requestFactory->createRequest('GET', $this->endpoint) + ->withHeader('Accept', 'text/event-stream'); + + if (null !== $this->sessionId) { + $request = $request->withHeader('Mcp-Session-Id', $this->sessionId); + } + if (null !== $version) { + $request = $request->withHeader('MCP-Protocol-Version', $version->value); + } + + foreach ($this->headers as $name => $value) { + $request = $request->withHeader($name, $value); + } + + try { + $response = $this->httpClient->sendRequest($request); + } catch (\Throwable $e) { + $this->logger->warning('Could not open the listening stream', ['exception' => $e]); + + return; + } + + if (405 === $response->getStatusCode()) { + $this->logger->info('Server offers no listening stream'); + + return; + } + + if (200 !== $response->getStatusCode() || !str_contains(strtolower($response->getHeaderLine('Content-Type')), 'text/event-stream')) { + $this->logger->warning('Server answered the listening stream with something else', [ + 'status' => $response->getStatusCode(), + 'content_type' => $response->getHeaderLine('Content-Type'), + ]); + + return; + } + + $stream = $this->nonBlocking($response->getBody(), $switched); + if (!$switched) { + $this->logger->warning('Not listening: the HTTP client returns a response body that cannot be read without blocking'); + + return; + } + + $this->listenStream = $stream; + $this->listenBuffer = ''; + $this->logger->info('Listening for server messages', ['session_id' => $this->sessionId]); + } + + /** + * The same body, made to return what it has instead of waiting for more. + * + * Always returns a stream to read the body with: the given one when it has + * no PHP stream behind it to switch, as detaching it would leave nothing to + * read with. + * + * @param-out bool $switched set to whether reads no longer block + */ + private function nonBlocking(StreamInterface $body, ?bool &$switched = null): StreamInterface + { + $switched = false; + + // A body backed by a PHP stream lists the stream's metadata. + if ([] === ($body->getMetadata() ?? [])) { + return $body; + } + + $resource = $body->detach(); + if (!\is_resource($resource)) { + return $body; + } + + $switched = stream_set_blocking($resource, false); + + return $this->streamFactory->createStreamFromResource($resource); + } + private function tick(): void { $this->checkInterruption(); $this->processSSEStream(); + $this->processListenStream(); $this->processProgress(); $this->checkInterruption(); $this->processFiber(); @@ -320,34 +441,75 @@ private function processSSEStream(): void return; } - if (!$this->activeStream->eof()) { - $chunk = $this->activeStream->read(4096); + $done = $this->pumpSse($this->activeStream, $this->sseBuffer, $this->abortSseStream(...)); + + if ($done) { + $this->sseBuffer = ''; + $this->activeStream = null; + } + } + + /** + * Read what arrived on the listening stream, if one is open. + * + * Its end is no error: the server may close it at any time, and every + * response still arrives on its own request. + */ + private function processListenStream(): void + { + if (null === $this->listenStream) { + return; + } + + $done = $this->pumpSse($this->listenStream, $this->listenBuffer, function (string $reason): void { + $this->logger->warning('Closing the listening stream: '.$reason, ['session_id' => $this->sessionId]); + }); + + if ($done) { + $this->listenBuffer = ''; + $this->listenStream = null; + $this->logger->info('Listening stream ended', ['session_id' => $this->sessionId]); + } + } + + /** + * Read a chunk of an SSE stream and dispatch every event it completes. + * + * @param callable(string $reason): void $onOverflow called when the buffer would exceed its cap + * + * @return bool whether the stream is done: ended, or given up for its size + */ + private function pumpSse(StreamInterface $stream, string &$buffer, callable $onOverflow): bool + { + if (!$stream->eof()) { + $chunk = $stream->read(4096); if ('' !== $chunk) { - if (\strlen($this->sseBuffer) + \strlen($chunk) > $this->maxSseBufferBytes) { - $this->abortSseStream(\sprintf('buffered %d bytes without a complete event, exceeding the %d byte limit', \strlen($this->sseBuffer) + \strlen($chunk), $this->maxSseBufferBytes)); + if (\strlen($buffer) + \strlen($chunk) > $this->maxSseBufferBytes) { + $onOverflow(\sprintf('buffered %d bytes without a complete event, exceeding the %d byte limit', \strlen($buffer) + \strlen($chunk), $this->maxSseBufferBytes)); - return; + return true; } - $this->sseBuffer .= $chunk; + $buffer .= $chunk; } } - while (null !== ($event = $this->extractSSEEvent())) { + while (null !== ($event = $this->extractSSEEvent($buffer))) { if (!empty(trim($event))) { $this->processSSEEvent($event); } } - if ($this->activeStream->eof()) { + if ($stream->eof()) { // The stream ended without a trailing blank line: dispatch what is left. - if (!empty(trim($this->sseBuffer))) { - $this->processSSEEvent($this->sseBuffer); + if (!empty(trim($buffer))) { + $this->processSSEEvent($buffer); } - $this->sseBuffer = ''; - $this->activeStream = null; + return true; } + + return false; } /** @@ -386,13 +548,13 @@ private function abortSseStream(string $reason): void * event is delimited by any pair of those. Servers built on sse-starlette * (the MCP Python SDK) use CRLF. */ - private function extractSSEEvent(): ?string + private function extractSSEEvent(string &$buffer): ?string { $position = null; $length = 0; foreach (["\r\n\r\n", "\n\n", "\r\r"] as $delimiter) { - $found = strpos($this->sseBuffer, $delimiter); + $found = strpos($buffer, $delimiter); if (false !== $found && (null === $position || $found < $position)) { $position = $found; @@ -404,8 +566,8 @@ private function extractSSEEvent(): ?string return null; } - $event = substr($this->sseBuffer, 0, $position); - $this->sseBuffer = substr($this->sseBuffer, $position + $length); + $event = substr($buffer, 0, $position); + $buffer = substr($buffer, $position + $length); return $event; } diff --git a/tests/Interop/Client/EverythingServerTestCase.php b/tests/Interop/Client/EverythingServerTestCase.php index b039569e..9842a5c2 100644 --- a/tests/Interop/Client/EverythingServerTestCase.php +++ b/tests/Interop/Client/EverythingServerTestCase.php @@ -147,14 +147,12 @@ public function __invoke(ListRootsRequest $request): ListRootsResult // The server asks for the client's roots shortly after the handshake, // outside any request. Wait for it here, so it is not recorded as part // of whichever scenario happens to run when it arrives. - if (static::receivesUnrelatedServerRequests()) { - $deadline = microtime(true) + 5; - while ([] === self::$serverRequests && microtime(true) < $deadline) { - usleep(100_000); - self::$client->ping(); - } - self::$serverRequests = []; + $deadline = microtime(true) + 5; + while ([] === self::$serverRequests && microtime(true) < $deadline) { + usleep(100_000); + self::$client->ping(); } + self::$serverRequests = []; } public static function tearDownAfterClass(): void @@ -368,7 +366,7 @@ private function runScenario(string $scenario, Client $client): mixed 'tools_call-long_running_operation' => $client->callTool('trigger-long-running-operation', ['duration' => 1, 'steps' => 2], $onProgress), 'tools_call-sampling' => $client->callTool('trigger-sampling-request', ['prompt' => 'Say hello', 'maxTokens' => 20]), 'tools_call-elicitation' => $client->callTool('trigger-elicitation-request'), - 'tools_call-roots' => $this->callRootsTool($client), + 'tools_call-roots' => $client->callTool('get-roots-list'), 'tools_call-unknown_tool' => $client->callTool('does-not-exist'), 'resources_list' => $client->listResources(), 'resources_templates_list' => $client->listResourceTemplates(), @@ -383,27 +381,6 @@ private function runScenario(string $scenario, Client $client): mixed }; } - /** - * Whether the server can reach the client with a request that is not part - * of a request the client sent. - */ - protected static function receivesUnrelatedServerRequests(): bool - { - return true; - } - - private function callRootsTool(Client $client): mixed - { - // The server asks for roots outside the tool call it is answering. Over - // Streamable HTTP that request travels on the standalone GET stream, - // which the HTTP transport does not open yet. - if (!static::receivesUnrelatedServerRequests()) { - $this->markTestSkipped('The client transport does not receive server requests sent outside of a client request.'); - } - - return $client->callTool('get-roots-list'); - } - private function assertMatchesSnapshot(string $scenario, mixed $data): void { $json = json_encode($data, \JSON_PRETTY_PRINT | \JSON_UNESCAPED_SLASHES | \JSON_UNESCAPED_UNICODE | \JSON_THROW_ON_ERROR); diff --git a/tests/Interop/Client/HttpEverythingServerTest.php b/tests/Interop/Client/HttpEverythingServerTest.php index a3cd5aa4..4ec8483c 100644 --- a/tests/Interop/Client/HttpEverythingServerTest.php +++ b/tests/Interop/Client/HttpEverythingServerTest.php @@ -55,11 +55,8 @@ public static function tearDownAfterClass(): void protected static function transport(): TransportInterface { - return new HttpTransport(self::$endpoint); - } - - protected static function receivesUnrelatedServerRequests(): bool - { - return false; + // The server asks for roots outside of any request, which over HTTP + // arrives on the standalone GET stream. + return new HttpTransport(self::$endpoint, listen: true); } } diff --git a/tests/Unit/Client/Transport/HttpTransportListenTest.php b/tests/Unit/Client/Transport/HttpTransportListenTest.php new file mode 100644 index 00000000..09c758f8 --- /dev/null +++ b/tests/Unit/Client/Transport/HttpTransportListenTest.php @@ -0,0 +1,211 @@ +client($server, listen: true); + + $result = $client->callTool('get-roots-list'); + + $this->assertInstanceOf(TextContent::class, $result->content[0]); + $this->assertSame('roots: file:///workspace/app', $result->content[0]->text); + + $get = $server->requests('GET')[0] ?? null; + $this->assertNotNull($get, 'the client opens the listening stream'); + $this->assertSame('text/event-stream', $get->getHeaderLine('Accept')); + $this->assertSame('session-1', $get->getHeaderLine('Mcp-Session-Id')); + $this->assertSame(ProtocolVersion::V2025_11_25->value, $get->getHeaderLine('MCP-Protocol-Version')); + + $client->disconnect(); + } + + #[TestDox('no listening stream is opened unless asked for')] + public function testNoListenStreamByDefault(): void + { + $server = new FakeListeningServer(); + $client = $this->client($server, listen: false); + + $this->assertSame([], $server->requests('GET')); + + $client->disconnect(); + } + + #[TestDox('a server without a listening stream (405) leaves the connection usable')] + public function testServerWithoutListenStream(): void + { + $server = new FakeListeningServer(getStatus: 405); + $client = $this->client($server, listen: true); + + $this->assertTrue($client->isConnected()); + $this->assertCount(1, $server->requests('GET')); + + $client->disconnect(); + } + + private function client(FakeListeningServer $server, bool $listen): Client + { + $factory = new Psr17Factory(); + + $client = Client::builder() + ->setClientInfo('test-client', '1.0.0') + ->setProtocolVersion(ProtocolVersion::V2025_11_25) + ->setInitTimeout(2) + ->setRequestTimeout(2) + ->setCapabilities(new ClientCapabilities(roots: true)) + ->addRequestHandler(new ListRootsRequestHandler(new class implements RootsCallbackInterface { + public function __invoke(ListRootsRequest $request): ListRootsResult + { + return new ListRootsResult([new Root('file:///workspace/app', 'Application')]); + } + })) + ->build(); + + $client->connect(new HttpTransport('http://localhost/mcp', [], $server, $factory, $factory, listen: $listen)); + + return $client; + } +} + +/** + * Answers `get-roots-list` the way the TypeScript reference server does: it asks + * for the roots on the listening stream, and finishes the tool call on the call's + * own stream once the client has answered. + */ +final class FakeListeningServer implements ClientInterface +{ + /** @var list */ + private array $requests = []; + + /** @var resource|null the server end of the tool call's stream */ + private $callStream; + + private int|string|null $callId = null; + + public function __construct( + private readonly int $getStatus = 200, + ) { + } + + /** + * @return list + */ + public function requests(string $method): array + { + return array_values(array_filter($this->requests, static fn (RequestInterface $r): bool => $method === $r->getMethod())); + } + + public function sendRequest(RequestInterface $request): ResponseInterface + { + $this->requests[] = $request; + + if ('GET' === $request->getMethod()) { + if (200 !== $this->getStatus) { + return new Response($this->getStatus); + } + + [$server, $client] = $this->socketPair(); + fwrite($server, self::event(['jsonrpc' => '2.0', 'id' => 'srv-1', 'method' => 'roots/list'])); + + return new Response(200, ['Content-Type' => 'text/event-stream'], Stream::create($client)); + } + + if ('DELETE' === $request->getMethod()) { + return new Response(200); + } + + $message = json_decode((string) $request->getBody(), true); + + if ('initialize' === ($message['method'] ?? null)) { + return new Response(200, ['Content-Type' => 'application/json', 'Mcp-Session-Id' => 'session-1'], json_encode([ + 'jsonrpc' => '2.0', + 'id' => $message['id'], + 'result' => [ + 'protocolVersion' => '2025-11-25', + 'capabilities' => ['tools' => new \stdClass()], + 'serverInfo' => ['name' => 'fake', 'version' => '1.0.0'], + ], + ])); + } + + if ('tools/call' === ($message['method'] ?? null)) { + // Held open with nothing on it until the roots arrive. + [$this->callStream, $client] = $this->socketPair(); + $this->callId = $message['id']; + + return new Response(200, ['Content-Type' => 'text/event-stream'], Stream::create($client)); + } + + if ('srv-1' === ($message['id'] ?? null) && null !== $this->callStream) { + $uri = $message['result']['roots'][0]['uri'] ?? '?'; + fwrite($this->callStream, self::event([ + 'jsonrpc' => '2.0', + 'id' => $this->callId, + 'result' => ['content' => [['type' => 'text', 'text' => 'roots: '.$uri]]], + ])); + fclose($this->callStream); + $this->callStream = null; + } + + return new Response(202); + } + + /** + * @return array{resource, resource} + */ + private function socketPair(): array + { + $pair = stream_socket_pair(\STREAM_PF_UNIX, \STREAM_SOCK_STREAM, \STREAM_IPPROTO_IP); + if (false === $pair) { + throw new \RuntimeException('Could not create a socket pair.'); + } + + return $pair; + } + + /** + * @param array $message + */ + private static function event(array $message): string + { + return 'data: '.json_encode($message)."\n\n"; + } +} From c414ff4d9f619a039907f156fbb5f5cc3b6c88d1 Mon Sep 17 00:00:00 2001 From: Christopher Hertel Date: Wed, 7 Oct 2026 22:19:11 +0200 Subject: [PATCH 2/3] [Client] Close the listening stream when dropping it, cover its edge cases Only read the POST stream without blocking while a listening stream is actually open, and close bodies that are given up on. --- src/Client/Transport/HttpTransport.php | 33 ++-- .../Transport/HttpTransportListenTest.php | 142 +++++++++++++++++- 2 files changed, 152 insertions(+), 23 deletions(-) diff --git a/src/Client/Transport/HttpTransport.php b/src/Client/Transport/HttpTransport.php index 0dde6885..8fd07be9 100644 --- a/src/Client/Transport/HttpTransport.php +++ b/src/Client/Transport/HttpTransport.php @@ -86,16 +86,9 @@ class HttpTransport extends BaseTransport implements HeaderAwareTransportInterfa * and exhaust client memory; reaching the cap aborts the * stream instead. Raise it for servers that legitimately * emit single events larger than the default. - * @param bool $listen Open the standalone GET stream after the handshake, on - * which a server sends requests and notifications that - * belong to none of the client's requests (2025 revisions - * only). It is read while a request of the client is in - * flight, so a message arriving while the client is idle - * waits for its next request. Needs a PSR-18 client that - * returns before the response body has ended and whose body - * can be read without blocking, such as the Psr18Client of - * symfony/http-client: one that buffers the whole body never - * returns from a stream the server keeps open. + * @param bool $listen Open the standalone GET stream after the handshake (2025 + * revisions only). Needs a PSR-18 client that streams + * response bodies; see docs/client/transports.md. */ public function __construct( private readonly string $endpoint, @@ -139,6 +132,7 @@ public function connect(): void $this->logger->info('HTTP client connected and initialized', ['endpoint' => $this->endpoint]); if ($this->listen) { + $this->closeListenStream(); $this->openListenStream(); } } @@ -213,7 +207,7 @@ public function send(string $data): void if (str_contains($contentType, 'text/event-stream')) { // While listening, a request on the GET stream can be what this // response waits for, so neither stream may block the other. - $this->activeStream = $this->listen + $this->activeStream = null !== $this->listenStream ? $this->nonBlocking($response->getBody()) : $response->getBody(); $this->sseBuffer = ''; @@ -281,9 +275,7 @@ public function close(): void $this->sessionId = null; $this->activeStream = null; - $this->listenStream?->close(); - $this->listenStream = null; - $this->listenBuffer = ''; + $this->closeListenStream(); $this->handleClose('Transport closed'); } @@ -346,12 +338,14 @@ private function openListenStream(): void } if (405 === $response->getStatusCode()) { + $response->getBody()->close(); $this->logger->info('Server offers no listening stream'); return; } if (200 !== $response->getStatusCode() || !str_contains(strtolower($response->getHeaderLine('Content-Type')), 'text/event-stream')) { + $response->getBody()->close(); $this->logger->warning('Server answered the listening stream with something else', [ 'status' => $response->getStatusCode(), 'content_type' => $response->getHeaderLine('Content-Type'), @@ -362,6 +356,7 @@ private function openListenStream(): void $stream = $this->nonBlocking($response->getBody(), $switched); if (!$switched) { + $stream->close(); $this->logger->warning('Not listening: the HTTP client returns a response body that cannot be read without blocking'); return; @@ -372,6 +367,13 @@ private function openListenStream(): void $this->logger->info('Listening for server messages', ['session_id' => $this->sessionId]); } + private function closeListenStream(): void + { + $this->listenStream?->close(); + $this->listenStream = null; + $this->listenBuffer = ''; + } + /** * The same body, made to return what it has instead of waiting for more. * @@ -466,8 +468,7 @@ private function processListenStream(): void }); if ($done) { - $this->listenBuffer = ''; - $this->listenStream = null; + $this->closeListenStream(); $this->logger->info('Listening stream ended', ['session_id' => $this->sessionId]); } } diff --git a/tests/Unit/Client/Transport/HttpTransportListenTest.php b/tests/Unit/Client/Transport/HttpTransportListenTest.php index 09c758f8..779a1e8e 100644 --- a/tests/Unit/Client/Transport/HttpTransportListenTest.php +++ b/tests/Unit/Client/Transport/HttpTransportListenTest.php @@ -24,6 +24,7 @@ use Nyholm\Psr7\Factory\Psr17Factory; use Nyholm\Psr7\Response; use Nyholm\Psr7\Stream; +use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\Attributes\TestDox; use PHPUnit\Framework\TestCase; use Psr\Http\Client\ClientInterface; @@ -81,13 +82,102 @@ public function testServerWithoutListenStream(): void $client->disconnect(); } - private function client(FakeListeningServer $server, bool $listen): Client + #[TestDox('nothing is opened on 2026-07-28, which has no standalone stream')] + public function testNoListenStreamOnModernRevision(): void { + $server = new FakeListeningServer(); + $client = $this->client($server, listen: true, version: ProtocolVersion::V2026_07_28); + + $this->assertTrue($client->isConnected()); + $this->assertSame([], $server->requests('GET')); + + $client->disconnect(); + } + + /** @return iterable */ + public static function unexpectedAnswerProvider(): iterable + { + yield 'unknown session (404)' => [404, 'application/json']; + yield 'JSON instead of a stream' => [200, 'application/json']; + } + + #[DataProvider('unexpectedAnswerProvider')] + #[TestDox('a listening stream answered otherwise leaves the connection usable: $_dataName')] + public function testUnexpectedAnswerToListenStream(int $status, string $contentType): void + { + $server = new FakeListeningServer(getStatus: $status, getContentType: $contentType); + $client = $this->client($server, listen: true); + + $this->assertSame('hello', $this->echo($client, 'hello')); + + $client->disconnect(); + } + + #[TestDox('the GET carries the configured headers')] + public function testListenStreamCarriesConfiguredHeaders(): void + { + $server = new FakeListeningServer(); + $client = $this->client($server, listen: true, headers: ['Authorization' => 'Bearer secret']); + + $this->assertSame('Bearer secret', $server->requests('GET')[0]->getHeaderLine('Authorization')); + + $client->disconnect(); + } + + #[TestDox('a listening stream the server ends leaves the connection usable')] + public function testListenStreamEndedByServer(): void + { + $server = new FakeListeningServer(rootsOnGet: false); + $client = $this->client($server, listen: true); + + fclose($server->getStream()); + + $this->assertSame('first', $this->echo($client, 'first')); + $this->assertSame('second', $this->echo($client, 'second')); + + $client->disconnect(); + } + + #[TestDox('an oversized event on the listening stream closes it, leaving the connection usable')] + public function testOversizedEventOnListenStreamClosesIt(): void + { + $server = new FakeListeningServer(rootsOnGet: false); + $client = $this->client($server, listen: true, maxSseBufferBytes: 256); + + $serverEnd = $server->getStream(); + fwrite($serverEnd, 'data: '.str_repeat('x', 1024)); + + $this->assertSame('hello', $this->echo($client, 'hello')); + + stream_set_blocking($serverEnd, false); + fread($serverEnd, 1); + $this->assertTrue(feof($serverEnd), 'the client closed its end of the listening stream'); + + $client->disconnect(); + } + + private function echo(Client $client, string $text): ?string + { + $content = $client->callTool('echo', ['text' => $text])->content[0] ?? null; + + return $content instanceof TextContent ? $content->text : null; + } + + /** + * @param array $headers + */ + private function client( + FakeListeningServer $server, + bool $listen, + ProtocolVersion $version = ProtocolVersion::V2025_11_25, + array $headers = [], + int $maxSseBufferBytes = 1_048_576, + ): Client { $factory = new Psr17Factory(); $client = Client::builder() ->setClientInfo('test-client', '1.0.0') - ->setProtocolVersion(ProtocolVersion::V2025_11_25) + ->setProtocolVersion($version) ->setInitTimeout(2) ->setRequestTimeout(2) ->setCapabilities(new ClientCapabilities(roots: true)) @@ -99,7 +189,7 @@ public function __invoke(ListRootsRequest $request): ListRootsResult })) ->build(); - $client->connect(new HttpTransport('http://localhost/mcp', [], $server, $factory, $factory, listen: $listen)); + $client->connect(new HttpTransport('http://localhost/mcp', $headers, $server, $factory, $factory, maxSseBufferBytes: $maxSseBufferBytes, listen: $listen)); return $client; } @@ -120,11 +210,24 @@ final class FakeListeningServer implements ClientInterface private int|string|null $callId = null; + /** @var resource|null the server end of the listening stream */ + private $getStream; + public function __construct( private readonly int $getStatus = 200, + private readonly string $getContentType = 'text/event-stream', + private readonly bool $rootsOnGet = true, ) { } + /** + * @return resource the server end of the listening stream + */ + public function getStream() + { + return $this->getStream ?? throw new \RuntimeException('No listening stream was opened.'); + } + /** * @return list */ @@ -138,12 +241,14 @@ public function sendRequest(RequestInterface $request): ResponseInterface $this->requests[] = $request; if ('GET' === $request->getMethod()) { - if (200 !== $this->getStatus) { - return new Response($this->getStatus); + if (200 !== $this->getStatus || 'text/event-stream' !== $this->getContentType) { + return new Response($this->getStatus, ['Content-Type' => $this->getContentType], '{}'); } - [$server, $client] = $this->socketPair(); - fwrite($server, self::event(['jsonrpc' => '2.0', 'id' => 'srv-1', 'method' => 'roots/list'])); + [$this->getStream, $client] = $this->socketPair(); + if ($this->rootsOnGet) { + fwrite($this->getStream, self::event(['jsonrpc' => '2.0', 'id' => 'srv-1', 'method' => 'roots/list'])); + } return new Response(200, ['Content-Type' => 'text/event-stream'], Stream::create($client)); } @@ -166,6 +271,29 @@ public function sendRequest(RequestInterface $request): ResponseInterface ])); } + if ('server/discover' === ($message['method'] ?? null)) { + return new Response(200, ['Content-Type' => 'application/json'], json_encode([ + 'jsonrpc' => '2.0', + 'id' => $message['id'], + 'result' => [ + 'resultType' => 'complete', + 'supportedVersions' => ['2026-07-28'], + 'capabilities' => ['tools' => new \stdClass()], + 'serverInfo' => ['name' => 'fake', 'version' => '1.0.0'], + ], + ])); + } + + // Streamed, so the client reads its streams in the meantime: a JSON + // answer completes the call before the listening stream is looked at. + if ('tools/call' === ($message['method'] ?? null) && 'echo' === ($message['params']['name'] ?? null)) { + return new Response(200, ['Content-Type' => 'text/event-stream'], self::event([ + 'jsonrpc' => '2.0', + 'id' => $message['id'], + 'result' => ['content' => [['type' => 'text', 'text' => $message['params']['arguments']['text'] ?? '']]], + ])); + } + if ('tools/call' === ($message['method'] ?? null)) { // Held open with nothing on it until the roots arrive. [$this->callStream, $client] = $this->socketPair(); From 9e3db7fada44ff93c8ae40c5e8d82d29d2090001 Mon Sep 17 00:00:00 2001 From: Christopher Hertel Date: Wed, 7 Oct 2026 22:36:45 +0200 Subject: [PATCH 3/3] [Tests] Run the HTTP interop suite with and without listening The default transport keeps its blocking read covered against the reference server, a listening variant runs every scenario on the non-blocking path, including roots. Setup now asserts the server's roots request arrived when the transport can receive it. --- .../Client/EverythingServerTestCase.php | 36 ++++++++++-- .../Client/HttpEverythingServerTest.php | 47 +++------------- .../Client/HttpEverythingServerTestCase.php | 56 +++++++++++++++++++ .../HttpListeningEverythingServerTest.php | 28 ++++++++++ 4 files changed, 122 insertions(+), 45 deletions(-) create mode 100644 tests/Interop/Client/HttpEverythingServerTestCase.php create mode 100644 tests/Interop/Client/HttpListeningEverythingServerTest.php diff --git a/tests/Interop/Client/EverythingServerTestCase.php b/tests/Interop/Client/EverythingServerTestCase.php index 9842a5c2..002fa4c0 100644 --- a/tests/Interop/Client/EverythingServerTestCase.php +++ b/tests/Interop/Client/EverythingServerTestCase.php @@ -147,12 +147,15 @@ public function __invoke(ListRootsRequest $request): ListRootsResult // The server asks for the client's roots shortly after the handshake, // outside any request. Wait for it here, so it is not recorded as part // of whichever scenario happens to run when it arrives. - $deadline = microtime(true) + 5; - while ([] === self::$serverRequests && microtime(true) < $deadline) { - usleep(100_000); - self::$client->ping(); + if (static::receivesUnrelatedServerRequests()) { + $deadline = microtime(true) + 5; + while ([] === self::$serverRequests && microtime(true) < $deadline) { + usleep(100_000); + self::$client->ping(); + } + self::assertNotSame([], self::$serverRequests, 'The server did not ask for the client\'s roots after the handshake.'); + self::$serverRequests = []; } - self::$serverRequests = []; } public static function tearDownAfterClass(): void @@ -366,7 +369,7 @@ private function runScenario(string $scenario, Client $client): mixed 'tools_call-long_running_operation' => $client->callTool('trigger-long-running-operation', ['duration' => 1, 'steps' => 2], $onProgress), 'tools_call-sampling' => $client->callTool('trigger-sampling-request', ['prompt' => 'Say hello', 'maxTokens' => 20]), 'tools_call-elicitation' => $client->callTool('trigger-elicitation-request'), - 'tools_call-roots' => $client->callTool('get-roots-list'), + 'tools_call-roots' => $this->callRootsTool($client), 'tools_call-unknown_tool' => $client->callTool('does-not-exist'), 'resources_list' => $client->listResources(), 'resources_templates_list' => $client->listResourceTemplates(), @@ -381,6 +384,27 @@ private function runScenario(string $scenario, Client $client): mixed }; } + /** + * Whether the server can reach the client with a request that is not part + * of a request the client sent. + */ + protected static function receivesUnrelatedServerRequests(): bool + { + return true; + } + + private function callRootsTool(Client $client): mixed + { + // The server asks for roots outside the tool call it is answering. Over + // Streamable HTTP that request travels on the standalone GET stream, + // which the HTTP transport only opens when listening. + if (!static::receivesUnrelatedServerRequests()) { + $this->markTestSkipped('The client transport does not receive server requests sent outside of a client request.'); + } + + return $client->callTool('get-roots-list'); + } + private function assertMatchesSnapshot(string $scenario, mixed $data): void { $json = json_encode($data, \JSON_PRETTY_PRINT | \JSON_UNESCAPED_SLASHES | \JSON_UNESCAPED_UNICODE | \JSON_THROW_ON_ERROR); diff --git a/tests/Interop/Client/HttpEverythingServerTest.php b/tests/Interop/Client/HttpEverythingServerTest.php index 4ec8483c..0f5e2234 100644 --- a/tests/Interop/Client/HttpEverythingServerTest.php +++ b/tests/Interop/Client/HttpEverythingServerTest.php @@ -13,50 +13,19 @@ use Mcp\Client\Transport\HttpTransport; use Mcp\Client\Transport\TransportInterface; -use Mcp\Tests\Interop\NodeTool; -use Mcp\Tests\Support\ServerProcess; -use Symfony\Component\Process\Process; -final class HttpEverythingServerTest extends EverythingServerTestCase +/** + * The default transport, which does not open the standalone GET stream. + */ +final class HttpEverythingServerTest extends HttpEverythingServerTestCase { - private static ?Process $server = null; - private static string $endpoint; - - public static function setUpBeforeClass(): void - { - // Started here rather than in transport() so a server that never comes - // up fails with its own output instead of as a client timeout. - $port = ServerProcess::findFreePort(); - - self::$server = new Process([NodeTool::path('mcp-server-everything'), 'streamableHttp'], env: ['PORT' => (string) $port]); - self::$server->start(); - - $deadline = microtime(true) + 10; - while (!@fsockopen('127.0.0.1', $port, $errno, $error, 0.1)) { - if (!self::$server->isRunning() || microtime(true) > $deadline) { - self::fail(\sprintf('The everything server did not start: %s', self::$server->getErrorOutput())); - } - - usleep(100_000); - } - - self::$endpoint = \sprintf('http://127.0.0.1:%d/mcp', $port); - - parent::setUpBeforeClass(); - } - - public static function tearDownAfterClass(): void + protected static function transport(): TransportInterface { - parent::tearDownAfterClass(); - - ServerProcess::stop(self::$server); - self::$server = null; + return new HttpTransport(self::$endpoint); } - protected static function transport(): TransportInterface + protected static function receivesUnrelatedServerRequests(): bool { - // The server asks for roots outside of any request, which over HTTP - // arrives on the standalone GET stream. - return new HttpTransport(self::$endpoint, listen: true); + return false; } } diff --git a/tests/Interop/Client/HttpEverythingServerTestCase.php b/tests/Interop/Client/HttpEverythingServerTestCase.php new file mode 100644 index 00000000..d0d41684 --- /dev/null +++ b/tests/Interop/Client/HttpEverythingServerTestCase.php @@ -0,0 +1,56 @@ + (string) $port]); + self::$server->start(); + + $deadline = microtime(true) + 10; + while (!@fsockopen('127.0.0.1', $port, $errno, $error, 0.1)) { + if (!self::$server->isRunning() || microtime(true) > $deadline) { + self::fail(\sprintf('The everything server did not start: %s', self::$server->getErrorOutput())); + } + + usleep(100_000); + } + + self::$endpoint = \sprintf('http://127.0.0.1:%d/mcp', $port); + + parent::setUpBeforeClass(); + } + + public static function tearDownAfterClass(): void + { + parent::tearDownAfterClass(); + + ServerProcess::stop(self::$server); + self::$server = null; + } +} diff --git a/tests/Interop/Client/HttpListeningEverythingServerTest.php b/tests/Interop/Client/HttpListeningEverythingServerTest.php new file mode 100644 index 00000000..1cb2baf8 --- /dev/null +++ b/tests/Interop/Client/HttpListeningEverythingServerTest.php @@ -0,0 +1,28 @@ +