diff --git a/CHANGELOG.md b/CHANGELOG.md index e8cf8f9b..f3ee7a5c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -36,6 +36,9 @@ All notable changes to `mcp/sdk` will be documented in this file. * 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. +* Serve both protocol eras over stdio: `StdioTransport` settles the era on the client's first request and serves `2026-07-28` requests, `subscriptions/listen` and `notifications/cancelled` on the one channel. +* [BC Break] `StatelessAwareTransportInterface` declares `setHandshakeVersions()`, so a server without the modern era names only the revisions it negotiates when refusing a `2026-07-28` request, e.g. the one set with `Builder::setProtocolVersion()`. +* Answer a bare `initialize` on a `2026-07-28`-only endpoint with `-32022` naming the served revisions, and a request without a session on the handshake leg with its id. 0.8.0 ----- diff --git a/docs/protocol-versions.md b/docs/protocol-versions.md index 994eec53..4dec7b83 100644 --- a/docs/protocol-versions.md +++ b/docs/protocol-versions.md @@ -21,6 +21,7 @@ map; the mechanics live with the task they belong to. | Change notifications | HTTP `GET` stream, `resources/subscribe` | `subscriptions/listen` | | Dispatcher | `Protocol` | `StatelessProtocol` | | HTTP entry | `StreamableHttpTransport` — the same one, for both | +| stdio entry | `StdioTransport` — the same one, settled by the client's first request | `ProtocolVersion::isModern()` tells the two apart, and `Mcp\Schema\Enum\ProtocolVersion::FIRST_MODERN_VERSION` is where the boundary sits. @@ -150,7 +151,9 @@ for a runnable version, described in [Examples](examples.md#modern-era-client). ## What was removed -Answered with `404` and `-32601` by a modern server: +Answered with `404` and `-32601` by a modern server — except a bare `initialize`, which is how +a client from before the modern era opens, and is refused with `-32022` naming the revisions +the server does speak: - `initialize`, `notifications/initialized` - `ping` diff --git a/docs/run/protocol-eras.md b/docs/run/protocol-eras.md index a5200a8c..fd9157d5 100644 --- a/docs/run/protocol-eras.md +++ b/docs/run/protocol-eras.md @@ -104,6 +104,33 @@ Both legs come from **one** builder configuration — one registry, one set of h instances, one session manager. A tool registered once is reachable from both, and a change made through one is visible to the other. +## Over stdio + +`StdioTransport` serves both eras too, but stdio carries one client per process, so the era is +settled once rather than per request: the client's **first request** decides it, by the same +body-primary rule as above. + +| Opening request | The connection | +| --- | --- | +| carries a modern revision in `params._meta` | is served by the modern dispatcher from then on | +| anything else — `initialize` above all | runs the handshake, as before `2026-07-28` | + +A request from the other era after that is refused rather than served: `initialize` on a modern +connection gets `-32022` naming the modern revisions, an enveloped request on a handshake one +gets `-32600`. That is what a client that probed, gave up waiting and fell back to the handshake +needs to learn that the server settled on the modern era after all. + +On a modern connection everything shares the one channel, and requests are served one at a +time: a request's progress and log messages are written as its handler emits them, ahead of its +result, and the next message is read once that result is out. A `subscriptions/listen` is the +long-lived exception: it stays open alongside other requests, each of its messages tagged with +the subscription id, until the client sends `notifications/cancelled` for it, since there is no +stream to close. `setSubscriptionLifetime()` does not apply here. stdio has no headers, so none +of the `Mcp-*` header rules apply. + +A server built `withoutModernEra()` refuses a modern opening with `-32022` naming the handshake +revisions, and still accepts the handshake that follows. + ## Middleware The [default middleware stack](http.md#default-middleware) runs at the edge, before the diff --git a/docs/run/server-builder.md b/docs/run/server-builder.md index c1337674..d1b29bde 100644 --- a/docs/run/server-builder.md +++ b/docs/run/server-builder.md @@ -324,7 +324,7 @@ $server = Server::builder() | `withoutInputRequiredShim()` | - | Do not fulfil an `InputRequiredResult` over a handshake-era connection | | `setCachePolicy()` | policy | Set the `ttlMs`/`cacheScope` hints on cacheable results | | `setNotificationBus()` | bus | Delivery for `subscriptions/listen` streams | -| `setSubscriptionLifetime()` | seconds | How long a subscription stream is held open (`0` = unbounded) | +| `setSubscriptionLifetime()` | seconds | How long a subscription stream is held open over HTTP (`0` = unbounded) | | `setHeaderValidator()` | enabled | Toggle the SEP-2243 standard-header check on `buildStateless()` | | `setDiscovery()` | basePath, scanDirs?, excludeDirs?, cache? | Configure attribute discovery | | `setSession()` | sessionStore?, sessionManager?, gcProbability?, gcDivisor? | Configure session management | diff --git a/docs/run/subscriptions.md b/docs/run/subscriptions.md index 4c3e989f..54e86bae 100644 --- a/docs/run/subscriptions.md +++ b/docs/run/subscriptions.md @@ -38,3 +38,6 @@ $bus->publish(new ResourceUpdatedNotification('file:///project/config.json')); `Builder::setSubscriptionLifetime()` bounds how long a stream is held before the server closes it gracefully. The real ceiling is the runtime's: under PHP-FPM a stream cannot outlive `max_execution_time`. Pass `0` for "until the client or the runtime ends it". + +Over stdio the lifetime does not apply: a stream stays open until the client sends +`notifications/cancelled` for it, see [Over stdio](protocol-eras.md#over-stdio). diff --git a/examples/server/bootstrap.php b/examples/server/bootstrap.php index 4b53e610..1c61fba0 100644 --- a/examples/server/bootstrap.php +++ b/examples/server/bootstrap.php @@ -30,10 +30,10 @@ /** * The transport every example runs on. * - * Over HTTP that is one endpoint serving both protocol eras: `StreamableHttpTransport` - * classifies each request and routes it to the lifecycle it belongs to, so every - * example here answers an `initialize` handshake and a 2026-07-28 envelope alike. - * Over stdio there is no such choice to make — that binding carries the handshake era. + * Either way it serves both protocol eras: over HTTP, `StreamableHttpTransport` + * classifies each request and routes it to the lifecycle it belongs to; over stdio, + * `StdioTransport` settles the era on the client's first request. So every example + * here answers an `initialize` handshake and a 2026-07-28 envelope alike. * * @return TransportInterface|TransportInterface */ diff --git a/src/Server.php b/src/Server.php index a19e8592..f623d86f 100644 --- a/src/Server.php +++ b/src/Server.php @@ -11,6 +11,7 @@ namespace Mcp; +use Mcp\Schema\Enum\ProtocolVersion; use Mcp\Server\Builder; use Mcp\Server\Protocol; use Mcp\Server\Stateless\StatelessProtocol; @@ -26,13 +27,15 @@ final class Server { /** - * @param StatelessProtocol|null $statelessProtocol the modern-era (SEP-2575) dispatcher, absent on a - * server that serves the handshake era alone + * @param StatelessProtocol|null $statelessProtocol the modern-era (SEP-2575) dispatcher, absent on a + * server that serves the handshake era alone + * @param non-empty-list|null $handshakeVersions revisions `initialize` negotiates, all of them by default */ public function __construct( private readonly Protocol $protocol, private readonly LoggerInterface $logger = new NullLogger(), private readonly ?StatelessProtocol $statelessProtocol = null, + private readonly ?array $handshakeVersions = null, ) { } @@ -56,9 +59,13 @@ public function run(TransportInterface $transport): mixed // The eras share the transport, not the dispatcher: a transport that // can tell them apart takes both and picks per request. One that - // cannot — stdio — carries the handshake era alone. - if (null !== $this->statelessProtocol && $transport instanceof StatelessAwareTransportInterface) { - $transport->connectStateless($this->statelessProtocol); + // cannot carries the handshake era alone. + if ($transport instanceof StatelessAwareTransportInterface) { + $transport->setHandshakeVersions($this->handshakeVersions ?? ProtocolVersion::handshakeVersions()); + + if (null !== $this->statelessProtocol) { + $transport->connectStateless($this->statelessProtocol); + } } $this->logger->info('Running server...'); diff --git a/src/Server/Builder.php b/src/Server/Builder.php index 035074c9..7b68924b 100644 --- a/src/Server/Builder.php +++ b/src/Server/Builder.php @@ -963,6 +963,7 @@ public function build(): Server $protocol, $parts['logger'], [] === $modernVersions ? null : $this->buildStateless($modernVersions), + $parts['configuration']->handshakeVersions(), ); } diff --git a/src/Server/Configuration.php b/src/Server/Configuration.php index 2f0bd7f0..07edd407 100644 --- a/src/Server/Configuration.php +++ b/src/Server/Configuration.php @@ -38,4 +38,22 @@ public function __construct( public readonly ?ProtocolVersion $protocolVersion = null, ) { } + + /** + * Versions this server is willing to negotiate over `initialize`. + * + * A configured version pins the handshake to exactly that revision. Modern + * revisions are never offered here: they have no `initialize` at all, so a + * client negotiating one could not use the connection. + * + * @return non-empty-list + */ + public function handshakeVersions(): array + { + if (null !== $this->protocolVersion && !$this->protocolVersion->isModern()) { + return [$this->protocolVersion]; + } + + return ProtocolVersion::handshakeVersions(); + } } diff --git a/src/Server/Handler/Request/InitializeHandler.php b/src/Server/Handler/Request/InitializeHandler.php index 4461ae41..af883cd3 100644 --- a/src/Server/Handler/Request/InitializeHandler.php +++ b/src/Server/Handler/Request/InitializeHandler.php @@ -84,23 +84,10 @@ private function negotiate(string $requested): ProtocolVersion } /** - * Versions this server is willing to negotiate over `initialize`. - * - * A version configured on the server pins the handshake to exactly that - * revision. Modern revisions are never offered here: they have no - * `initialize` at all, so a client that reached this handler cannot speak - * one, and answering with it would leave the connection unusable. - * * @return non-empty-list */ private function supportedVersions(): array { - $configured = $this->configuration?->protocolVersion; - - if (null !== $configured && !$configured->isModern()) { - return [$configured]; - } - - return ProtocolVersion::handshakeVersions(); + return $this->configuration?->handshakeVersions() ?? ProtocolVersion::handshakeVersions(); } } diff --git a/src/Server/Protocol.php b/src/Server/Protocol.php index c6fea38b..81cb2957 100644 --- a/src/Server/Protocol.php +++ b/src/Server/Protocol.php @@ -758,7 +758,14 @@ private function resolveSession(TransportInterface $transport, ?Uuid $sessionId, } if (!$sessionId) { - $error = Error::forInvalidRequest('A valid session id is REQUIRED for non-initialize requests.'); + // Echo the id so a client probing with `server/discover` can correlate the refusal. + $id = match (true) { + 1 !== \count($messages) => null, + $messages[0] instanceof Request => $messages[0]->getId(), + $messages[0] instanceof InvalidInputMessageException => $messages[0]->getRequestId(), + default => null, + }; + $error = Error::forInvalidRequest('A valid session id is REQUIRED for non-initialize requests.', $id); $this->sendResponse($transport, $error, null, ['status_code' => 400]); return null; diff --git a/src/Server/Stateless/StatelessProtocol.php b/src/Server/Stateless/StatelessProtocol.php index fcee2fd9..63aad19a 100644 --- a/src/Server/Stateless/StatelessProtocol.php +++ b/src/Server/Stateless/StatelessProtocol.php @@ -156,6 +156,28 @@ private function requiresTransportHeaders(): bool * @param AccessToken|null $accessToken the token the request was authorized with, if the transport authorizes */ public function handle(string $body, array $headers = [], ?AccessToken $accessToken = null): StatelessResult + { + return $this->answer($body, $headers, true, $accessToken); + } + + /** + * Answers one JSON-RPC message read off a transport without a header layer. + * + * stdio carries the request metadata inline (see the stdio binding's + * "Request Metadata"), so there is no header to require or cross-check, and + * its one channel always carries a request's notifications. A long-lived + * stream is left for the caller to pace, since it interleaves it with + * everything else arriving on that channel. + */ + public function handleInline(string $message): StatelessResult + { + return $this->answer($message, [], false); + } + + /** + * @param array $headers + */ + private function answer(string $body, array $headers, bool $headerLayer, ?AccessToken $accessToken = null): StatelessResult { try { /** @var array|null $decoded */ @@ -205,19 +227,26 @@ public function handle(string $body, array $headers = [], ?AccessToken $accessTo return StatelessResult::error(Error::forInvalidRequest('A JSON-RPC request id must be a string or a number.'), 400); } + // A pre-modern client opens with a bare initialize: name the served revisions, not the missing envelope. + if ('initialize' === $method && !isset($params['_meta'][RequestMeta::PROTOCOL_VERSION])) { + $offered = $params['protocolVersion'] ?? null; + + return StatelessResult::error(Error::forUnsupportedProtocolVersion(\is_string($offered) ? $offered : '', $this->supportedVersions, $id), 400); + } + try { $meta = RequestMeta::fromParams($params, $headers); } catch (MissingRequestMetaException $e) { return StatelessResult::error(Error::forInvalidParams($e->getMessage(), $id), 400); } - if (null !== $versionError = $this->checkVersion($meta, $headers, $id)) { + if (null !== $versionError = $this->checkVersion($meta, $headers, $id, $headerLayer)) { return $versionError; } // After the version check: a peer on the wrong revision has a more // fundamental problem than headers that disagree with its body. - if (null !== $headerError = $this->headerValidator?->validate($method, $params, $headers)) { + if ($headerLayer && null !== $headerError = $this->headerValidator?->validate($method, $params, $headers)) { return StatelessResult::error(Error::forHeaderMismatch($headerError, $id), 400); } @@ -226,7 +255,7 @@ public function handle(string $body, array $headers = [], ?AccessToken $accessTo return $this->encode($method, $id, $this->discover()); } - return $this->listen($params, $id); + return $this->listen($params, $id, paced: $headerLayer); } if (\in_array($method, self::REMOVED_METHODS, true)) { @@ -236,7 +265,7 @@ public function handle(string $body, array $headers = [], ?AccessToken $accessTo ); } - return $this->dispatch($method, $decoded, $meta, $id, self::acceptsEventStream($headers), $accessToken); + return $this->dispatch($method, $decoded, $meta, $id, !$headerLayer || self::acceptsEventStream($headers), $accessToken); } /** @@ -269,14 +298,14 @@ private function acknowledge(string $method): StatelessResult * * @param array $headers */ - private function checkVersion(RequestMeta $meta, array $headers, string|int|null $id): ?StatelessResult + private function checkVersion(RequestMeta $meta, array $headers, string|int|null $id, bool $headerLayer = true): ?StatelessResult { $headerVersion = $this->header($headers, 'MCP-Protocol-Version'); // REQUIRED on every POST. The 2025-03-26 fallback for a header-less // request exists only for servers choosing to serve pre-2025-06-18 // clients, which a modern-only endpoint is not. - if (null === $headerVersion && $this->requiresTransportHeaders()) { + if (null === $headerVersion && $headerLayer && $this->requiresTransportHeaders()) { return StatelessResult::error( Error::forHeaderMismatch( \sprintf('Missing required MCP-Protocol-Version header (_meta declares "%s").', $meta->protocolVersion), @@ -309,8 +338,11 @@ private function checkVersion(RequestMeta $meta, array $headers, string|int|null * JSON-RPC id of this request, so there is none to mint. * * @param array|null $params + * @param bool $paced whether the stream sleeps between polls itself and ends after + * the subscription lifetime, or its consumer paces it by how often + * it asks for the next frame and ends it by dropping it */ - private function listen(?array $params, string|int $id): StatelessResult + private function listen(?array $params, string|int $id, bool $paced = true): StatelessResult { $notifications = \is_array($params['notifications'] ?? null) ? $params['notifications'] : null; $agreed = NotificationFilter::fromParams($notifications)->intersect($this->configuration->capabilities); @@ -319,7 +351,7 @@ private function listen(?array $params, string|int $id): StatelessResult $bus = $this->notificationBus; $codec = $this->codec; - return StatelessResult::stream(static function () use ($agreed, $id, $lifetime, $bus, $codec): \Generator { + return StatelessResult::stream(static function () use ($agreed, $id, $lifetime, $bus, $codec, $paced): \Generator { // MUST be the first message carrying this subscription's id, and // MUST precede any notification on it. yield [ @@ -337,7 +369,8 @@ private function listen(?array $params, string|int $id): StatelessResult // The tick is not optional: PHP spots a dropped peer by writing, // and a sleeping loop would pin an FPM worker for the full lifetime. - $deadline = 0.0 >= $lifetime ? \INF : microtime(true) + $lifetime; + // An unpaced stream is ended by its consumer, so it is not bounded. + $deadline = !$paced || 0.0 >= $lifetime ? \INF : microtime(true) + $lifetime; while (microtime(true) < $deadline) { if (null !== $bus) { @@ -354,6 +387,10 @@ private function listen(?array $params, string|int $id): StatelessResult yield null; + if (!$paced) { + continue; + } + if (connection_aborted()) { return; } diff --git a/src/Server/Transport/StatelessAwareTransportInterface.php b/src/Server/Transport/StatelessAwareTransportInterface.php index 52815b3e..966ca13d 100644 --- a/src/Server/Transport/StatelessAwareTransportInterface.php +++ b/src/Server/Transport/StatelessAwareTransportInterface.php @@ -11,6 +11,7 @@ namespace Mcp\Server\Transport; +use Mcp\Schema\Enum\ProtocolVersion; use Mcp\Server\Stateless\StatelessProtocol; /** @@ -26,4 +27,12 @@ interface StatelessAwareTransportInterface { public function connectStateless(StatelessProtocol $protocol): void; + + /** + * Revisions the handshake dispatcher negotiates, named when a modern-era + * request reaches a server without the modern era. + * + * @param non-empty-list $versions + */ + public function setHandshakeVersions(array $versions): void; } diff --git a/src/Server/Transport/StdioTransport.php b/src/Server/Transport/StdioTransport.php index 92a49887..4f25d5e9 100644 --- a/src/Server/Transport/StdioTransport.php +++ b/src/Server/Transport/StdioTransport.php @@ -12,32 +12,73 @@ namespace Mcp\Server\Transport; use Mcp\Exception\InvalidArgumentException; +use Mcp\JsonRpc\MessageFactory; +use Mcp\Schema\Enum\ProtocolVersion; use Mcp\Schema\JsonRpc\Error; +use Mcp\Schema\Notification\CancelledNotification; +use Mcp\Server\Stateless\StatelessProtocol; use Mcp\Server\Transport\Stdio\RunnerControl; use Mcp\Server\Transport\Stdio\RunnerControlInterface; use Mcp\Server\Transport\Stdio\RunnerState; +use Mcp\Server\Wire\InboundClassifier; use Psr\Log\LoggerInterface; /** + * Serves one client over the standard streams, in whichever protocol era it + * opens with. + * + * The client's first request decides, once, for the life of the process: a + * request carrying the 2026-07-28 per-request `_meta` envelope opens a modern + * connection, anything else — the `initialize` handshake above all — a + * handshake-era one. A later request from the other era is refused rather than + * served, so a client that probed with `server/discover`, timed out and fell + * back to the handshake learns the connection is already modern. + * + * @see https://modelcontextprotocol.io/specification/2026-07-28/basic/versioning#backward-compatibility-with-initialization-based-versions + * * @extends BaseTransport * * @phpstan-import-type McpFiber from TransportInterface * * @author Kyrian Obikwelu */ -class StdioTransport extends BaseTransport +class StdioTransport extends BaseTransport implements StatelessAwareTransportInterface { /** * Default cap on the bytes read for a single input line. */ public const DEFAULT_MAX_LINE_BYTES = 4 * 1024 * 1024; + private const CANCELLED_NOTIFICATION = 'notifications/cancelled'; + /** Whether the current over-length line is still being drained and discarded. */ private bool $discardingLine = false; /** @var positive-int */ private readonly int $maxLineBytes; + private ?StatelessProtocol $stateless = null; + + /** @var non-empty-list|null null names every handshake revision */ + private ?array $handshakeVersions = null; + + private readonly InboundClassifier $classifier; + + /** Null until the client's first request settles the era. */ + private ?bool $modern = null; + + /** + * Modern-era answers still being written, by the id of the request they + * answer: a `subscriptions/listen` for as long as it lasts, and a request + * whose handler streams notifications until its result is in. + * + * Keyed by {@see self::streamKey()}, since PHP would fold the ids `"5"` + * and `5` into one key, and JSON-RPC tells them apart. + * + * @var array> + */ + private array $streams = []; + /** * @param resource $input * @param resource $output @@ -55,6 +96,8 @@ public function __construct( ) { parent::__construct($logger); + $this->classifier = new InboundClassifier(); + if ($maxLineBytes < 1) { throw new InvalidArgumentException(\sprintf('The maximum line size must be a positive number of bytes, got %d.', $maxLineBytes)); } @@ -62,6 +105,16 @@ public function __construct( $this->maxLineBytes = $maxLineBytes; } + public function connectStateless(StatelessProtocol $protocol): void + { + $this->stateless = $protocol; + } + + public function setHandshakeVersions(array $versions): void + { + $this->handshakeVersions = $versions; + } + public function send(string $data, array $context): void { if (isset($context['session_id'])) { @@ -79,6 +132,7 @@ public function listen(): int while (!feof($this->input) && RunnerState::RUNNING === $this->runnerControl->getState()) { $this->processInput(); $this->processFiber(); + $this->processStreams(); $this->flushOutgoingMessages(); } @@ -125,8 +179,166 @@ protected function processInput(): void $trimmedLine = trim($line); if (!empty($trimmedLine)) { - $this->handleMessage($trimmedLine, $this->sessionId); + $this->route($trimmedLine); + } + } + + /** + * Hands one message to the era it belongs to, settling the connection's + * era on the first request. + */ + private function route(string $message): void + { + $classification = $this->classifier->classify('POST', $message); + + if ($classification->isRejected()) { + \assert(null !== $classification->error); + $this->writeError($classification->error); + + return; } + + $decoded = json_decode($message, true); + // Only a request with a usable id may settle the era; anything else is left to the dispatcher. + $request = \is_array($decoded) && !array_is_list($decoded) && \is_string($decoded['method'] ?? null) + && (\is_string($decoded['id'] ?? null) || \is_int($decoded['id'] ?? null)) ? $decoded : null; + + // Only the handshake era sends requests to the client, so elsewhere a response answers nothing. + if (false !== $this->modern && \is_array($decoded) && !array_is_list($decoded) && !isset($decoded['method']) + && (\array_key_exists('result', $decoded) || \array_key_exists('error', $decoded))) { + $this->logger->warning('StdioTransport ignored a response outside the handshake era.', [ + 'id' => $decoded['id'] ?? null, + ]); + + return; + } + + if (null === $this->modern && null !== $request) { + if ($classification->modern && null === $this->stateless) { + // Handshake only: name its revisions, like the HTTP entry, and leave the era open. + $this->writeError(Error::forUnsupportedProtocolVersion((string) $classification->claimedVersion, $this->handshakeVersions ?? ProtocolVersion::handshakeVersions(), $request['id'])); + + return; + } + + if ($classification->modern && !\in_array(ProtocolVersion::tryFrom((string) $classification->claimedVersion), $this->stateless->supportedVersions(), true)) { + // Refused by the modern leg; the era stays open for a fallback to the handshake. + $this->routeModern($message, $decoded, $request); + + return; + } + + $this->modern = $classification->modern; + + $this->logger->info('StdioTransport settled the connection era.', [ + 'era' => $this->modern ? 'modern' : 'handshake', + 'opened_with' => $request['method'], + ]); + } + + if (true === $this->modern || (null === $this->modern && $classification->modern && null !== $this->stateless)) { + $this->routeModern($message, $decoded, $request); + + return; + } + + if (false === $this->modern && $classification->modern && null !== $request) { + $this->writeError(Error::forInvalidRequest('This connection opened with the "initialize" handshake; a request carrying a per-request protocol version cannot follow it.', $request['id'])); + + return; + } + + $this->handleMessage($message, $this->sessionId); + } + + /** + * @param mixed $decoded the message, decoded + * @param array|null $request the message when it is a request + */ + private function routeModern(string $message, mixed $decoded, ?array $request): void + { + \assert(null !== $this->stateless); + + // stdio has no stream to close, so this is how a client ends a subscriptions/listen. + if (\is_array($decoded) && self::CANCELLED_NOTIFICATION === ($decoded['method'] ?? null) && !isset($decoded['id'])) { + [$cancellation] = MessageFactory::make()->create($message); + + if (!$cancellation instanceof CancelledNotification) { + $this->logger->debug('StdioTransport ignored a malformed cancellation.'); + + return; + } + + $requestId = $cancellation->requestId; + + if (isset($this->streams[$key = self::streamKey($requestId)])) { + unset($this->streams[$key]); + $this->logger->debug('StdioTransport dropped a cancelled request.', ['request_id' => $requestId]); + } + + return; + } + + $result = $this->stateless->handleInline($message); + + if ($result->isEmpty()) { + return; + } + + if ($result->isStream()) { + \assert(null !== $result->frames && null !== $request); + $this->streams[self::streamKey($request['id'])] = ($result->frames)(); + + return; + } + + $this->writeLine($result->toJson()); + } + + /** + * Writes what each open stream has ready: every frame up to its next idle + * poll, so a listen stream polls once per tick and a handler's + * notifications go out as it emits them. + */ + private function processStreams(): void + { + foreach ($this->streams as $id => $frames) { + try { + while ($frames->valid()) { + $frame = $frames->current(); + + // Written before resuming: resuming runs the handler on to its next frame. + if (null !== $frame) { + $this->writeLine(json_encode($frame, \JSON_THROW_ON_ERROR | \JSON_UNESCAPED_SLASHES)); + } + + $frames->next(); + + if (null === $frame) { + break; + } + } + } catch (\Throwable $e) { + $this->logger->error('StdioTransport ended a stream that failed.', ['stream' => $id, 'exception' => $e]); + unset($this->streams[$id]); + + continue; + } + + if (!$frames->valid()) { + unset($this->streams[$id]); + } + } + } + + private static function streamKey(string|int $id): string + { + return (\is_int($id) ? 'i:' : 's:').$id; + } + + private function writeError(Error $error): void + { + $this->writeLine(json_encode($error, \JSON_THROW_ON_ERROR | \JSON_UNESCAPED_SLASHES)); } private function processFiber(): void diff --git a/src/Server/Transport/StreamableHttpTransport.php b/src/Server/Transport/StreamableHttpTransport.php index 2530cebf..d201b959 100644 --- a/src/Server/Transport/StreamableHttpTransport.php +++ b/src/Server/Transport/StreamableHttpTransport.php @@ -74,6 +74,9 @@ class StreamableHttpTransport extends BaseTransport implements StatelessAwareTra private ?StatelessProtocol $stateless = null; + /** @var non-empty-list|null null names every handshake revision */ + private ?array $handshakeVersions = null; + private ?string $immediateResponse = null; private ?int $immediateStatusCode = null; @@ -171,6 +174,11 @@ public function connectStateless(StatelessProtocol $protocol): void $this->stateless = $protocol; } + public function setHandshakeVersions(array $versions): void + { + $this->handshakeVersions = $versions; + } + public function send(string $data, array $context): void { if (isset($context['session_id'])) { @@ -493,7 +501,7 @@ private function handleModernRequest(ServerRequestInterface $request, string $bo { if (null === $this->stateless) { return $this->responder->error( - Error::forUnsupportedProtocolVersion($claimedVersion, ProtocolVersion::handshakeVersions()), + Error::forUnsupportedProtocolVersion($claimedVersion, $this->handshakeVersions ?? ProtocolVersion::handshakeVersions()), 400, ); } diff --git a/tests/Unit/Server/ProtocolTest.php b/tests/Unit/Server/ProtocolTest.php index 0a13b9e4..89788605 100644 --- a/tests/Unit/Server/ProtocolTest.php +++ b/tests/Unit/Server/ProtocolTest.php @@ -251,7 +251,9 @@ public function testNonInitializeRequestWithoutSessionIdReturnsError(): void $this->callback(static function ($data) { $decoded = json_decode($data, true); + // Echoing the id lets a probing client correlate the refusal. return isset($decoded['error']) + && 1 === ($decoded['id'] ?? null) && str_contains($decoded['error']['message'], 'session id is REQUIRED'); }), $this->callback(static function ($context) { diff --git a/tests/Unit/Server/Stateless/StatelessProtocolTest.php b/tests/Unit/Server/Stateless/StatelessProtocolTest.php index 5770b178..78c4ad4f 100644 --- a/tests/Unit/Server/Stateless/StatelessProtocolTest.php +++ b/tests/Unit/Server/Stateless/StatelessProtocolTest.php @@ -278,7 +278,6 @@ public function testEmptyCapabilitiesReportNone(): void */ public static function removedMethods(): iterable { - yield 'initialize' => ['initialize', []]; yield 'ping' => ['ping', []]; yield 'logging/setLevel' => ['logging/setLevel', ['level' => 'info']]; yield 'resources/subscribe' => ['resources/subscribe', ['uri' => 'test://static']]; @@ -298,6 +297,25 @@ public function testRemovedMethodsAreUnknown(string $method, array $params): voi $this->assertSame(Error::METHOD_NOT_FOUND, $answer['body']['error']['code']); } + #[TestDox('a handshake-era client opening with "initialize" is told which revisions are served')] + public function testInitializeNamesTheServedRevisions(): void + { + $result = self::protocol()->handle(json_encode([ + 'jsonrpc' => '2.0', + 'id' => 1, + 'method' => 'initialize', + 'params' => ['protocolVersion' => '2025-11-25', 'capabilities' => new \stdClass(), 'clientInfo' => ['name' => 'legacy', 'version' => '1.0.0']], + ], \JSON_THROW_ON_ERROR)); + + $answer = json_decode($result->toJson(), true); + + $this->assertSame(400, $result->httpStatus); + $this->assertSame(Error::UNSUPPORTED_PROTOCOL_VERSION, $answer['error']['code']); + $this->assertSame(1, $answer['id']); + $this->assertSame('2025-11-25', $answer['error']['data']['requested']); + $this->assertSame([ProtocolVersion::V2026_07_28->value], $answer['error']['data']['supported']); + } + /** * Drains a streaming result into the frames it would write. * @@ -1016,6 +1034,111 @@ public function testListenStreamDeliversSubscribedNotifications(): void $this->assertSame('complete', $frames[3]['result']['resultType']); } + /** + * A request the way stdio carries it: the metadata inline, no headers. + * + * @param array $params + */ + private static function inlineRequest(string $method, array $params = []): string + { + $params['_meta'] = [ + RequestMeta::PROTOCOL_VERSION => ProtocolVersion::V2026_07_28->value, + RequestMeta::CLIENT_CAPABILITIES => new \stdClass(), + ...($params['_meta'] ?? []), + ]; + + return json_encode(['jsonrpc' => '2.0', 'id' => 9, 'method' => $method, 'params' => $params], \JSON_THROW_ON_ERROR); + } + + #[TestDox('a message off a transport without headers is answered without them')] + public function testInlineMessageNeedsNoHeaders(): void + { + $protocol = self::protocol(); + $message = self::inlineRequest('tools/call', ['name' => 'plain_tool', 'arguments' => []]); + + // The same message over HTTP is missing headers it has to carry. + $this->assertSame(Error::HEADER_MISMATCH, json_decode($protocol->handle($message)->toJson(), true)['error']['code']); + + $result = $protocol->handleInline($message); + + $this->assertSame(200, $result->httpStatus); + $this->assertSame('ok', json_decode($result->toJson(), true)['result']['content'][0]['text']); + } + + #[TestDox('an inline request streams its progress without being asked to')] + public function testInlineProgressIsStreamed(): void + { + $result = self::protocol()->handleInline(self::inlineRequest('tools/call', [ + 'name' => 'progress_tool', + 'arguments' => [], + '_meta' => ['progressToken' => 'tok-1'], + ])); + + $this->assertTrue($result->isStream()); + + $frames = self::frames($result); + + $this->assertSame(['notifications/progress', 'notifications/progress'], [$frames[0]['method'], $frames[1]['method']]); + $this->assertSame(9, $frames[2]['id']); + } + + #[TestDox('an inline listen stream leaves the pacing to its consumer')] + public function testInlineListenIsNotPaced(): void + { + $protocol = Server::builder() + ->setServerInfo('test-server', '1.0.0') + ->setCapabilities(new ServerCapabilities(toolsListChanged: true)) + ->setNotificationBus(new InMemoryNotificationBus()) + ->setSubscriptionLifetime(60) + ->buildStateless([ProtocolVersion::V2026_07_28]); + + $result = $protocol->handleInline(self::inlineRequest('subscriptions/listen', ['notifications' => ['toolsListChanged' => true]])); + + $this->assertTrue($result->isStream()); + + $stream = $result->frames; + $this->assertNotNull($stream); + + $frames = $stream(); + $started = microtime(true); + + $this->assertSame('notifications/subscriptions/acknowledged', $frames->current()['method']); + + // A paced stream sleeps between polls, stalling the shared stdio channel. + for ($i = 0; $i < 5; ++$i) { + $frames->next(); + $this->assertNull($frames->current()); + } + + $this->assertLessThan(0.2, microtime(true) - $started); + } + + #[TestDox('an inline listen stream outlives the subscription lifetime, since its consumer ends it')] + public function testInlineListenIsNotBoundedByTheLifetime(): void + { + $protocol = Server::builder() + ->setServerInfo('test-server', '1.0.0') + ->setCapabilities(new ServerCapabilities(toolsListChanged: true)) + ->setNotificationBus(new InMemoryNotificationBus()) + ->setSubscriptionLifetime(0.01) + ->buildStateless([ProtocolVersion::V2026_07_28]); + + $result = $protocol->handleInline(self::inlineRequest('subscriptions/listen', ['notifications' => ['toolsListChanged' => true]])); + + $stream = $result->frames; + $this->assertNotNull($stream); + + $frames = $stream(); + $frames->current(); + usleep(20_000); + + for ($i = 0; $i < 3; ++$i) { + $frames->next(); + $this->assertTrue($frames->valid()); + $this->assertNull($frames->current()); + } + } + #[TestDox('the acknowledgment drops types the server cannot honour')] public function testAcknowledgmentReflectsWhatTheServerCanDo(): void { diff --git a/tests/Unit/Server/Transport/DualEraRoutingTest.php b/tests/Unit/Server/Transport/DualEraRoutingTest.php index d73d94b4..82d2bc8f 100644 --- a/tests/Unit/Server/Transport/DualEraRoutingTest.php +++ b/tests/Unit/Server/Transport/DualEraRoutingTest.php @@ -183,6 +183,21 @@ public function testHandshakeOnlyServerRefusesModernTraffic(): void $this->assertNotContains(ProtocolVersion::V2026_07_28->value, $answer['body']['error']['data']['supported']); } + #[TestDox('a server without the modern era pinned to one revision names only that one when refusing a modern claim')] + public function testPinnedHandshakeOnlyServerNamesItsRevision(): void + { + $server = Server::builder() + ->setServerInfo('dual-era-server', '1.0.0') + ->setProtocolVersion(ProtocolVersion::V2025_06_18) + ->withoutModernEra() + ->build(); + + $answer = $this->post($server, $this->enveloped('server/discover')); + + $this->assertSame(-32022, $answer['body']['error']['code']); + $this->assertSame([ProtocolVersion::V2025_06_18->value], $answer['body']['error']['data']['supported']); + } + #[TestDox('a DELETE still ends a handshake-era session')] public function testDeleteReachesTheHandshakeLeg(): void { diff --git a/tests/Unit/Server/Transport/StdioDualEraTest.php b/tests/Unit/Server/Transport/StdioDualEraTest.php new file mode 100644 index 00000000..04c5bfb2 --- /dev/null +++ b/tests/Unit/Server/Transport/StdioDualEraTest.php @@ -0,0 +1,440 @@ +serve(self::builder(), [ + self::modern(1, 'server/discover'), + self::modern(2, 'tools/call', ['name' => 'echo', 'arguments' => ['text' => 'hi']]), + ]); + + $this->assertSame([ProtocolVersion::V2026_07_28->value], $answers[1]['result']['supportedVersions']); + $this->assertSame('test-server', $answers[1]['result']['_meta'][RequestMeta::SERVER_INFO]['name']); + $this->assertSame('hi', $answers[2]['result']['content'][0]['text']); + } + + #[TestDox('a client opening with the handshake is served the handshake era')] + public function testHandshakeOpening(): void + { + $answers = $this->serve(self::builder(), [ + self::initialize(1), + ['jsonrpc' => '2.0', 'method' => 'notifications/initialized'], + ['jsonrpc' => '2.0', 'id' => 2, 'method' => 'tools/call', 'params' => ['name' => 'echo', 'arguments' => ['text' => 'hi']]], + ]); + + $this->assertSame(ProtocolVersion::V2025_11_25->value, $answers[1]['result']['protocolVersion']); + $this->assertSame('hi', $answers[2]['result']['content'][0]['text']); + } + + #[TestDox('a handshake after a modern opening is refused, naming the modern revisions')] + public function testHandshakeAfterModernOpeningIsRefused(): void + { + $answers = $this->serve(self::builder(), [ + self::modern(1, 'server/discover'), + self::initialize(2), + ]); + + $this->assertSame(Error::UNSUPPORTED_PROTOCOL_VERSION, $answers[2]['error']['code']); + $this->assertSame([ProtocolVersion::V2026_07_28->value], $answers[2]['error']['data']['supported']); + } + + #[TestDox('a modern request after a handshake opening is refused')] + public function testModernRequestAfterHandshakeOpeningIsRefused(): void + { + $answers = $this->serve(self::builder(), [ + self::initialize(1), + self::modern(2, 'server/discover'), + ]); + + $this->assertSame(ProtocolVersion::V2025_11_25->value, $answers[1]['result']['protocolVersion']); + $this->assertSame(Error::INVALID_REQUEST, $answers[2]['error']['code']); + } + + #[TestDox('a handshake-only server refuses a modern probe and still accepts the handshake')] + public function testHandshakeOnlyServerRefusesTheProbe(): void + { + $answers = $this->serve(self::builder()->withoutModernEra(), [ + self::modern(1, 'server/discover'), + self::initialize(2), + ]); + + $this->assertSame(Error::UNSUPPORTED_PROTOCOL_VERSION, $answers[1]['error']['code']); + $this->assertNotContains(ProtocolVersion::V2026_07_28->value, $answers[1]['error']['data']['supported']); + $this->assertSame(ProtocolVersion::V2025_11_25->value, $answers[2]['result']['protocolVersion']); + } + + #[TestDox('a handshake-only server pinned to one revision names only that one when refusing a modern probe')] + public function testPinnedHandshakeOnlyServerNamesItsRevision(): void + { + $answers = $this->serve(self::builder()->withoutModernEra()->setProtocolVersion(ProtocolVersion::V2025_06_18), [ + self::modern(1, 'server/discover'), + ]); + + $this->assertSame([ProtocolVersion::V2025_06_18->value], $answers[1]['error']['data']['supported']); + } + + #[TestDox('a probe claiming an unserved revision is refused, naming the served ones, and the handshake still follows')] + public function testUnservedModernProbeLeavesTheEraOpenForTheHandshake(): void + { + $answers = $this->serve(self::builder(), [ + self::modern(1, 'server/discover', version: '2027-01-01'), + self::initialize(2), + ]); + + $this->assertSame(Error::UNSUPPORTED_PROTOCOL_VERSION, $answers[1]['error']['code']); + $this->assertSame([ProtocolVersion::V2026_07_28->value], $answers[1]['error']['data']['supported']); + $this->assertSame(ProtocolVersion::V2025_11_25->value, $answers[2]['result']['protocolVersion']); + } + + #[TestDox('a probe claiming an unserved revision leaves the modern era open to a served one')] + public function testUnservedModernProbeLeavesTheEraOpenForAServedRevision(): void + { + $answers = $this->serve(self::builder(), [ + self::modern(1, 'server/discover', version: '2027-01-01'), + self::modern(2, 'tools/call', ['name' => 'echo', 'arguments' => ['text' => 'hi']]), + ]); + + $this->assertSame(Error::UNSUPPORTED_PROTOCOL_VERSION, $answers[1]['error']['code']); + $this->assertSame('hi', $answers[2]['result']['content'][0]['text']); + } + + /** + * @return iterable + */ + public static function provideServers(): iterable + { + yield 'a server speaking both eras' => [self::builder()]; + yield 'a server without the modern era' => [self::builder()->withoutModernEra()]; + } + + #[TestDox('a request with an id that is neither a string nor a number is refused without taking the server down, on $_dataName')] + #[DataProvider('provideServers')] + public function testMalformedIdDoesNotStopTheServer(Builder $builder): void + { + $malformed = self::modern(1, 'server/discover'); + $malformed['id'] = []; + + $lines = $this->exchange($builder, [$malformed, self::initialize(2)]); + + $this->assertSame(Error::INVALID_REQUEST, $lines[0]['error']['code'] ?? null); + $this->assertSame(2, $lines[1]['id'] ?? null); + $this->assertSame(ProtocolVersion::V2025_11_25->value, $lines[1]['result']['protocolVersion'] ?? null); + } + + #[TestDox('a probe without an envelope gets an error it can correlate, not silence')] + public function testUnenvelopedRequestBeforeTheHandshakeIsAnswered(): void + { + $answers = $this->serve(self::builder(), [ + ['jsonrpc' => '2.0', 'id' => 1, 'method' => 'server/discover'], + ]); + + $this->assertSame(Error::INVALID_REQUEST, $answers[1]['error']['code']); + } + + #[TestDox('a response is not a request, so one arriving first leaves the era open')] + public function testLeadingResponseDoesNotSettleTheEra(): void + { + $lines = $this->exchange(self::builder(), [ + ['jsonrpc' => '2.0', 'id' => 'stray', 'result' => []], + self::modern(1, 'server/discover'), + ]); + + $this->assertCount(1, $lines); + $this->assertSame([ProtocolVersion::V2026_07_28->value], $lines[0]['result']['supportedVersions']); + } + + #[TestDox('a modern connection sends no requests, so a response on it is ignored rather than refused')] + public function testResponseOnModernConnectionIsIgnored(): void + { + $lines = $this->exchange(self::builder(), [ + self::modern(1, 'server/discover'), + ['jsonrpc' => '2.0', 'id' => 'stray', 'result' => []], + self::modern(2, 'server/discover'), + ]); + + $this->assertSame([1, 2], array_column($lines, 'id')); + } + + #[TestDox('a modern request streams its progress on the shared channel before its result')] + public function testModernProgressIsStreamed(): void + { + $lines = $this->exchange(self::builder(), [ + self::modern(1, 'tools/call', ['name' => 'count', 'arguments' => [], '_meta' => ['progressToken' => 'p']]), + ]); + + $this->assertSame(['notifications/progress', 'notifications/progress'], [$lines[0]['method'], $lines[1]['method']]); + $this->assertSame(1, $lines[2]['id']); + $this->assertSame('counted', $lines[2]['result']['content'][0]['text']); + } + + #[TestDox('a modern request is served to its result before the next message is read, so a cancel for it comes too late')] + public function testRequestCompletesBeforeItsCancelIsRead(): void + { + $lines = $this->exchange(self::builder(), [ + self::modern(1, 'tools/call', ['name' => 'count', 'arguments' => [], '_meta' => ['progressToken' => 'p']]), + ['jsonrpc' => '2.0', 'method' => 'notifications/cancelled', 'params' => ['requestId' => 1]], + ]); + + $this->assertCount(3, $lines); + $this->assertSame('counted', $lines[2]['result']['content'][0]['text']); + } + + #[TestDox('a listen stream shares the channel, acknowledged and tagged with its subscription')] + public function testListenStreamIsAcknowledged(): void + { + $lines = $this->exchange(self::builder()->setCapabilities(new ServerCapabilities(toolsListChanged: true)), [ + self::modern(5, 'subscriptions/listen', ['notifications' => ['toolsListChanged' => true]]), + self::modern(6, 'tools/call', ['name' => 'echo', 'arguments' => ['text' => 'still served']]), + ['jsonrpc' => '2.0', 'method' => 'notifications/cancelled', 'params' => ['requestId' => 5]], + ]); + + $this->assertSame('notifications/subscriptions/acknowledged', $lines[0]['method']); + $this->assertSame(5, $lines[0]['params']['_meta'][RequestMeta::SUBSCRIPTION_ID]); + + // The open subscription doesn't hold up the next request; its cancel gets no answer. + $this->assertSame(6, $lines[1]['id']); + $this->assertSame('still served', $lines[1]['result']['content'][0]['text']); + $this->assertCount(2, $lines); + } + + #[TestDox('a notification published while a listen stream is open reaches the client on a later tick, tagged with its subscription')] + public function testListenStreamDeliversLaterNotifications(): void + { + $input = fopen('php://memory', 'r'); + $output = fopen('php://memory', 'r+'); + $this->assertNotFalse($input); + $this->assertNotFalse($output); + $bus = new InMemoryNotificationBus(); + + $transport = new StdioTransport($input, $output); + $transport->connectStateless(self::builder() + ->setCapabilities(new ServerCapabilities(toolsListChanged: true)) + ->setNotificationBus($bus) + ->setSubscriptionLifetime(0.01) + ->buildStateless([ProtocolVersion::V2026_07_28])); + + $tick = new \ReflectionMethod($transport, 'processStreams'); + (new \ReflectionMethod($transport, 'route'))->invoke($transport, json_encode( + self::modern(5, 'subscriptions/listen', ['notifications' => ['toolsListChanged' => true]]), + \JSON_THROW_ON_ERROR, + )); + $tick->invoke($transport); + + // Past the lifetime, which an HTTP stream would have closed on. + usleep(20_000); + $bus->publish(new ToolListChangedNotification()); + + // A tick ends past its idle poll, so what that poll finds is written on the next tick. + $tick->invoke($transport); + $tick->invoke($transport); + + rewind($output); + $lines = array_map( + static fn (string $line): array => json_decode($line, true, flags: \JSON_THROW_ON_ERROR), + array_values(array_filter(explode("\n", (string) stream_get_contents($output)))), + ); + + $this->assertSame(['notifications/subscriptions/acknowledged', 'notifications/tools/list_changed'], array_column($lines, 'method')); + $this->assertSame(5, $lines[1]['params']['_meta'][RequestMeta::SUBSCRIPTION_ID]); + $this->assertCount(2, $lines); + } + + #[TestDox('a string and an integer request id are different requests, so cancelling one leaves the other')] + public function testStreamsAreKeyedByIdType(): void + { + $input = fopen('php://temp', 'r+'); + $output = fopen('php://temp', 'r+'); + $this->assertNotFalse($input); + $this->assertNotFalse($output); + + foreach ([ + self::modern(5, 'subscriptions/listen', ['notifications' => ['toolsListChanged' => true]]), + ['jsonrpc' => '2.0', 'id' => '5', 'method' => 'subscriptions/listen', 'params' => self::modern(0, 'x', ['notifications' => ['toolsListChanged' => true]])['params']], + ['jsonrpc' => '2.0', 'method' => 'notifications/cancelled', 'params' => ['requestId' => 5]], + ] as $message) { + fwrite($input, json_encode($message, \JSON_THROW_ON_ERROR)."\n"); + } + + rewind($input); + + $transport = new StdioTransport($input, $output); + self::builder()->setCapabilities(new ServerCapabilities(toolsListChanged: true))->build()->run($transport); + + $streams = (new \ReflectionProperty($transport, 'streams'))->getValue($transport); + + $this->assertSame(['s:5'], array_keys($streams)); + } + + #[TestDox('a malformed cancellation leaves the stream it names open')] + public function testMalformedCancellationIsIgnored(): void + { + $input = fopen('php://temp', 'r+'); + $output = fopen('php://temp', 'r+'); + $this->assertNotFalse($input); + $this->assertNotFalse($output); + + foreach ([ + self::modern(5, 'subscriptions/listen', ['notifications' => ['toolsListChanged' => true]]), + ['jsonrpc' => '2.0', 'method' => 'notifications/cancelled', 'params' => ['requestId' => 5, 'reason' => []]], + ] as $message) { + fwrite($input, json_encode($message, \JSON_THROW_ON_ERROR)."\n"); + } + + rewind($input); + + $transport = new StdioTransport($input, $output); + self::builder()->setCapabilities(new ServerCapabilities(toolsListChanged: true))->build()->run($transport); + + $streams = (new \ReflectionProperty($transport, 'streams'))->getValue($transport); + + $this->assertSame(['i:5'], array_keys($streams)); + } + + #[TestDox('a frame is written before its stream is resumed, so a slow handler does not hold back its progress')] + public function testFrameIsWrittenBeforeTheStreamResumes(): void + { + $output = fopen('php://memory', 'r+'); + $input = fopen('php://memory', 'r'); + $this->assertNotFalse($output); + $this->assertNotFalse($input); + $writtenBeforeResuming = null; + + $frames = (static function () use ($output, &$writtenBeforeResuming): \Generator { + yield ['jsonrpc' => '2.0', 'method' => 'notifications/progress', 'params' => ['progressToken' => 'p', 'progress' => 1]]; + + // Where a handler would carry on with its slow work. + $writtenBeforeResuming = ftell($output) > 0; + + yield ['jsonrpc' => '2.0', 'id' => 1, 'result' => []]; + })(); + + $transport = new StdioTransport($input, $output); + (new \ReflectionProperty($transport, 'streams'))->setValue($transport, ['i:1' => $frames]); + (new \ReflectionMethod($transport, 'processStreams'))->invoke($transport); + + $this->assertTrue($writtenBeforeResuming); + } + + private static function builder(): Builder + { + return Server::builder() + ->setServerInfo('test-server', '1.0.0') + ->addTool(static fn (string $text): string => $text, name: 'echo', description: 'Echoes') + ->addTool(static function (RequestContext $context): string { + $context->getClientGateway()->progress(1, 2); + $context->getClientGateway()->progress(2, 2); + + return 'counted'; + }, name: 'count', description: 'Reports progress'); + } + + /** + * @param array $params + * + * @return array + */ + private static function modern(int $id, string $method, array $params = [], ?string $version = null): array + { + $params['_meta'] = [ + RequestMeta::PROTOCOL_VERSION => $version ?? ProtocolVersion::V2026_07_28->value, + RequestMeta::CLIENT_CAPABILITIES => new \stdClass(), + ...($params['_meta'] ?? []), + ]; + + return ['jsonrpc' => '2.0', 'id' => $id, 'method' => $method, 'params' => $params]; + } + + /** + * @return array + */ + private static function initialize(int $id): array + { + return ['jsonrpc' => '2.0', 'id' => $id, 'method' => 'initialize', 'params' => [ + 'protocolVersion' => ProtocolVersion::V2025_11_25->value, + 'capabilities' => new \stdClass(), + 'clientInfo' => ['name' => 'test-client', 'version' => '1.0.0'], + ]]; + } + + /** + * Runs the server over the given input until it is exhausted, and returns + * every answer by the id it answers. + * + * @param list> $messages + * + * @return array> + */ + private function serve(Builder $builder, array $messages): array + { + $answers = []; + + foreach ($this->exchange($builder, $messages) as $line) { + if (\array_key_exists('id', $line)) { + $answers[$line['id']] = $line; + } + } + + return $answers; + } + + /** + * @param list> $messages + * + * @return list> + */ + private function exchange(Builder $builder, array $messages): array + { + $input = fopen('php://temp', 'r+'); + $this->assertNotFalse($input); + // A file rather than memory: running the server closes its streams. + $outputFile = tempnam(sys_get_temp_dir(), 'mcp-stdio'); + $output = fopen($outputFile, 'w'); + $this->assertNotFalse($output); + + foreach ($messages as $message) { + fwrite($input, json_encode($message, \JSON_THROW_ON_ERROR)."\n"); + } + + rewind($input); + + $builder->build()->run(new StdioTransport($input, $output)); + + $written = (string) file_get_contents($outputFile); + unlink($outputFile); + + return array_map( + static fn (string $line): array => json_decode($line, true, flags: \JSON_THROW_ON_ERROR), + array_values(array_filter(explode("\n", $written))), + ); + } +}