Skip to content
Merged
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
-----
Expand Down
20 changes: 20 additions & 0 deletions docs/client/transports.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:**

Expand All @@ -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
Expand Down
197 changes: 180 additions & 17 deletions src/Client/Transport/HttpTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand All @@ -80,6 +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 (2025
* revisions only). Needs a PSR-18 client that streams
* response bodies; see docs/client/transports.md.
*/
public function __construct(
private readonly string $endpoint,
Expand All @@ -89,6 +98,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);

Expand Down Expand Up @@ -120,6 +130,11 @@ public function connect(): void
}

$this->logger->info('HTTP client connected and initialized', ['endpoint' => $this->endpoint]);

if ($this->listen) {
$this->closeListenStream();
$this->openListenStream();
}
}

public function onHeaders(callable $callback): void
Expand Down Expand Up @@ -190,7 +205,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 = null !== $this->listenStream
? $this->nonBlocking($response->getBody())
: $response->getBody();
$this->sseBuffer = '';
} elseif (str_contains($contentType, 'application/json')) {
$body = $response->getBody()->getContents();
Expand Down Expand Up @@ -256,6 +275,7 @@ public function close(): void

$this->sessionId = null;
$this->activeStream = null;
$this->closeListenStream();
$this->handleClose('Transport closed');
}

Expand All @@ -280,10 +300,113 @@ 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()) {
$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'),
]);

return;
}

$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;
}

$this->listenStream = $stream;
$this->listenBuffer = '';
$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.
*
* 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();
Expand Down Expand Up @@ -320,34 +443,74 @@ 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->closeListenStream();
$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;
}

/**
Expand Down Expand Up @@ -386,13 +549,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;
Expand All @@ -404,8 +567,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;
}
Expand Down
3 changes: 2 additions & 1 deletion tests/Interop/Client/EverythingServerTestCase.php
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,7 @@ public function __invoke(ListRootsRequest $request): ListRootsResult
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 = [];
}
}
Expand Down Expand Up @@ -396,7 +397,7 @@ 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.
// 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.');
}
Expand Down
Loading
Loading