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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ All notable changes to `mcp/sdk` will be documented in this file.
* Fix stateless SSE streams holding back frames until close when PHP output buffering is enabled.
* Reject a recognized `Mcp-Param-*` header whose mirrored argument is absent from the body with `-32020`, instead of accepting the request (SEP-2243).
* Fix `RequestEvent`, `ResponseEvent` and `ErrorEvent` not being dispatched for `2026-07-28` requests.
* Fix the client's `HttpTransport` waiting out the request timeout when a response stream or JSON body ends, is cut or fails to read without the answer: the call now fails at once with a `ConnectionException`.
* [BC Break] Validate a tool result's `structuredContent` against the tool's `outputSchema`, which the specification requires the server to honour. A mismatch is answered with a `CallToolResult` carrying `isError: true` instead of the non-conforming value, matching the TypeScript, Python and Java SDKs. Skipped when the tool declares no `outputSchema`, when the result carries no `structuredContent`, and when the result is already an error. Return `new \stdClass()` for an empty object, since `[]` is sent as an array.
* Stop the server `Protocol` from logging full JSON-RPC payloads (tool arguments, client replies) at info level: info records now carry only the method and id, the raw message is logged at debug level.
* Add `PassthroughMiddleware` to opt `StreamableHttpTransport` out of its default middleware without the warning an empty `$middleware` list logs.
Expand Down
94 changes: 92 additions & 2 deletions src/Client/Transport/HttpTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,12 @@ class HttpTransport extends BaseTransport implements HeaderAwareTransportInterfa
/** @var string Buffer for incomplete SSE data */
private string $sseBuffer = '';

/** The request whose POST response opened the active SSE stream. */
private int|string|null $sseRequestId = null;

/** Whether the active SSE stream delivered the answer to that request. */
private bool $sseRequestCompleted = false;

/** @var StreamInterface|null The standalone GET stream the server may send on unprompted */
private ?StreamInterface $listenStream = null;

Expand Down Expand Up @@ -215,11 +221,24 @@ public function send(string $data): void
? $this->nonBlocking($response->getBody())
: $response->getBody();
$this->sseBuffer = '';
$this->sseRequestId = self::requestId($data);
$this->sseRequestCompleted = false;
} elseif (str_contains($contentType, 'application/json')) {
$body = $response->getBody()->getContents();
try {
$body = $response->getBody()->getContents();
} catch (\Throwable $e) {
throw self::unanswered($e);
}

if (!empty($body)) {
$this->handleMessage($body);
}

// A cut or truncated body would otherwise leave the request waiting
// for an answer that can no longer arrive.
if (null !== ($requestId = self::requestId($data)) && !self::answers($body, $requestId)) {
throw self::unanswered();
}
}
}

Expand All @@ -234,6 +253,42 @@ private static function isNotification(string $data): bool
return \is_array($payload) && \array_key_exists('method', $payload) && !\array_key_exists('id', $payload);
}

/**
* The id of the outgoing message if it is a request, so an answer is owed.
*/
private static function requestId(string $data): int|string|null
{
$payload = json_decode($data, true);
$id = \is_array($payload) && \array_key_exists('method', $payload) ? ($payload['id'] ?? null) : null;

return \is_int($id) || \is_string($id) ? $id : null;
}

/**
* Whether the received message, or a message of the received batch, answers the request.
*/
private static function answers(string $data, int|string $requestId): bool
{
$decoded = json_decode($data, true);
if (!\is_array($decoded)) {
return false;
}

foreach (array_is_list($decoded) ? $decoded : [$decoded] as $message) {
if (\is_array($message) && !\array_key_exists('method', $message) && ($message['id'] ?? null) === $requestId
&& (\array_key_exists('result', $message) || \array_key_exists('error', $message))) {
return true;
}
}

return false;
}

private static function unanswered(?\Throwable $previous = null): ConnectionException
{
return new ConnectionException(\sprintf('The response ended without answering the request%s.', null !== $previous ? ': '.$previous->getMessage() : ''), 0, $previous);
}

/**
* Fails a request refused at the HTTP level at once, rather than at its timeout.
*/
Expand Down Expand Up @@ -296,6 +351,7 @@ public function runRequest(\Fiber $fiber, ?callable $onProgress = null): Respons
$this->activeStream?->close();
$this->activeStream = null;
$this->sseBuffer = '';
$this->sseRequestId = null;
}
}

Expand Down Expand Up @@ -486,14 +542,43 @@ private function processSSEStream(): void
return;
}

