Skip to content
Closed
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 @@ -22,6 +22,7 @@ All notable changes to `mcp/sdk` will be documented in this file.
* [BC Break] `AbstractSchemaDefinition` declares an abstract `getDefault()`, which a custom schema definition has to implement.
* Fix `Client::getPrompt()` without arguments sending `"arguments": []`, which servers validating the spec's object type (e.g. the TypeScript SDK) reject: empty arguments are now omitted.
* Add `Tool::$execution` (`ToolExecution` with a `TaskSupport` enum) and `ServerCapabilities::$tasks`, which were dropped when parsing a 2025-11-25 server's `tools/list` and `initialize` results.
* Fix a request to the client (sampling, elicitation, roots) intermittently timing out over HTTP when the server runs in several processes: the stream waiting for the answer saved the session on every poll, and could overwrite the answer another process had just stored.

0.8.0
-----
Expand Down
5 changes: 5 additions & 0 deletions src/Server/Protocol.php
Original file line number Diff line number Diff line change
Expand Up @@ -518,6 +518,11 @@ public function consumeOutgoingMessages(Uuid $sessionId): array
{
$session = $this->sessionManager->createWithId($sessionId);
$queue = $session->get(self::SESSION_OUTGOING_QUEUE, []);

if ([] === $queue) {
return [];
}

$session->set(self::SESSION_OUTGOING_QUEUE, []);
$session->save();

Expand Down
110 changes: 110 additions & 0 deletions tests/Unit/Server/ProtocolSessionRaceTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
<?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\InMemorySessionStore;
use Mcp\Server\Session\SessionManager;
use Mcp\Server\Session\SessionStoreInterface;
use Mcp\Server\Transport\TransportInterface;
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(new InMemorySessionStore());
$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->afterNextRead(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));
}
}

/**
* Runs a callback once, right after the next read, to interleave another
* worker's write with the reader's read-modify-write.
*/
final class InterleavingSessionStore implements SessionStoreInterface
{
/** @var (\Closure(): void)|null */
private ?\Closure $afterNextRead = null;

public function __construct(
private readonly SessionStoreInterface $store,
) {
}

/**
* @param \Closure(): void $callback
*/
public function afterNextRead(\Closure $callback): void
{
$this->afterNextRead = $callback;
}

public function exists(Uuid $id): bool
{
return $this->store->exists($id);
}

public function read(Uuid $id): string|false
{
$data = $this->store->read($id);

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

return $data;
}

public function write(Uuid $id, string $data): bool
{
return $this->store->write($id, $data);
}

public function destroy(Uuid $id): bool
{
return $this->store->destroy($id);
}

public function gc(): array
{
return $this->store->gc();
}
}
Loading