From 6dcd6204b4c6ded15fcf0f7103a305e533155e8c Mon Sep 17 00:00:00 2001 From: Guillaume Sainthillier Date: Tue, 6 Oct 2026 09:24:55 +0200 Subject: [PATCH 1/5] [Server] Fix lost responses on concurrent requests of the same session (Streamable HTTP) Fixes #275 and #467. A response to a POST went through the session's outgoing queue, which concurrent requests of the same session read and write back whole. One request could then answer with another's response, and the other with an empty 202. StreamableHttpTransport now implements InlineResponseTransportInterface: Protocol hands it the responses through send() and the transport answers the POST with them. The queue keeps server-initiated requests and notifications. --- CHANGELOG.md | 1 + src/Server/Protocol.php | 18 +++- .../InlineResponseTransportInterface.php | 27 ++++++ .../Transport/StreamableHttpTransport.php | 32 ++++++- src/Server/Transport/TransportInterface.php | 3 +- .../Fixture/InterleavingSessionStore.php | 68 +++++++++++++ .../Transport/StreamableHttpTransportTest.php | 95 +++++++++++++++++++ 7 files changed, 235 insertions(+), 9 deletions(-) create mode 100644 src/Server/Transport/InlineResponseTransportInterface.php create mode 100644 tests/Unit/Server/Transport/Fixture/InterleavingSessionStore.php diff --git a/CHANGELOG.md b/CHANGELOG.md index 458177c8f..e8cf8f9b3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -35,6 +35,7 @@ All notable changes to `mcp/sdk` will be documented in this file. * [BC Break] Add `ScopePolicy` as third argument of `AuthorizationMiddleware`, answering `403 insufficient_scope` per method and tool, with scope hierarchies; the `resource_metadata` challenge URL comes from the configured resource instead of the `Host` header. * Expose `WWW-Authenticate` in the default `CorsMiddleware`. * [BC Break] Fix concurrent Streamable HTTP streams on one session resuming each other's fibers: each stream now polls only the client request its own fiber sent, so an elicitation answer reaches the tool call that asked for it. `Protocol::handleFiberYield()` returns the ID of the request it sent. +* Fix lost responses on concurrent requests of one session over Streamable HTTP: a POST is answered with its own responses instead of taking them from the session's outgoing queue. Adds `InlineResponseTransportInterface` for transports that answer each request on the exchange that carried it. 0.8.0 ----- diff --git a/src/Server/Protocol.php b/src/Server/Protocol.php index e31afcc0e..4af596932 100644 --- a/src/Server/Protocol.php +++ b/src/Server/Protocol.php @@ -32,6 +32,7 @@ use Mcp\Server\Stateless\RequestStateCodec; use Mcp\Server\Suspension\NotificationSuspension; use Mcp\Server\Suspension\RequestSuspension; +use Mcp\Server\Transport\InlineResponseTransportInterface; use Mcp\Server\Transport\TransportInterface; use Psr\EventDispatcher\EventDispatcherInterface; use Psr\Log\LoggerInterface; @@ -488,7 +489,10 @@ public function sendNotification(Notification $notification, SessionInterface $s */ private function sendResponse(TransportInterface $transport, Response|Error $response, ?SessionInterface $session, array $context = []): void { - if (null === $session) { + // Queued in the session, a response can be overwritten or taken by a concurrent + // request of the same session: a transport that can answer on the request's + // own exchange gets it directly. + if (null === $session || $transport instanceof InlineResponseTransportInterface) { $this->logger->info('Sending immediate response', [ 'response_id' => $response->getId(), ]); @@ -511,6 +515,10 @@ private function sendResponse(TransportInterface $transport, Response|Error $res } $context['type'] = 'response'; + if (null !== $session) { + $context['session_id'] = $session->getId(); + } + $transport->send($encoded, $context); } else { $this->logger->info('Queueing server response', [ @@ -556,8 +564,12 @@ public function consumeOutgoingMessages(Uuid $sessionId): array { $session = $this->sessionManager->createWithId($sessionId); $queue = $session->get(self::SESSION_OUTGOING_QUEUE, []); - $session->set(self::SESSION_OUTGOING_QUEUE, []); - $session->save(); + + // Saving an unchanged session would only overwrite what a concurrent request saved in the meantime. + if ([] !== $queue) { + $session->set(self::SESSION_OUTGOING_QUEUE, []); + $session->save(); + } return $queue; } diff --git a/src/Server/Transport/InlineResponseTransportInterface.php b/src/Server/Transport/InlineResponseTransportInterface.php new file mode 100644 index 000000000..ec4c2277a --- /dev/null +++ b/src/Server/Transport/InlineResponseTransportInterface.php @@ -0,0 +1,27 @@ + */ -class StreamableHttpTransport extends BaseTransport implements StatelessAwareTransportInterface +class StreamableHttpTransport extends BaseTransport implements StatelessAwareTransportInterface, InlineResponseTransportInterface { use ReadsBoundedBody; @@ -77,6 +77,9 @@ class StreamableHttpTransport extends BaseTransport implements StatelessAwareTra private ?string $immediateResponse = null; private ?int $immediateStatusCode = null; + /** @var list responses to the requests of the current POST, see {@see InlineResponseTransportInterface} */ + private array $inlineResponses = []; + /** @var list|null null until {@see self::listen()} resolves the defaults */ private ?array $middleware; @@ -170,6 +173,12 @@ public function connectStateless(StatelessProtocol $protocol): void public function send(string $data, array $context): void { + if (isset($context['session_id'])) { + $this->inlineResponses[] = $data; + + return; + } + $this->immediateResponse = $data; $this->immediateStatusCode = $context['status_code'] ?? 200; } @@ -205,6 +214,8 @@ protected function handlePostRequest(string $body, ?AccessToken $accessToken = n $this->immediateStatusCode = null; if (null !== $immediateResponse) { + $this->inlineResponses = []; + return $this->responseFactory->createResponse($immediateStatusCode ?? 200) ->withHeader('Content-Type', 'application/json') ->withBody($this->streamFactory->createStream($immediateResponse)); @@ -232,14 +243,14 @@ protected function handleDeleteRequest(): ResponseInterface protected function createJsonResponse(): ResponseInterface { - $outgoingMessages = $this->getOutgoingMessages($this->sessionId); + $messages = [...array_column($this->getOutgoingMessages($this->sessionId), 'message'), ...$this->inlineResponses]; + $this->inlineResponses = []; - if (empty($outgoingMessages)) { + if ([] === $messages) { return $this->responseFactory->createResponse(202) ->withHeader('Content-Type', 'application/json'); } - $messages = array_column($outgoingMessages, 'message'); $responseBody = 1 === \count($messages) ? $messages[0] : '['.implode(',', $messages).']'; $response = $this->responseFactory->createResponse(200) @@ -257,7 +268,11 @@ protected function createStreamedResponse(): ResponseInterface { $fiber = $this->sessionFiber; - $callback = function () use ($fiber): void { + // The other requests of a batch whose handler did not suspend. + $inlineResponses = $this->inlineResponses; + $this->inlineResponses = []; + + $callback = function () use ($fiber, $inlineResponses): void { if (null === $fiber) { return; } @@ -265,6 +280,13 @@ protected function createStreamedResponse(): ResponseInterface try { $this->logger->info('SSE: Starting request processing loop'); + foreach ($inlineResponses as $message) { + echo "event: message\n"; + echo "data: {$message}\n\n"; + @ob_flush(); + flush(); + } + while ($fiber->isSuspended()) { $this->flushOutgoingMessages($this->sessionId); diff --git a/src/Server/Transport/TransportInterface.php b/src/Server/Transport/TransportInterface.php index d35fd3e25..00d8b41f1 100644 --- a/src/Server/Transport/TransportInterface.php +++ b/src/Server/Transport/TransportInterface.php @@ -48,7 +48,8 @@ public function listen(): mixed; /** * Send a message to the client immediately (bypassing session queue). * - * Used for session resolution errors when no session is available. + * Used for session resolution errors when no session is available, and for + * every response on a {@see InlineResponseTransportInterface}. * The transport decides HOW to send based on context. * * @param array $context Context about this message: diff --git a/tests/Unit/Server/Transport/Fixture/InterleavingSessionStore.php b/tests/Unit/Server/Transport/Fixture/InterleavingSessionStore.php new file mode 100644 index 000000000..53a6a246c --- /dev/null +++ b/tests/Unit/Server/Transport/Fixture/InterleavingSessionStore.php @@ -0,0 +1,68 @@ +interleaved = $interleaved; + $this->readBeforeWrite = $readBeforeWrite; + } + + public function read(Uuid $id): string|false + { + if (null !== $data = $this->staleRead) { + $this->staleRead = null; + + return $data; + } + + return parent::read($id); + } + + public function write(Uuid $id, string $data): bool + { + $before = parent::read($id); + $written = parent::write($id, $data); + + if (null !== $interleaved = $this->interleaved) { + $this->interleaved = null; + if ($this->readBeforeWrite) { + $this->staleRead = $before; + } + + $interleaved(); + } + + return $written; + } +} diff --git a/tests/Unit/Server/Transport/StreamableHttpTransportTest.php b/tests/Unit/Server/Transport/StreamableHttpTransportTest.php index 9905be5e8..50af13c7c 100644 --- a/tests/Unit/Server/Transport/StreamableHttpTransportTest.php +++ b/tests/Unit/Server/Transport/StreamableHttpTransportTest.php @@ -13,13 +13,17 @@ use Mcp\Exception\InvalidArgumentException; use Mcp\Schema\JsonRpc\Error; +use Mcp\Server; +use Mcp\Server\RequestContext; use Mcp\Server\Transport\Http\Middleware\CorsMiddleware; use Mcp\Server\Transport\Http\Middleware\DnsRebindingProtectionMiddleware; use Mcp\Server\Transport\Http\Middleware\PassthroughMiddleware; use Mcp\Server\Transport\Http\Middleware\ProtocolVersionMiddleware; use Mcp\Server\Transport\StreamableHttpTransport; use Mcp\Server\Transport\TransportInterface; +use Mcp\Tests\Unit\Server\Transport\Fixture\InterleavingSessionStore; use Nyholm\Psr7\Factory\Psr17Factory; +use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\Attributes\TestDox; use PHPUnit\Framework\TestCase; use Psr\Clock\ClockInterface; @@ -463,6 +467,97 @@ public function now(): \DateTimeImmutable $this->assertInstanceOf(Error::class, $received); } + /** + * @return iterable + */ + public static function provideInterleavings(): iterable + { + yield 'B runs between A saving its session and A answering' => [false]; + yield 'B loaded the session before A saved it (lost update)' => [true]; + } + + #[TestDox('concurrent POSTs of one session each get their own response: $_dataName')] + #[DataProvider('provideInterleavings')] + public function testConcurrentPostsOfOneSessionEachGetTheirOwnResponse(bool $readBeforeWrite): void + { + $store = new InterleavingSessionStore(); + $sessionId = $this->post($store, '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"test","version":"1.0"}}}') + ->getHeaderLine(StreamableHttpTransport::SESSION_HEADER); + + $responseB = null; + $store->interleaveOnNextWrite(function () use ($store, $sessionId, &$responseB): void { + $responseB = $this->post($store, '{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"echo","arguments":{"text":"b"}}}', $sessionId); + }, $readBeforeWrite); + + $responseA = $this->post($store, '{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"echo","arguments":{"text":"a"}}}', $sessionId); + + $this->assertInstanceOf(ResponseInterface::class, $responseB); + foreach ([2 => $responseA, 3 => $responseB] as $id => $response) { + $this->assertSame(200, $response->getStatusCode(), \sprintf('Request %d was answered %d.', $id, $response->getStatusCode())); + $this->assertSame($sessionId, $response->getHeaderLine(StreamableHttpTransport::SESSION_HEADER)); + $this->assertSame($id, json_decode((string) $response->getBody(), true)['id'] ?? null, \sprintf('Request %d got: %s', $id, $response->getBody())); + } + } + + #[TestDox('a batch streamed over SSE still carries the responses that did not suspend')] + public function testStreamedBatchCarriesInlineResponses(): void + { + $store = new InterleavingSessionStore(); + $sessionId = $this->post($store, '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"test","version":"1.0"}}}') + ->getHeaderLine(StreamableHttpTransport::SESSION_HEADER); + + $response = $this->post($store, '[{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"progress","arguments":{},"_meta":{"progressToken":"p"}}},{"jsonrpc":"2.0","id":3,"method":"ping"}]', $sessionId); + + $this->assertSame('text/event-stream', $response->getHeaderLine('Content-Type')); + + // The stream calls ob_flush() itself, so the output is captured by a handler, not a plain buffer. + $output = ''; + ob_start(static function (string $chunk) use (&$output): string { + $output .= $chunk; + + return ''; + }); + try { + $response->getBody()->getContents(); + } finally { + ob_end_flush(); + } + + $this->assertMatchesRegularExpression('/"id":3,"result".*"progressToken":"p".*"id":2,"result"/s', $output); + } + + /** + * Sends one POST to a fresh server sharing $store, like a PHP worker would. + */ + private function post(InterleavingSessionStore $store, string $body, string $sessionId = ''): ResponseInterface + { + $request = $this->factory + ->createServerRequest('POST', 'http://localhost/') + ->withHeader('Host', 'localhost') + ->withHeader('Content-Type', 'application/json') + ->withHeader('Accept', 'application/json, text/event-stream') + ->withBody($this->factory->createStream($body)); + + if ('' !== $sessionId) { + $request = $request + ->withHeader(StreamableHttpTransport::SESSION_HEADER, $sessionId) + ->withHeader(StreamableHttpTransport::PROTOCOL_VERSION_HEADER, '2025-06-18'); + } + + $server = Server::builder() + ->setServerInfo('test', '1.0') + ->setSession($store) + ->addTool(static fn (string $text): string => $text, 'echo') + ->addTool(static function (RequestContext $context): string { + $context->getClientGateway()->progress(0.5); + + return 'done'; + }, 'progress') + ->build(); + + return $server->run(new StreamableHttpTransport($request, $this->factory, $this->factory)); + } + private function stubAuth401(): MiddlewareInterface { return new class($this->factory) implements MiddlewareInterface { From 9a0ef2ec36fa71a9e1db583f0e15c42782e0c4fa Mon Sep 17 00:00:00 2001 From: Guillaume Sainthillier Date: Thu, 8 Oct 2026 02:30:31 +0200 Subject: [PATCH 2/5] [Server] Cover the polling race fixed by skipping empty-queue saves Ported from #545: one worker polls the session for a client's answer while another stores it. Saving the session on every poll used to overwrite that answer. The store fixture moves to Session/Fixture and gains a hook after the next read. Co-authored-by: Christopher Hertel --- tests/Unit/Server/ProtocolSessionRaceTest.php | 54 +++++++++++++++++++ .../Fixture/InterleavingSessionStore.php | 30 ++++++++--- .../Transport/StreamableHttpTransportTest.php | 2 +- 3 files changed, 78 insertions(+), 8 deletions(-) create mode 100644 tests/Unit/Server/ProtocolSessionRaceTest.php rename tests/Unit/Server/{Transport => Session}/Fixture/InterleavingSessionStore.php (65%) diff --git a/tests/Unit/Server/ProtocolSessionRaceTest.php b/tests/Unit/Server/ProtocolSessionRaceTest.php new file mode 100644 index 000000000..73f9e2683 --- /dev/null +++ b/tests/Unit/Server/ProtocolSessionRaceTest.php @@ -0,0 +1,54 @@ +createWithId($sessionId)->save(); + + $waiting = new Protocol([], [], MessageFactory::make(), $sessions); + $answering = new Protocol([], [], MessageFactory::make(), $sessions); + $transport = $this->createMock(TransportInterface::class); + + // The answer lands right after the waiting worker read the session, + // before anything it does next could write the session back. + $store->interleaveAfterNextRead(static function () use ($answering, $transport, $sessionId): void { + $answering->processInput($transport, '{"jsonrpc": "2.0", "id": 7, "result": {"ok": true}}', $sessionId); + }); + + // One turn of the waiting worker's loop, with nothing queued to send. + $this->assertSame([], $waiting->consumeOutgoingMessages($sessionId)); + + $this->assertInstanceOf(Response::class, $waiting->checkResponse(7, $sessionId)); + } +} diff --git a/tests/Unit/Server/Transport/Fixture/InterleavingSessionStore.php b/tests/Unit/Server/Session/Fixture/InterleavingSessionStore.php similarity index 65% rename from tests/Unit/Server/Transport/Fixture/InterleavingSessionStore.php rename to tests/Unit/Server/Session/Fixture/InterleavingSessionStore.php index 53a6a246c..e464ed44b 100644 --- a/tests/Unit/Server/Transport/Fixture/InterleavingSessionStore.php +++ b/tests/Unit/Server/Session/Fixture/InterleavingSessionStore.php @@ -9,23 +9,32 @@ * file that was distributed with this source code. */ -namespace Mcp\Tests\Unit\Server\Transport\Fixture; +namespace Mcp\Tests\Unit\Server\Session\Fixture; use Mcp\Server\Session\InMemorySessionStore; use Symfony\Component\Uid\Uuid; /** - * A session store that runs a second request while the first one saves its session. + * A session store that runs another request in the middle of one reading or saving its session. * * Replays, in one process and in a fixed order, what two PHP workers serving the * same session do when their requests overlap. */ final class InterleavingSessionStore extends InMemorySessionStore { - private ?\Closure $interleaved = null; + private ?\Closure $afterNextRead = null; + private ?\Closure $afterNextWrite = null; private bool $readBeforeWrite = false; private string|false|null $staleRead = null; + /** + * Runs $interleaved right after the next read, before the reader can write the session back. + */ + public function interleaveAfterNextRead(\Closure $interleaved): void + { + $this->afterNextRead = $interleaved; + } + /** * Runs $interleaved right after the next write. * @@ -34,7 +43,7 @@ final class InterleavingSessionStore extends InMemorySessionStore */ public function interleaveOnNextWrite(\Closure $interleaved, bool $readBeforeWrite = false): void { - $this->interleaved = $interleaved; + $this->afterNextWrite = $interleaved; $this->readBeforeWrite = $readBeforeWrite; } @@ -46,7 +55,14 @@ public function read(Uuid $id): string|false return $data; } - return parent::read($id); + $data = parent::read($id); + + if (null !== $interleaved = $this->afterNextRead) { + $this->afterNextRead = null; + $interleaved(); + } + + return $data; } public function write(Uuid $id, string $data): bool @@ -54,8 +70,8 @@ public function write(Uuid $id, string $data): bool $before = parent::read($id); $written = parent::write($id, $data); - if (null !== $interleaved = $this->interleaved) { - $this->interleaved = null; + if (null !== $interleaved = $this->afterNextWrite) { + $this->afterNextWrite = null; if ($this->readBeforeWrite) { $this->staleRead = $before; } diff --git a/tests/Unit/Server/Transport/StreamableHttpTransportTest.php b/tests/Unit/Server/Transport/StreamableHttpTransportTest.php index 50af13c7c..9716b3357 100644 --- a/tests/Unit/Server/Transport/StreamableHttpTransportTest.php +++ b/tests/Unit/Server/Transport/StreamableHttpTransportTest.php @@ -21,7 +21,7 @@ use Mcp\Server\Transport\Http\Middleware\ProtocolVersionMiddleware; use Mcp\Server\Transport\StreamableHttpTransport; use Mcp\Server\Transport\TransportInterface; -use Mcp\Tests\Unit\Server\Transport\Fixture\InterleavingSessionStore; +use Mcp\Tests\Unit\Server\Session\Fixture\InterleavingSessionStore; use Nyholm\Psr7\Factory\Psr17Factory; use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\Attributes\TestDox; From 023887357c1651e512cd6bac765a140cc0208d59 Mon Sep 17 00:00:00 2001 From: Christopher Hertel Date: Fri, 9 Oct 2026 23:10:28 +0200 Subject: [PATCH 3/5] [Server] Cover inline responses in ProtocolTest --- tests/Unit/Server/ProtocolTest.php | 44 ++++++++++++++++++++++++++++++ 1 file changed, 44 insertions(+) diff --git a/tests/Unit/Server/ProtocolTest.php b/tests/Unit/Server/ProtocolTest.php index 54dd81ad2..0a13b9e4e 100644 --- a/tests/Unit/Server/ProtocolTest.php +++ b/tests/Unit/Server/ProtocolTest.php @@ -34,6 +34,8 @@ use Mcp\Server\Session\SessionManagerInterface; use Mcp\Server\Suspension\NotificationSuspension; use Mcp\Server\Suspension\RequestSuspension; +use Mcp\Server\Transport\InlineResponseTransportInterface; +use Mcp\Server\Transport\InMemoryTransport; use Mcp\Server\Transport\TransportInterface; use Mcp\Tests\Unit\Fixtures\PollingLoopTransport; use Mcp\Tests\Unit\Fixtures\ThrowingRequest; @@ -1034,6 +1036,48 @@ public function testSuccessfulRequestReturnsResponseWithSessionId(): void $this->assertEquals(['status' => 'ok'], $message['result']); } + #[TestDox('An inline response transport gets the response directly, not through the session queue')] + public function testInlineResponseTransportGetsResponseDirectly(): void + { + $handler = $this->createMock(RequestHandlerInterface::class); + $handler->method('supports')->willReturn(true); + $handler->method('handle')->willReturn(new Response(1, ['status' => 'ok'])); + + $sessions = new SessionManager(new InMemorySessionStore(), gcProbability: 0); + $sessionId = Uuid::v4(); + $sessions->createWithId($sessionId)->save(); + + $transport = new class extends InMemoryTransport implements InlineResponseTransportInterface { + /** @var list}> */ + public array $sent = []; + + public function send(string $data, array $context): void + { + $this->sent[] = [$data, $context]; + } + }; + + $protocol = new Protocol( + requestHandlers: [$handler], + notificationHandlers: [], + messageFactory: MessageFactory::make(), + sessionManager: $sessions, + ); + + $protocol->processInput( + $transport, + '{"jsonrpc": "2.0", "id": 1, "method": "tools/list"}', + $sessionId + ); + + $this->assertCount(1, $transport->sent); + [$data, $context] = $transport->sent[0]; + $this->assertSame(['status' => 'ok'], json_decode($data, true)['result']); + $this->assertSame('response', $context['type']); + $this->assertEquals($sessionId, $context['session_id']); + $this->assertSame([], $protocol->consumeOutgoingMessages($sessionId)); + } + #[TestDox('Batch requests are processed and send multiple responses')] public function testBatchRequestsAreProcessed(): void { From 483f7dce0ffa99fd68e32c3488f57fa5dbd596f8 Mon Sep 17 00:00:00 2001 From: Christopher Hertel Date: Fri, 9 Oct 2026 23:10:28 +0200 Subject: [PATCH 4/5] [Server] Log immediate responses at debug level --- src/Server/Protocol.php | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/Server/Protocol.php b/src/Server/Protocol.php index 4af596932..c6fea38b9 100644 --- a/src/Server/Protocol.php +++ b/src/Server/Protocol.php @@ -493,7 +493,7 @@ private function sendResponse(TransportInterface $transport, Response|Error $res // request of the same session: a transport that can answer on the request's // own exchange gets it directly. if (null === $session || $transport instanceof InlineResponseTransportInterface) { - $this->logger->info('Sending immediate response', [ + $this->logger->debug('Sending immediate response', [ 'response_id' => $response->getId(), ]); From abc599f051cb3d7546ef6ec1b8bf9a6f0e72e203 Mon Sep 17 00:00:00 2001 From: Christopher Hertel Date: Fri, 9 Oct 2026 23:14:34 +0200 Subject: [PATCH 5/5] [Server] Cover a JSON batch carrying queued notifications and inline responses --- .../Transport/StreamableHttpTransportTest.php | 31 +++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/tests/Unit/Server/Transport/StreamableHttpTransportTest.php b/tests/Unit/Server/Transport/StreamableHttpTransportTest.php index 9716b3357..2e0a7a288 100644 --- a/tests/Unit/Server/Transport/StreamableHttpTransportTest.php +++ b/tests/Unit/Server/Transport/StreamableHttpTransportTest.php @@ -15,6 +15,7 @@ use Mcp\Schema\JsonRpc\Error; use Mcp\Server; use Mcp\Server\RequestContext; +use Mcp\Server\Session\Session; use Mcp\Server\Transport\Http\Middleware\CorsMiddleware; use Mcp\Server\Transport\Http\Middleware\DnsRebindingProtectionMiddleware; use Mcp\Server\Transport\Http\Middleware\PassthroughMiddleware; @@ -526,6 +527,36 @@ public function testStreamedBatchCarriesInlineResponses(): void $this->assertMatchesRegularExpression('/"id":3,"result".*"progressToken":"p".*"id":2,"result"/s', $output); } + #[TestDox('a batch answered as JSON carries the queued notifications first, then its responses, in one array')] + public function testJsonBatchCarriesQueuedNotificationsAndInlineResponses(): void + { + $store = new InterleavingSessionStore(); + $sessionId = $this->post($store, '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"test","version":"1.0"}}}') + ->getHeaderLine(StreamableHttpTransport::SESSION_HEADER); + + // A notification another request of the session queued, e.g. a resource update. + $session = new Session($store, Uuid::fromString($sessionId)); + $session->set('_mcp.outgoing_queue', [[ + 'message' => '{"jsonrpc":"2.0","method":"notifications/resources/updated","params":{"uri":"file:///a"}}', + 'context' => ['type' => 'notification'], + ]]); + $session->save(); + + $response = $this->post($store, '[{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"echo","arguments":{"text":"a"}}},{"jsonrpc":"2.0","id":3,"method":"ping"}]', $sessionId); + + $this->assertSame(200, $response->getStatusCode()); + $this->assertSame('application/json', $response->getHeaderLine('Content-Type')); + $this->assertSame($sessionId, $response->getHeaderLine(StreamableHttpTransport::SESSION_HEADER)); + + $messages = json_decode((string) $response->getBody(), true); + $this->assertIsArray($messages); + $this->assertTrue(array_is_list($messages)); + $this->assertSame( + ['notifications/resources/updated', 2, 3], + array_map(static fn (array $message): string|int => $message['method'] ?? $message['id'], $messages), + ); + } + /** * Sends one POST to a fresh server sharing $store, like a PHP worker would. */