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 @@ -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.
Expand Down
139 changes: 116 additions & 23 deletions src/Server/Protocol.php
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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<TransportInterface<mixed>, list<array{message: string, context: array<string, mixed>}>>
*/
private \WeakMap $fiberOutgoingMessages;

/**
* @param array<int, RequestHandlerInterface<ResultInterface|array<string, mixed>>> $requestHandlers
* @param array<int, NotificationHandlerInterface> $notificationHandlers
Expand All @@ -96,6 +107,7 @@ public function __construct(
private readonly ?RequestStateCodec $requestStateCodec = null,
) {
$this->awaitedRequestIds = new \WeakMap();
$this->fiberOutgoingMessages = new \WeakMap();
}

/**
Expand All @@ -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);
}
});
Expand Down Expand Up @@ -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
Expand All @@ -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();
Expand Down Expand Up @@ -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,
Expand All @@ -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;
}

Expand Down Expand Up @@ -552,6 +583,55 @@ private function sendResponse(TransportInterface $transport, Response|Error $res
* @param array<string, mixed> $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<mixed> $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<mixed>|null $transport
*
* @return list<array{message: string, context: array<string, mixed>}>
*/
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<string, mixed> $context
*
* @return array{message: string, context: array<string, mixed>}|null
*/
private function encodeOutgoing(Request|Notification $message, array $context): ?array
{
try {
$encoded = json_encode($message, \JSON_THROW_ON_ERROR);
Expand All @@ -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);
}

/**
Expand Down Expand Up @@ -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<mixed>|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.');
Expand All @@ -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;
}

/**
Expand Down
8 changes: 8 additions & 0 deletions tests/Unit/Fixtures/PollingLoopTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,14 @@ public function getPendingRequestIds(): array
return array_keys($this->getPendingRequests($this->sessionId));
}

/**
* @return array<int, array<mixed>> 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
*/
Expand Down
Loading
Loading