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 @@ -45,6 +45,7 @@ All notable changes to `mcp/sdk` will be documented in this file.
* On a `2026-07-28` connection, `Client::setLoggingLevel()` stamps the level on every following request, `Client::ping()` sends `server/discover` and `Client::sendRootsListChanged()` sends nothing.
* Fail a client request at once when the HTTP server refuses it with an error status or the stdio server process exits, instead of waiting out the timeout.
* [BC Break] Bump `MessageInterface::PROTOCOL_VERSION` to `2026-07-28`. Use `ProtocolVersion::latestHandshake()` where a handshake revision is needed, e.g. in an `initialize` answer.
* [BC Break] Fix concurrent stdio tool calls that report progress, log, elicit or sample never finishing: `StdioTransport` keeps every suspended fiber and resumes each with the answer to its own request. A transport passes the suspended `\Fiber` as last argument to the pending-requests provider and the fiber yield handler.

0.8.0
-----
Expand Down
35 changes: 14 additions & 21 deletions src/Server/Protocol.php
Original file line number Diff line number Diff line change
Expand Up @@ -74,10 +74,10 @@ class Protocol
private const INTERNAL_ERROR_MESSAGE = 'Internal server error.';

/**
* The client request each transport's fiber is suspended on. Pending requests are stored in the
* session, which concurrent streams share, so a stream must only poll the one its fiber sent.
* The client request each suspended fiber waits on. Pending requests are stored in the session,
* which concurrent calls share, so a fiber must only be resumed with the answer to the one it sent.
*
* @var \WeakMap<TransportInterface<mixed>, int>
* @var \WeakMap<McpFiber, int>
*/
private \WeakMap $awaitedRequestIds;

Expand Down Expand Up @@ -113,19 +113,12 @@ public function connect(TransportInterface $transport): void

$transport->setOutgoingMessagesProvider($this->consumeOutgoingMessages(...));

// The transport keeps these callbacks, so they reference it weakly to not keep it alive.
$transportRef = \WeakReference::create($transport);

$transport->setPendingRequestsProvider(fn (Uuid $sessionId): array => $this->getAwaitedPendingRequests($transportRef->get(), $sessionId));
$transport->setPendingRequestsProvider(fn (Uuid $sessionId, \Fiber $fiber): array => $this->getAwaitedPendingRequests($fiber, $sessionId));

$transport->setResponseFinder($this->checkResponse(...));

