diff --git a/CHANGELOG.md b/CHANGELOG.md index e34852a9..e1b69b61 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,7 @@ All notable changes to `mcp/sdk` will be documented in this file. * [BC Break] Reject a `Tool` input schema whose `properties` is not an object or whose `required` is neither a list nor `null`, instead of silently replacing the member. Reject a `completion/complete` whose `argument` is missing `name` or `value`, instead of completing against an empty prefix. * Add `HttpTransport::getSessionId()` to read the server-minted `Mcp-Session-Id`: a request-scoped caller can persist it and pass it back through the constructor's `$headers` on a later transport. Always `null` on `2026-07-28`, which removed protocol-level sessions. * Fix OIDC discovery rejecting issuers with a trailing slash (e.g. Authentik, Auth0). +* Fix parallel elicitations on one session getting each other's answers: requests to the client get random ids instead of a session counter, and a request or notification a handler sends goes out on the stream of its own call. * 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. diff --git a/src/Server/Protocol.php b/src/Server/Protocol.php index bbd25019..4827f58d 100644 --- a/src/Server/Protocol.php +++ b/src/Server/Protocol.php @@ -50,8 +50,11 @@ */ class Protocol { - /** Session key for request ID counter */ - private const SESSION_REQUEST_ID_COUNTER = '_mcp.request_id_counter'; + /** + * Largest id a request to the client gets: JavaScript clients read JSON numbers as doubles, + * which hold integers exactly only up to 2^53 - 1. + */ + private const MAX_REQUEST_ID = 9007199254740991; /** Session key for pending outgoing requests */ private const SESSION_PENDING_REQUESTS = '_mcp.pending_requests'; @@ -81,6 +84,14 @@ class Protocol */ private \WeakMap $awaitedRequestIds; + /** + * What each transport's fiber sends the client, kept out of the session for the same reason: + * a request or notification must go out on the stream of the call that produced it. + * + * @var \WeakMap, list}>> + */ + private \WeakMap $fiberOutgoingMessages; + /** * @param array>> $requestHandlers * @param array $notificationHandlers @@ -96,6 +107,7 @@ public function __construct( private readonly ?RequestStateCodec $requestStateCodec = null, ) { $this->awaitedRequestIds = new \WeakMap(); + $this->fiberOutgoingMessages = new \WeakMap(); } /** @@ -111,19 +123,23 @@ public function connect(TransportInterface $transport): void $transport->onSessionEnd($this->destroySession(...)); - $transport->setOutgoingMessagesProvider($this->consumeOutgoingMessages(...)); - // The transport keeps these callbacks, so they reference it weakly to not keep it alive. $transportRef = \WeakReference::create($transport); + $transport->setOutgoingMessagesProvider(fn (Uuid $sessionId): array => [ + ...$this->takeFiberOutgoingMessages($transportRef->get()), + ...$this->consumeOutgoingMessages($sessionId), + ]); + $transport->setPendingRequestsProvider(fn (Uuid $sessionId): array => $this->getAwaitedPendingRequests($transportRef->get(), $sessionId)); $transport->setResponseFinder($this->checkResponse(...)); $transport->setFiberYieldHandler(function (mixed $yieldedValue, ?Uuid $sessionId) use ($transportRef): void { - $requestId = $this->handleFiberYield($yieldedValue, $sessionId); + $transport = $transportRef->get(); + $requestId = $this->handleFiberYield($yieldedValue, $sessionId, $transport); - if (null !== $transport = $transportRef->get()) { + if (null !== $transport) { $this->trackAwaitedRequest($transport, $requestId); } }); @@ -351,10 +367,12 @@ private function handleRequest(TransportInterface $transport, Request $request, $beforeSuspension = $session->all(); $awaitedRequestId = null; + $outgoing = null; if ($result instanceof NotificationSuspension) { - $this->sendNotification($result->notification, $session); + $outgoing = $result->notification; } elseif ($result instanceof RequestSuspension) { - $awaitedRequestId = $this->sendRequest($result->request, $result->timeout, $session); + $awaitedRequestId = $this->registerRequest($result->request, $result->timeout, $session); + $outgoing = $result->request->withId($awaitedRequestId); } // The transport resumes the fiber from what the session holds: it must @@ -377,6 +395,10 @@ private function handleRequest(TransportInterface $transport, Request $request, $this->trackAwaitedRequest($transport, $awaitedRequestId); $transport->attachFiberToSession($fiber, $session->getId()); + if (null !== $outgoing) { + $this->queueFiberOutgoing($transport, $outgoing); + } + return; } $finalResult = $fiber->getReturn(); @@ -468,11 +490,22 @@ private function handleNotification(Notification $notification, SessionInterface */ public function sendRequest(Request $request, int $timeout, SessionInterface $session): int { - $counter = $session->get(self::SESSION_REQUEST_ID_COUNTER, 1000); - $requestId = $counter++; - $session->set(self::SESSION_REQUEST_ID_COUNTER, $counter); + $requestId = $this->registerRequest($request, $timeout, $session); - $requestWithId = $request->withId($requestId); + $this->queueOutgoing($request->withId($requestId), ['type' => 'request'], $session); + + return $requestId; + } + + /** + * Picks the id of a request to the client and stores the request as pending in the session. + * + * The id is random, not counted in the session: concurrent requests of a client load + * the session at the same time and would count up to the same id. + */ + private function registerRequest(Request $request, int $timeout, SessionInterface $session): int + { + $requestId = random_int(1, self::MAX_REQUEST_ID); $this->logger->info('Queueing server request to client', [ 'request_id' => $requestId, @@ -487,8 +520,6 @@ public function sendRequest(Request $request, int $timeout, SessionInterface $se ]; $session->set(self::SESSION_PENDING_REQUESTS, $pending); - $this->queueOutgoing($requestWithId, ['type' => 'request'], $session); - return $requestId; } @@ -552,6 +583,55 @@ private function sendResponse(TransportInterface $transport, Response|Error $res * @param array $context */ private function queueOutgoing(Request|Notification $message, array $context, SessionInterface $session): void + { + if (null === $outgoing = $this->encodeOutgoing($message, $context)) { + return; + } + + $queue = $session->get(self::SESSION_OUTGOING_QUEUE, []); + $queue[] = $outgoing; + $session->set(self::SESSION_OUTGOING_QUEUE, $queue); + } + + /** + * Queues what a transport's fiber sends the client, for that transport alone. + * + * @param TransportInterface $transport + */ + private function queueFiberOutgoing(TransportInterface $transport, Request|Notification $message): void + { + if (null === $outgoing = $this->encodeOutgoing($message, ['type' => $message instanceof Request ? 'request' : 'notification'])) { + return; + } + + $queue = $this->fiberOutgoingMessages[$transport] ?? []; + $queue[] = $outgoing; + $this->fiberOutgoingMessages[$transport] = $queue; + } + + /** + * @param TransportInterface|null $transport + * + * @return list}> + */ + private function takeFiberOutgoingMessages(?TransportInterface $transport): array + { + if (null === $transport || !isset($this->fiberOutgoingMessages[$transport])) { + return []; + } + + $queue = $this->fiberOutgoingMessages[$transport]; + unset($this->fiberOutgoingMessages[$transport]); + + return $queue; + } + + /** + * @param array $context + * + * @return array{message: string, context: array}|null + */ + private function encodeOutgoing(Request|Notification $message, array $context): ?array { try { $encoded = json_encode($message, \JSON_THROW_ON_ERROR); @@ -560,15 +640,13 @@ private function queueOutgoing(Request|Notification $message, array $context, Se 'exception' => $e, ]); - return; + return null; } - $queue = $session->get(self::SESSION_OUTGOING_QUEUE, []); - $queue[] = [ + return [ 'message' => $encoded, 'context' => $context, ]; - $session->set(self::SESSION_OUTGOING_QUEUE, $queue); } /** @@ -651,11 +729,15 @@ public function getPendingRequests(Uuid $sessionId): array /** * Handle values yielded by Fibers during transport-managed resumes. * - * @param FiberSuspend|null $yieldedValue + * With the transport the fiber runs on, what it sends goes out on that transport alone; + * without one, it is queued for whichever stream of the session asks first. + * + * @param FiberSuspend|null $yieldedValue + * @param TransportInterface|null $transport * * @return int|null the ID of the request sent to the client, which the fiber now waits on */ - public function handleFiberYield(mixed $yieldedValue, ?Uuid $sessionId): ?int + public function handleFiberYield(mixed $yieldedValue, ?Uuid $sessionId, ?TransportInterface $transport = null): ?int { if (!$sessionId) { $this->logger->warning('Fiber yielded value without associated session context.'); @@ -681,17 +763,28 @@ public function handleFiberYield(mixed $yieldedValue, ?Uuid $sessionId): ?int ]); } + $requestId = null; + try { if ($yieldedValue instanceof RequestSuspension) { - return $this->sendRequest($yieldedValue->request, $yieldedValue->timeout, $session); + $requestId = $this->registerRequest($yieldedValue->request, $yieldedValue->timeout, $session); + $outgoing = $yieldedValue->request->withId($requestId); + } else { + $outgoing = $yieldedValue->notification; } - $this->sendNotification($yieldedValue->notification, $session); + if (null === $transport) { + $this->queueOutgoing($outgoing, ['type' => null === $requestId ? 'notification' : 'request'], $session); + } } finally { $session->save(); } - return null; + if (null !== $transport) { + $this->queueFiberOutgoing($transport, $outgoing); + } + + return $requestId; } /** diff --git a/tests/Unit/Fixtures/PollingLoopTransport.php b/tests/Unit/Fixtures/PollingLoopTransport.php index 34317a18..b27ede04 100644 --- a/tests/Unit/Fixtures/PollingLoopTransport.php +++ b/tests/Unit/Fixtures/PollingLoopTransport.php @@ -32,6 +32,14 @@ public function getPendingRequestIds(): array return array_keys($this->getPendingRequests($this->sessionId)); } + /** + * @return array> the messages this stream sends the client next, decoded + */ + public function takeOutgoingMessages(): array + { + return array_map(static fn (array $message): array => json_decode($message['message'], true), $this->getOutgoingMessages($this->sessionId)); + } + /** * @param FiberSuspend $yielded */ diff --git a/tests/Unit/Server/ProtocolTest.php b/tests/Unit/Server/ProtocolTest.php index 26af24fd..f7015af6 100644 --- a/tests/Unit/Server/ProtocolTest.php +++ b/tests/Unit/Server/ProtocolTest.php @@ -38,6 +38,7 @@ use Mcp\Tests\Unit\Fixtures\PollingLoopTransport; use Mcp\Tests\Unit\Fixtures\RecordingTransport; use Mcp\Tests\Unit\Fixtures\ThrowingRequest; +use Mcp\Tests\Unit\Server\Session\Fixture\InterleavingSessionStore; use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\Attributes\TestDox; use PHPUnit\Framework\MockObject\MockObject; @@ -919,10 +920,8 @@ public function testOutboundRequestFailureIsAnsweredUnderInboundRequestId(): voi $this->assertSame($exception, $errorEvents[0]->getThrowable()); $this->assertSame(1, $errorEvents[0]->getError()->getId()); - // The outbound request was queued for the client before the failure; the error is the response. - $outgoing = array_map(static fn (array $outgoingMessage): array => json_decode($outgoingMessage['message'], true), $protocol->consumeOutgoingMessages($sessionId)); - $this->assertCount(1, $outgoing); - $this->assertSame('ping', $outgoing[0]['method']); + // Nobody would resume the fiber, so the outbound request is not sent; the error is the response. + $this->assertSame([], $protocol->consumeOutgoingMessages($sessionId)); $this->assertCount(1, $sent); $this->assertSame(1, $sent[0]['id']); @@ -933,46 +932,123 @@ public function testOutboundRequestFailureIsAnsweredUnderInboundRequestId(): voi public function testConcurrentStreamsPollOnlyTheirOwnPendingRequest(): void { [$protocol, $sessionId, $firstStream, $secondStream] = $this->startTwoStreamsWaitingOnClient(); + $firstIds = $firstStream->getPendingRequestIds(); + $secondIds = $secondStream->getPendingRequestIds(); - $protocol->processInput($secondStream, '{"jsonrpc": "2.0", "id": 1001, "result": {}}', $sessionId); + $protocol->processInput($secondStream, \sprintf('{"jsonrpc": "2.0", "id": %d, "result": {}}', $secondIds[0]), $sessionId); - $this->assertSame([1000], $firstStream->getPendingRequestIds()); - $this->assertSame([1001], $secondStream->getPendingRequestIds()); + $this->assertCount(1, $firstIds); + $this->assertCount(1, $secondIds); + $this->assertNotSame($firstIds, $secondIds); + $this->assertSame($firstIds, $firstStream->getPendingRequestIds()); + $this->assertSame($secondIds, $secondStream->getPendingRequestIds()); } #[TestDox('A client request a fiber sends after resuming is polled only by its own stream')] public function testRequestYieldedOnResumeIsPolledOnlyByItsOwnStream(): void { [, $sessionId, $firstStream, $secondStream] = $this->startTwoStreamsWaitingOnClient(); + $firstIds = $firstStream->getPendingRequestIds(); + $secondIds = $secondStream->getPendingRequestIds(); $firstStream->yieldFromFiber(new RequestSuspension(new PingRequest(), $sessionId->toRfc4122(), 5)); - $this->assertSame([1002], $firstStream->getPendingRequestIds()); - $this->assertSame([1001], $secondStream->getPendingRequestIds()); + $this->assertCount(1, $firstStream->getPendingRequestIds()); + $this->assertNotSame($firstIds, $firstStream->getPendingRequestIds()); + $this->assertNotSame($secondIds, $firstStream->getPendingRequestIds()); + $this->assertSame($secondIds, $secondStream->getPendingRequestIds()); } #[TestDox('A stream whose fiber resumes and sends a notification no longer polls the request it was waiting on')] public function testNotificationYieldedOnResumeClearsTheAwaitedRequest(): void { [, $sessionId, $firstStream, $secondStream] = $this->startTwoStreamsWaitingOnClient(); + $secondIds = $secondStream->getPendingRequestIds(); $firstStream->yieldFromFiber(new NotificationSuspension(new LoggingMessageNotification(LoggingLevel::Info, 'hello'), $sessionId->toRfc4122())); $this->assertSame([], $firstStream->getPendingRequestIds()); - $this->assertSame([1001], $secondStream->getPendingRequestIds()); + $this->assertSame($secondIds, $secondStream->getPendingRequestIds()); } #[TestDox('A stream whose fiber first suspends on a notification polls none of the session\'s pending requests')] public function testStreamSuspendedOnNotificationPollsNoPendingRequest(): void { [$protocol, $sessionId, , $secondStream] = $this->startTwoStreamsWaitingOnClient(); + $secondIds = $secondStream->getPendingRequestIds(); $thirdStream = new PollingLoopTransport(); $protocol->connect($thirdStream); $protocol->processInput($thirdStream, '{"jsonrpc": "2.0", "id": 3, "method": "ping"}', $sessionId); $this->assertSame([], $thirdStream->getPendingRequestIds()); - $this->assertSame([1001], $secondStream->getPendingRequestIds()); + $this->assertSame($secondIds, $secondStream->getPendingRequestIds()); + } + + #[TestDox('Concurrent streams on one session each send the client request their own fiber sent')] + public function testConcurrentStreamsSendOnlyTheirOwnRequest(): void + { + [, , $firstStream, $secondStream] = $this->startTwoStreamsWaitingOnClient(); + + $firstOutgoing = $firstStream->takeOutgoingMessages(); + $secondOutgoing = $secondStream->takeOutgoingMessages(); + + $this->assertCount(1, $firstOutgoing); + $this->assertSame($firstStream->getPendingRequestIds(), [$firstOutgoing[0]['id']]); + $this->assertCount(1, $secondOutgoing); + $this->assertSame($secondStream->getPendingRequestIds(), [$secondOutgoing[0]['id']]); + } + + #[TestDox('A client request a fiber sends after resuming goes out on its own stream')] + public function testRequestYieldedOnResumeIsSentOnItsOwnStream(): void + { + [, $sessionId, $firstStream, $secondStream] = $this->startTwoStreamsWaitingOnClient(); + $firstStream->takeOutgoingMessages(); + $secondStream->takeOutgoingMessages(); + + $firstStream->yieldFromFiber(new RequestSuspension(new PingRequest(), $sessionId->toRfc4122(), 5)); + + $this->assertSame([], $secondStream->takeOutgoingMessages()); + $firstOutgoing = $firstStream->takeOutgoingMessages(); + $this->assertCount(1, $firstOutgoing); + $this->assertSame($firstStream->getPendingRequestIds(), [$firstOutgoing[0]['id']]); + } + + #[TestDox('Client requests of calls handled by different workers on one session get distinct ids')] + public function testConcurrentWorkersSendClientRequestsUnderDistinctIds(): void + { + $handler = $this->createMock(RequestHandlerInterface::class); + $handler->method('supports')->willReturn(true); + $handler->method('handle')->willReturnCallback(static function (Request $request, SessionInterface $session): Response { + \Fiber::suspend(new RequestSuspension(new PingRequest(), $session->getId()->toRfc4122(), 5)); + + return new Response(1, []); + }); + + $store = new InterleavingSessionStore(); + $sessionManager = new SessionManager($store, gcProbability: 0); + $session = $sessionManager->create(); + $session->save(); + $sessionId = $session->getId(); + + $firstWorker = new Protocol([$handler], [], MessageFactory::make(), $sessionManager); + $secondWorker = new Protocol([$handler], [], MessageFactory::make(), $sessionManager); + $firstStream = new PollingLoopTransport(); + $secondStream = new PollingLoopTransport(); + $firstWorker->connect($firstStream); + $secondWorker->connect($secondStream); + + // The second call is handled after the first worker loaded the session, before it saved it. + $secondIds = null; + $store->interleaveAfterNextRead(static function () use ($secondWorker, $secondStream, $sessionId, &$secondIds): void { + $secondWorker->processInput($secondStream, '{"jsonrpc": "2.0", "id": 2, "method": "ping"}', $sessionId); + $secondIds = $secondStream->getPendingRequestIds(); + }); + $firstWorker->processInput($firstStream, '{"jsonrpc": "2.0", "id": 1, "method": "ping"}', $sessionId); + + $this->assertCount(1, $firstStream->getPendingRequestIds()); + $this->assertCount(1, $secondIds ?? []); + $this->assertNotSame($secondIds, $firstStream->getPendingRequestIds()); } /** @@ -1913,11 +1989,12 @@ public function testFiberYieldedRequestSuspensionIsQueued(): void $this->assertSame(45, $suspension->timeout); $this->assertNull($suspension->inputKey); - $protocol->handleFiberYield($suspension, $sessionId); + $requestId = $protocol->handleFiberYield($suspension, $sessionId); + $this->assertNotNull($requestId); $pending = $protocol->getPendingRequests($sessionId); $this->assertCount(1, $pending); - $this->assertSame(45, $pending[1000]['timeout']); + $this->assertSame(45, $pending[$requestId]['timeout']); $outgoing = $protocol->consumeOutgoingMessages($sessionId); $this->assertCount(1, $outgoing); @@ -1925,7 +2002,7 @@ public function testFiberYieldedRequestSuspensionIsQueued(): void $message = json_decode($outgoing[0]['message'], true); $this->assertSame('roots/list', $message['method']); - $this->assertSame(1000, $message['id']); + $this->assertSame($requestId, $message['id']); } #[TestDox('A fiber yield that is not a suspension object is dropped without touching the session')]