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 @@ -34,6 +34,7 @@ All notable changes to `mcp/sdk` will be documented in this file.
* [BC Break] `ProtectedResourceMetadata` requires `$resource`, serves at the path derived from it (RFC 9728 §3.1) and requires https except for loopback hosts; drops localized, policy, ToS, extra fields and `$metadataPaths`.
* [BC Break] Add `ScopePolicy` as third argument of `AuthorizationMiddleware`, answering `403 insufficient_scope` per method and tool, with scope hierarchies; the `resource_metadata` challenge URL comes from the configured resource instead of the `Host` header.
* Expose `WWW-Authenticate` in the default `CorsMiddleware`.
* Fix lost responses on concurrent requests of one session over Streamable HTTP: a POST is answered with its own responses instead of taking them from the session's outgoing queue. Adds `InlineResponseTransportInterface` for transports that answer each request on the exchange that carried it.

0.8.0
-----
Expand Down
18 changes: 15 additions & 3 deletions src/Server/Protocol.php
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
use Mcp\Server\Session\SessionManagerInterface;
use Mcp\Server\Stateless\InputContext;
use Mcp\Server\Stateless\RequestStateCodec;
use Mcp\Server\Transport\InlineResponseTransportInterface;
use Mcp\Server\Transport\TransportInterface;
use Psr\EventDispatcher\EventDispatcherInterface;
use Psr\Log\LoggerInterface;
Expand Down Expand Up @@ -473,7 +474,10 @@ public function sendNotification(Notification $notification, SessionInterface $s
*/
private function sendResponse(TransportInterface $transport, Response|Error $response, ?SessionInterface $session, array $context = []): void
{
if (null === $session) {
// Queued in the session, a response can be overwritten or taken by a concurrent
// request of the same session: a transport that can answer on the request's
// own exchange gets it directly.
if (null === $session || $transport instanceof InlineResponseTransportInterface) {
$this->logger->info('Sending immediate response', [
'response_id' => $response->getId(),
]);
Expand All @@ -496,6 +500,10 @@ private function sendResponse(TransportInterface $transport, Response|Error $res
}

$context['type'] = 'response';
if (null !== $session) {
$context['session_id'] = $session->getId();
}

$transport->send($encoded, $context);
} else {
$this->logger->info('Queueing server response', [
Expand Down Expand Up @@ -541,8 +549,12 @@ public function consumeOutgoingMessages(Uuid $sessionId): array
{
$session = $this->sessionManager->createWithId($sessionId);
$queue = $session->get(self::SESSION_OUTGOING_QUEUE, []);
$session->set(self::SESSION_OUTGOING_QUEUE, []);
$session->save();

// Saving an unchanged session would only overwrite what a concurrent request saved in the meantime.
if ([] !== $queue) {
$session->set(self::SESSION_OUTGOING_QUEUE, []);
$session->save();
}

return $queue;
}
Expand Down
27 changes: 27 additions & 0 deletions src/Server/Transport/InlineResponseTransportInterface.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
<?php

/*
* This file is part of the official PHP MCP SDK.
*
* A collaboration between Symfony and the PHP Foundation.
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

namespace Mcp\Server\Transport;

/**
* A transport that answers each request on the exchange that carried it.
*
* {@see \Mcp\Server\Protocol} hands such a transport its responses through
* {@see TransportInterface::send()}, with the session in the `session_id`
* context key, instead of queueing them in the session. The session is shared
* by the concurrent requests of a client and written back whole, so a queued
* response can be overwritten by another request, or taken by it.
*
* Server-initiated requests and notifications still go through the session queue.
*/
interface InlineResponseTransportInterface
{
}
32 changes: 27 additions & 5 deletions src/Server/Transport/StreamableHttpTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@
*
* @author Kyrian Obikwelu <koshnawaza@gmail.com>
*/
class StreamableHttpTransport extends BaseTransport implements StatelessAwareTransportInterface
class StreamableHttpTransport extends BaseTransport implements StatelessAwareTransportInterface, InlineResponseTransportInterface
{
use ReadsBoundedBody;

Expand All @@ -77,6 +77,9 @@ class StreamableHttpTransport extends BaseTransport implements StatelessAwareTra
private ?string $immediateResponse = null;
private ?int $immediateStatusCode = null;

/** @var list<string> responses to the requests of the current POST, see {@see InlineResponseTransportInterface} */
private array $inlineResponses = [];

/** @var list<MiddlewareInterface>|null null until {@see self::listen()} resolves the defaults */
private ?array $middleware;

Expand Down Expand Up @@ -170,6 +173,12 @@ public function connectStateless(StatelessProtocol $protocol): void

public function send(string $data, array $context): void
{
if (isset($context['session_id'])) {
$this->inlineResponses[] = $data;

return;
}

$this->immediateResponse = $data;
$this->immediateStatusCode = $context['status_code'] ?? 200;
}
Expand Down Expand Up @@ -205,6 +214,8 @@ protected function handlePostRequest(string $body, ?AccessToken $accessToken = n
$this->immediateStatusCode = null;

if (null !== $immediateResponse) {
$this->inlineResponses = [];

return $this->responseFactory->createResponse($immediateStatusCode ?? 200)
->withHeader('Content-Type', 'application/json')
->withBody($this->streamFactory->createStream($immediateResponse));
Expand Down Expand Up @@ -232,14 +243,14 @@ protected function handleDeleteRequest(): ResponseInterface

protected function createJsonResponse(): ResponseInterface
{
$outgoingMessages = $this->getOutgoingMessages($this->sessionId);
$messages = [...array_column($this->getOutgoingMessages($this->sessionId), 'message'), ...$this->inlineResponses];
$this->inlineResponses = [];

if (empty($outgoingMessages)) {
if ([] === $messages) {
return $this->responseFactory->createResponse(202)
->withHeader('Content-Type', 'application/json');
}

$messages = array_column($outgoingMessages, 'message');
$responseBody = 1 === \count($messages) ? $messages[0] : '['.implode(',', $messages).']';

$response = $this->responseFactory->createResponse(200)
Expand All @@ -257,14 +268,25 @@ protected function createStreamedResponse(): ResponseInterface
{
$fiber = $this->sessionFiber;

$callback = function () use ($fiber): void {
// The other requests of a batch whose handler did not suspend.
$inlineResponses = $this->inlineResponses;
$this->inlineResponses = [];

$callback = function () use ($fiber, $inlineResponses): void {
if (null === $fiber) {
return;
}

try {
$this->logger->info('SSE: Starting request processing loop');

foreach ($inlineResponses as $message) {
echo "event: message\n";
echo "data: {$message}\n\n";
@ob_flush();
flush();
}

while ($fiber->isSuspended()) {
$this->flushOutgoingMessages($this->sessionId);

Expand Down
3 changes: 2 additions & 1 deletion src/Server/Transport/TransportInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,8 @@ public function listen(): mixed;
/**
* Send a message to the client immediately (bypassing session queue).
*
* Used for session resolution errors when no session is available.
* Used for session resolution errors when no session is available, and for
* every response on a {@see InlineResponseTransportInterface}.
* The transport decides HOW to send based on context.
*
* @param array<string, mixed> $context Context about this message:
Expand Down
54 changes: 54 additions & 0 deletions tests/Unit/Server/ProtocolSessionRaceTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
<?php

/*
* This file is part of the official PHP MCP SDK.
*
* A collaboration between Symfony and the PHP Foundation.
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

namespace Mcp\Tests\Unit\Server;

use Mcp\JsonRpc\MessageFactory;
use Mcp\Schema\JsonRpc\Response;
use Mcp\Server\Protocol;
use Mcp\Server\Session\SessionManager;
use Mcp\Server\Transport\TransportInterface;
use Mcp\Tests\Unit\Server\Session\Fixture\InterleavingSessionStore;
use PHPUnit\Framework\Attributes\TestDox;
use PHPUnit\Framework\TestCase;
use Symfony\Component\Uid\Uuid;

/**
* Two workers sharing one session: one holds a tool call's SSE stream open and
* polls the session for the client's answer, the other receives that answer as
* a separate POST and stores it.
*/
final class ProtocolSessionRaceTest extends TestCase
{
#[TestDox('polling the session does not overwrite a client response stored meanwhile')]
public function testPollingDoesNotOverwriteAConcurrentResponse(): void
{
$store = new InterleavingSessionStore();
$sessions = new SessionManager($store, gcProbability: 0);
$sessionId = Uuid::v4();
$sessions->createWithId($sessionId)->save();

$waiting = new Protocol([], [], MessageFactory::make(), $sessions);
$answering = new Protocol([], [], MessageFactory::make(), $sessions);
$transport = $this->createMock(TransportInterface::class);

// The answer lands right after the waiting worker read the session,
// before anything it does next could write the session back.
$store->interleaveAfterNextRead(static function () use ($answering, $transport, $sessionId): void {
$answering->processInput($transport, '{"jsonrpc": "2.0", "id": 7, "result": {"ok": true}}', $sessionId);
});

// One turn of the waiting worker's loop, with nothing queued to send.
$this->assertSame([], $waiting->consumeOutgoingMessages($sessionId));

$this->assertInstanceOf(Response::class, $waiting->checkResponse(7, $sessionId));
}
}
84 changes: 84 additions & 0 deletions tests/Unit/Server/Session/Fixture/InterleavingSessionStore.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
<?php

/*
* This file is part of the official PHP MCP SDK.
*
* A collaboration between Symfony and the PHP Foundation.
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

namespace Mcp\Tests\Unit\Server\Session\Fixture;

use Mcp\Server\Session\InMemorySessionStore;
use Symfony\Component\Uid\Uuid;

/**
* A session store that runs another request in the middle of one reading or saving its session.
*
* Replays, in one process and in a fixed order, what two PHP workers serving the
* same session do when their requests overlap.
*/
final class InterleavingSessionStore extends InMemorySessionStore
{
private ?\Closure $afterNextRead = null;
private ?\Closure $afterNextWrite = null;
private bool $readBeforeWrite = false;
private string|false|null $staleRead = null;

/**
* Runs $interleaved right after the next read, before the reader can write the session back.
*/
public function interleaveAfterNextRead(\Closure $interleaved): void
{
$this->afterNextRead = $interleaved;
}

/**
* Runs $interleaved right after the next write.
*
* With $readBeforeWrite, the interleaved request reads the session as it was
* before that write, as if it had loaded it before the first request saved.
*/
public function interleaveOnNextWrite(\Closure $interleaved, bool $readBeforeWrite = false): void
{
$this->afterNextWrite = $interleaved;
$this->readBeforeWrite = $readBeforeWrite;
}

public function read(Uuid $id): string|false
{
if (null !== $data = $this->staleRead) {
$this->staleRead = null;

return $data;
}

$data = parent::read($id);

if (null !== $interleaved = $this->afterNextRead) {
$this->afterNextRead = null;
$interleaved();
}

return $data;
}

public function write(Uuid $id, string $data): bool
{
$before = parent::read($id);
$written = parent::write($id, $data);

if (null !== $interleaved = $this->afterNextWrite) {
$this->afterNextWrite = null;
if ($this->readBeforeWrite) {
$this->staleRead = $before;
}

$interleaved();
}

return $written;
}
}
Loading
Loading