$transport->setFiberYieldHandler(function (mixed $yieldedValue, ?Uuid $sessionId) use ($transportRef): void {
$requestId = $this->handleFiberYield($yieldedValue, $sessionId);

if (null !== $transport = $transportRef->get()) {
$this->trackAwaitedRequest($transport, $requestId);
}
$transport->setFiberYieldHandler(function (mixed $yieldedValue, ?Uuid $sessionId, \Fiber $fiber): void {
$this->trackAwaitedRequest($fiber, $this->handleFiberYield($yieldedValue, $sessionId));
});

$this->logger->info('Protocol connected to transport', ['transport' => $transport::class]);
Expand Down Expand Up @@ -374,7 +367,7 @@ private function handleRequest(TransportInterface $transport, Request $request,
throw new RuntimeException('Failed to save the session of a suspended request.');
}

$this->trackAwaitedRequest($transport, $awaitedRequestId);
$this->trackAwaitedRequest($fiber, $awaitedRequestId);
$transport->attachFiberToSession($fiber, $session->getId());

return;
Expand Down Expand Up @@ -695,27 +688,27 @@ public function handleFiberYield(mixed $yieldedValue, ?Uuid $sessionId): ?int
}

/**
* @param TransportInterface<mixed> $transport
* @param McpFiber $fiber
*/
private function trackAwaitedRequest(TransportInterface $transport, ?int $requestId): void
private function trackAwaitedRequest(\Fiber $fiber, ?int $requestId): void
{
if (null === $requestId) {
unset($this->awaitedRequestIds[$transport]);
unset($this->awaitedRequestIds[$fiber]);

return;
}

$this->awaitedRequestIds[$transport] = $requestId;
$this->awaitedRequestIds[$fiber] = $requestId;
}

/**
* @param TransportInterface<mixed>|null $transport
* @param McpFiber $fiber
*
* @return array<int, mixed>
*/
private function getAwaitedPendingRequests(?TransportInterface $transport, Uuid $sessionId): array
private function getAwaitedPendingRequests(\Fiber $fiber, Uuid $sessionId): array
{
$requestId = null !== $transport ? $this->awaitedRequestIds[$transport] ?? null : null;
$requestId = $this->awaitedRequestIds[$fiber] ?? null;
if (null === $requestId) {
return [];
}
Expand Down
11 changes: 7 additions & 4 deletions src/Server/Transport/BaseTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -86,12 +86,14 @@ protected function getOutgoingMessages(?Uuid $sessionId): array
}

/**
* @param McpFiber $fiber
*
* @return array<int, array<string, mixed>>
*/
protected function getPendingRequests(?Uuid $sessionId): array
protected function getPendingRequests(?Uuid $sessionId, \Fiber $fiber): array
{
if ($sessionId && \is_callable($this->pendingRequestsProvider)) {
return ($this->pendingRequestsProvider)($sessionId);
return ($this->pendingRequestsProvider)($sessionId, $fiber);
}

return [];
Expand All @@ -111,15 +113,16 @@ protected function checkForResponse(int $requestId, ?Uuid $sessionId): Response|

/**
* @param FiberSuspend|null $yielded
* @param McpFiber $fiber the fiber that yielded
*/
protected function handleFiberYield(mixed $yielded, ?Uuid $sessionId): void
protected function handleFiberYield(mixed $yielded, ?Uuid $sessionId, \Fiber $fiber): void
{
if (null === $yielded || !\is_callable($this->fiberYieldHandler)) {
return;
}

try {
($this->fiberYieldHandler)($yielded, $sessionId);
($this->fiberYieldHandler)($yielded, $sessionId, $fiber);
} catch (\Throwable $e) {
$this->logger->error('Fiber yield handler failed.', [
'exception' => $e,
Expand Down
7 changes: 4 additions & 3 deletions src/Server/Transport/ManagesTransportCallbacks.php
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
* @phpstan-import-type FiberReturn from \Mcp\Server\Transport\TransportInterface
* @phpstan-import-type FiberResume from \Mcp\Server\Transport\TransportInterface
* @phpstan-import-type FiberSuspend from \Mcp\Server\Transport\TransportInterface
* @phpstan-import-type McpFiber from \Mcp\Server\Transport\TransportInterface
*
* @author Kyrian Obikwelu <koshnawaza@gmail.com>
* */
Expand All @@ -36,13 +37,13 @@ trait ManagesTransportCallbacks
/** @var callable(Uuid): array<int, array{message: string, context: array<string, mixed>}> */
protected $outgoingMessagesProvider;

/** @var callable(Uuid): array<int, array<string, mixed>> */
/** @var callable(Uuid, McpFiber): array<int, array<string, mixed>> */
protected $pendingRequestsProvider;

/** @var (callable(int, Uuid): (Response<array<string, mixed>>|Error|null))|null */
protected $responseFinder;

/** @var callable(FiberSuspend|null, ?Uuid): void */
/** @var callable(FiberSuspend|null, ?Uuid, McpFiber): void */
protected $fiberYieldHandler;

public function onMessage(callable $listener): void
Expand Down Expand Up @@ -74,7 +75,7 @@ public function setResponseFinder(callable $finder): void
}

/**
* @param callable(FiberSuspend|null, ?Uuid): void $handler
* @param callable(FiberSuspend|null, ?Uuid, McpFiber): void $handler
*/
public function setFiberYieldHandler(callable $handler): void
{
Expand Down
65 changes: 43 additions & 22 deletions src/Server/Transport/StdioTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
use Mcp\Server\Transport\Stdio\RunnerState;
use Mcp\Server\Wire\InboundClassifier;
use Psr\Log\LoggerInterface;
use Symfony\Component\Uid\Uuid;

/**
* Serves one client over the standard streams, in whichever protocol era it
Expand Down Expand Up @@ -79,6 +80,14 @@ class StdioTransport extends BaseTransport implements StatelessAwareTransportInt
*/
private array $streams = [];

/**
* Handshake-era handlers suspended on the client, by object id: the
* client pipelines its requests, so several can be in flight at once.
*
* @var array<int, McpFiber>
*/
private array $fibers = [];

/**
* @param resource $input
* @param resource $output
Expand Down Expand Up @@ -124,14 +133,20 @@ public function send(string $data, array $context): void
$this->writeLine($data);
}

public function attachFiberToSession(\Fiber $fiber, Uuid $sessionId): void
{
$this->fibers[spl_object_id($fiber)] = $fiber;
$this->sessionId = $sessionId;
}

public function listen(): int
{
$this->logger->info('StdioTransport is listening for messages on STDIN...');
stream_set_blocking($this->input, false);

while (!feof($this->input) && RunnerState::RUNNING === $this->runnerControl->getState()) {
$this->processInput();
$this->processFiber();
$this->processFibers();
$this->processStreams();
$this->flushOutgoingMessages();
}
Expand Down Expand Up @@ -341,27 +356,35 @@ private function writeError(Error $error): void
$this->writeLine(json_encode($error, \JSON_THROW_ON_ERROR | \JSON_UNESCAPED_SLASHES));
}

private function processFiber(): void
/**
* Resumes every suspended handler that can go on: one that sent a
* notification right away, one that sent a request once its answer is in
* or it timed out.
*/
private function processFibers(): void
{
if (null === $this->sessionFiber) {
return;
}

if ($this->sessionFiber->isTerminated()) {
$this->handleFiberTermination($this->sessionFiber);

return;
}
foreach ($this->fibers as $key => $fiber) {
if ($fiber->isSuspended()) {
$this->resumeFiber($fiber);
}

if (!$this->sessionFiber->isSuspended()) {
return;
if ($fiber->isTerminated()) {
unset($this->fibers[$key]);
$this->handleFiberTermination($fiber);
}
}
}

$pendingRequests = $this->getPendingRequests($this->sessionId);
/**
* @param McpFiber $fiber
*/
private function resumeFiber(\Fiber $fiber): void
{
$pendingRequests = $this->getPendingRequests($this->sessionId, $fiber);

if (empty($pendingRequests)) {
$yielded = $this->sessionFiber->resume();
$this->handleFiberYield($yielded, $this->sessionId);
$yielded = $fiber->resume();
$this->handleFiberYield($yielded, $this->sessionId, $fiber);

return;
}
Expand All @@ -374,16 +397,16 @@ private function processFiber(): void
$response = $this->checkForResponse($requestId, $this->sessionId);

if (null !== $response) {
$yielded = $this->sessionFiber->resume($response);
$this->handleFiberYield($yielded, $this->sessionId);
$yielded = $fiber->resume($response);
$this->handleFiberYield($yielded, $this->sessionId, $fiber);

return;
}

if (time() - $timestamp >= $timeout) {
$error = Error::forInternalError('Request timed out', $requestId);
$yielded = $this->sessionFiber->resume($error);
$this->handleFiberYield($yielded, $this->sessionId);
$yielded = $fiber->resume($error);
$this->handleFiberYield($yielded, $this->sessionId, $fiber);

return;
}
Expand All @@ -405,8 +428,6 @@ private function handleFiberTermination(\Fiber $fiber): void
$this->logger->error('STDIO: Failed to encode final Fiber result.', ['exception' => $e]);
}
}

$this->sessionFiber = null;
}

private function flushOutgoingMessages(): void
Expand Down
8 changes: 4 additions & 4 deletions src/Server/Transport/StreamableHttpTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -299,11 +299,11 @@ protected function createStreamedResponse(): ResponseInterface
while ($fiber->isSuspended()) {
$this->flushOutgoingMessages($this->sessionId);

$pendingRequests = $this->getPendingRequests($this->sessionId);
$pendingRequests = $this->getPendingRequests($this->sessionId, $fiber);

if (empty($pendingRequests)) {
$yielded = $fiber->resume();
$this->handleFiberYield($yielded, $this->sessionId);
$this->handleFiberYield($yielded, $this->sessionId, $fiber);
continue;
}

Expand All @@ -317,15 +317,15 @@ protected function createStreamedResponse(): ResponseInterface

if (null !== $response) {
$yielded = $fiber->resume($response);
$this->handleFiberYield($yielded, $this->sessionId);
$this->handleFiberYield($yielded, $this->sessionId, $fiber);
$resumed = true;
break;
}

if ($this->clock->now()->getTimestamp() - $timestamp >= $timeout) {
$error = Error::forInternalError('Request timed out', $requestId);
$yielded = $fiber->resume($error);
$this->handleFiberYield($yielded, $this->sessionId);
$this->handleFiberYield($yielded, $this->sessionId, $fiber);
$resumed = true;
break;
}
Expand Down
6 changes: 3 additions & 3 deletions src/Server/Transport/TransportInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -97,11 +97,11 @@ public function onSessionEnd(callable $listener): void;
public function setOutgoingMessagesProvider(callable $provider): void;

/**
* Set a provider function to retrieve all pending server-initiated requests.
* Set a provider function to retrieve the pending server-initiated request a suspended Fiber waits on.
*
* The transport calls this to decide if it should wait for a client response before resuming a Fiber.
*
* @param callable(Uuid $sessionId): array<int, array<string, mixed>> $provider
* @param callable(Uuid $sessionId, McpFiber $fiber): array<int, array<string, mixed>> $provider
*/
public function setPendingRequestsProvider(callable $provider): void;

Expand All @@ -118,7 +118,7 @@ public function setResponseFinder(callable $finder): void;
* The transport calls this to let the Protocol handle new requests/notifications
* that are yielded from a Fiber's execution.
*
* @param callable(FiberSuspend|null, ?Uuid $sessionId): void $handler
* @param callable(FiberSuspend|null, ?Uuid $sessionId, McpFiber $fiber): void $handler
*/
public function setFiberYieldHandler(callable $handler): void;

Expand Down
8 changes: 6 additions & 2 deletions tests/Unit/Fixtures/PollingLoopTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -29,14 +29,18 @@ final class PollingLoopTransport extends InMemoryTransport
*/
public function getPendingRequestIds(): array
{
return array_keys($this->getPendingRequests($this->sessionId));
\assert(null !== $this->sessionFiber);

return array_keys($this->getPendingRequests($this->sessionId, $this->sessionFiber));
}

/**
* @param FiberSuspend $yielded
*/
public function yieldFromFiber(NotificationSuspension|RequestSuspension $yielded): void
{
$this->handleFiberYield($yielded, $this->sessionId);
\assert(null !== $this->sessionFiber);

$this->handleFiberYield($yielded, $this->sessionId, $this->sessionFiber);
}
}
Loading
Loading