$done = $this->pumpSse($this->activeStream, $this->sseBuffer, $this->abortSseStream(...));
try {
$done = $this->pumpSse($this->activeStream, $this->sseBuffer, $this->abortSseStream(...));
} catch (\Throwable $e) {
$this->sseBuffer = '';
$this->activeStream = null;
$this->failUnansweredRequest($e);

return;
}

if ($done) {
$this->sseBuffer = '';
$this->activeStream = null;
$this->failUnansweredRequest();
}
}

/**
* Fail the request the ended SSE stream was opened for, unless it was answered.
*
* Nothing else can deliver the answer once its stream is gone, so the
* waiting fiber fails now instead of spinning until the request timeout.
*/
private function failUnansweredRequest(?\Throwable $previous = null): void
{
$requestId = $this->sseRequestId;
$this->sseRequestId = null;

if (null === $requestId || $this->sseRequestCompleted || !isset($this->state?->getPendingRequests()[$requestId])
|| !$this->activeFiber?->isSuspended()) {
return;
}

$this->logger->warning('Response stream ended without answering the request', ['request_id' => $requestId, 'exception' => $previous]);
$this->activeSuspend = $this->activeFiber->throw(self::unanswered($previous));
}

/**
* Read what arrived on the listening stream, if one is open.
*
Expand Down Expand Up @@ -567,6 +652,7 @@ private function abortSseStream(string $reason): void
$bufferedBytes = \strlen($this->sseBuffer);
$this->sseBuffer = '';
$this->activeStream = null;
$this->sseRequestId = null;

$this->logger->warning('Aborting SSE stream: '.$reason, [
'session_id' => $this->sessionId,
Expand Down Expand Up @@ -630,6 +716,10 @@ private function processSSEEvent(string $event): void
}

if (!empty($data)) {
if (null !== $this->sseRequestId && self::answers($data, $this->sseRequestId)) {
$this->sseRequestCompleted = true;
}

$this->handleMessage($data);
}
}
Expand Down
52 changes: 52 additions & 0 deletions tests/Unit/Client/Transport/HttpTransportTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -567,6 +567,58 @@ public function isCancellationRequested(): bool
$client->disconnect();
}

/**
* A null body is one whose read fails, like a connection reset.
*
* @return iterable<string, array{string, ?string}>
*/
public static function cutResponseProvider(): iterable
{
yield 'an SSE stream ending before the answer' => ['text/event-stream', 'data: {"jsonrpc":"2.0","method":"notifications/message","params":{"level":"info","data":"working"}}'."\n\n"];
yield 'an SSE stream cut mid-event' => ['text/event-stream', 'data: {"jsonrpc":"2.0","id":%d,"result":{"cont'];
yield 'an SSE stream failing to read' => ['text/event-stream', null];
yield 'a truncated JSON body' => ['application/json', '{"jsonrpc":"2.0","id":%d,"result":{"cont'];
yield 'an empty JSON body' => ['application/json', ''];
yield 'a JSON body failing to read' => ['application/json', null];
}

#[DataProvider('cutResponseProvider')]
#[TestDox('a response that ends without the answer fails the call at once: $_dataName')]
public function testCutResponseFailsTheCallAtOnce(string $contentType, ?string $body): void
{
$httpClient = new RecordingHttpClient(function (array $message) use ($contentType, $body): ?ResponseInterface {
if ('cut' !== ($message['params']['name'] ?? null)) {
return null;
}

if (null !== $body) {
return new Response(200, ['Content-Type' => $contentType], \sprintf($body, $message['id']));
}

$stream = $this->createMock(StreamInterface::class);
$stream->method('eof')->willReturn(false);
$stream->method('read')->willThrowException(new \RuntimeException('Connection reset by peer'));
$stream->method('getContents')->willThrowException(new \RuntimeException('Connection reset by peer'));

return new Response(200, ['Content-Type' => $contentType], $stream);
});

$client = Client::builder()->setClientInfo('test', '1')->setProtocolVersion(ProtocolVersion::V2025_11_25)->setRequestTimeout(30)->build();
$client->connect(new HttpTransport('http://localhost/mcp', [], $httpClient, $this->factory, $this->factory));

$started = microtime(true);
try {
$client->callTool('cut');
$this->fail('Expected the call to fail.');
} catch (ConnectionException $e) {
$this->assertStringContainsString('without answering the request', $e->getMessage());
}

$this->assertLessThan(1, microtime(true) - $started, 'the cut must not be waited out like a slow answer');
$this->assertSame('next call', $client->callTool('fast')->content[0]->text ?? null);
$client->disconnect();
}

/** @return iterable<string, array{ProtocolVersion, bool}> */
public static function revisionProvider(): iterable
{
Expand Down
Loading