Skip to content
Merged
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
16 changes: 4 additions & 12 deletions src/Server/ClientGateway.php
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,8 @@
use Mcp\Schema\Tool;
use Mcp\Schema\ToolChoice;
use Mcp\Server\Session\SessionInterface;
use Mcp\Server\Suspension\NotificationSuspension;
use Mcp\Server\Suspension\RequestSuspension;

/**
* @final
Expand Down Expand Up @@ -91,11 +93,7 @@ public function __construct(
*/
public function notify(Notification $notification): void
{
\Fiber::suspend([
'type' => 'notification',
'notification' => $notification,
'session_id' => $this->session->getId()->toRfc4122(),
]);
\Fiber::suspend(new NotificationSuspension($notification, $this->session->getId()->toRfc4122()));
}

/**
Expand Down Expand Up @@ -450,13 +448,7 @@ public function request(Request $request, int $timeout = 120): Response|Error
*/
private function suspend(Request $request, int $timeout, ?string $key = null): Response|Error
{
$response = \Fiber::suspend([
'type' => 'request',
'request' => $request,
'session_id' => $this->session->getId()->toRfc4122(),
'timeout' => $timeout,
'input_key' => $key,
]);
$response = \Fiber::suspend(new RequestSuspension($request, $this->session->getId()->toRfc4122(), $timeout, $key));

if (!$response instanceof Response && !$response instanceof Error) {
throw new RuntimeException('Transport returned an unexpected payload; expected a Response or Error message.');
Expand Down
56 changes: 13 additions & 43 deletions src/Server/Protocol.php
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@
use Mcp\Server\Session\SessionManagerInterface;
use Mcp\Server\Stateless\InputContext;
use Mcp\Server\Stateless\RequestStateCodec;
use Mcp\Server\Suspension\NotificationSuspension;
use Mcp\Server\Suspension\RequestSuspension;
use Mcp\Server\Transport\TransportInterface;
use Psr\EventDispatcher\EventDispatcherInterface;
use Psr\Log\LoggerInterface;
Expand Down Expand Up @@ -322,17 +324,10 @@ private function handleRequest(TransportInterface $transport, Request $request,
$result = $fiber->start();

if ($fiber->isSuspended()) {
if (\is_array($result) && isset($result['type'])) {
if ('notification' === $result['type']) {
$notification = $result['notification'];
$this->sendNotification($notification, $session);
} elseif ('request' === $result['type']) {
// Keep $request untouched: it is the inbound request the catch
// blocks below answer under, not this outbound one.
$outboundRequest = $result['request'];
$timeout = $result['timeout'] ?? 120;
$this->sendRequest($outboundRequest, $timeout, $session);
}
if ($result instanceof NotificationSuspension) {
$this->sendNotification($result->notification, $session);
} elseif ($result instanceof RequestSuspension) {
$this->sendRequest($result->request, $result->timeout, $session);
Comment thread
chr-hertel marked this conversation as resolved.
}

$transport->attachFiberToSession($fiber, $session->getId());
Expand Down Expand Up @@ -618,7 +613,7 @@ public function handleFiberYield(mixed $yieldedValue, ?Uuid $sessionId): void
return;
}

if (!\is_array($yieldedValue) || !isset($yieldedValue['type'])) {
if (!$yieldedValue instanceof NotificationSuspension && !$yieldedValue instanceof RequestSuspension) {
$this->logger->warning('Fiber yielded unexpected payload.', [
'payload' => $yieldedValue,
'session_id' => $sessionId->toRfc4122(),
Expand All @@ -629,43 +624,18 @@ public function handleFiberYield(mixed $yieldedValue, ?Uuid $sessionId): void

$session = $this->sessionManager->createWithId($sessionId);

$payloadSessionId = $yieldedValue['session_id'] ?? null;
if (\is_string($payloadSessionId) && $payloadSessionId !== $sessionId->toRfc4122()) {
if ($yieldedValue->sessionId !== $sessionId->toRfc4122()) {
$this->logger->warning('Fiber yielded payload with mismatched session ID.', [
'payload_session_id' => $payloadSessionId,
'payload_session_id' => $yieldedValue->sessionId,
'expected_session_id' => $sessionId->toRfc4122(),
]);
}

try {
if ('notification' === $yieldedValue['type']) {
$notification = $yieldedValue['notification'] ?? null;
if (!$notification instanceof Notification) {
$this->logger->warning('Fiber yielded notification without Notification instance.', [
'payload' => $yieldedValue,
]);

return;
}

$this->sendNotification($notification, $session);
} elseif ('request' === $yieldedValue['type']) {
$request = $yieldedValue['request'] ?? null;
if (!$request instanceof Request) {
$this->logger->warning('Fiber yielded request without Request instance.', [
'payload' => $yieldedValue,
]);

return;
}

$timeout = isset($yieldedValue['timeout']) ? (int) $yieldedValue['timeout'] : 120;
$this->sendRequest($request, $timeout, $session);
} else {
$this->logger->warning('Fiber yielded unknown operation type.', [
'type' => $yieldedValue['type'],
]);
}
match (true) {
$yieldedValue instanceof NotificationSuspension => $this->sendNotification($yieldedValue->notification, $session),
$yieldedValue instanceof RequestSuspension => $this->sendRequest($yieldedValue->request, $yieldedValue->timeout, $session),
};
} finally {
$session->save();
}
Expand Down
24 changes: 7 additions & 17 deletions src/Server/Stateless/StatelessProtocol.php
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@
use Mcp\Server\Session\InMemorySessionStore;
use Mcp\Server\Session\Session;
use Mcp\Server\Subscription\NotificationBusInterface;
use Mcp\Server\Suspension\NotificationSuspension;
use Mcp\Server\Suspension\RequestSuspension;
use Mcp\Server\Wire\CachePolicy;
use Mcp\Server\Wire\InboundClassifier;
use Mcp\Server\Wire\Rev2026Codec;
Expand Down Expand Up @@ -675,8 +677,8 @@ private function run(RequestHandlerInterface $handler, Request $request, Session
*/
private function readNotification(mixed $suspended, RequestMeta $meta): ?Notification
{
if (!\is_array($suspended) || 'notification' !== ($suspended['type'] ?? null)) {
if (\is_array($suspended) && 'request' === ($suspended['type'] ?? null)) {
if (!$suspended instanceof NotificationSuspension) {
if ($suspended instanceof RequestSuspension) {
// Elicitation never reaches here — it is answered in run(). What
// is left are the kinds this revision removed outright, and no
// multi round-trip shape brings them back.
Expand All @@ -686,11 +688,7 @@ private function readNotification(mixed $suspended, RequestMeta $meta): ?Notific
return null;
}

$notification = $suspended['notification'] ?? null;

if (!$notification instanceof Notification) {
return null;
}
$notification = $suspended->notification;

// The client opts into logs per request; with no level named the server
// MUST NOT send any, which is why an absent level drops rather than
Expand All @@ -713,19 +711,11 @@ private function readNotification(mixed $suspended, RequestMeta $meta): ?Notific
*/
private static function readElicitation(mixed $suspended): ?array
{
if (!\is_array($suspended) || 'request' !== ($suspended['type'] ?? null)) {
return null;
}

$request = $suspended['request'] ?? null;

if (!$request instanceof ElicitRequest) {
if (!$suspended instanceof RequestSuspension || !$suspended->request instanceof ElicitRequest) {
return null;
}

$key = $suspended['input_key'] ?? null;

return [\is_string($key) ? $key : null, $request];
return [$suspended->inputKey, $suspended->request];
}

/**
Expand Down
28 changes: 28 additions & 0 deletions src/Server/Suspension/NotificationSuspension.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
<?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\Suspension;

use Mcp\Schema\JsonRpc\Notification;

/**
* Payload a handler fiber suspends with to send a notification to the client.
*
* @author Christopher Hertel <mail@christopher-hertel.de>
*/
final class NotificationSuspension
{
public function __construct(
public readonly Notification $notification,
public readonly string $sessionId,
) {
}
}
38 changes: 38 additions & 0 deletions src/Server/Suspension/RequestSuspension.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
<?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\Suspension;

use Mcp\Schema\JsonRpc\Request;

/**
* Payload a handler fiber suspends with to send a request to the client
* and wait for its answer.
*
* @author Christopher Hertel <mail@christopher-hertel.de>
*/
final class RequestSuspension
{
/**
* @param int $timeout maximum time to wait for the response (seconds)
* @param string|null $inputKey the name an elicitation's answer is filed under when
* the revision serving the call answers by asking
* ({@see \Mcp\Server\Stateless\ElicitationReplay});
* ignored by every leg that has a live client to ask
*/
public function __construct(
public readonly Request $request,
public readonly string $sessionId,
public readonly int $timeout = 120,
public readonly ?string $inputKey = null,
) {
}
}
5 changes: 1 addition & 4 deletions src/Server/Transport/TransportInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,7 @@
*
* @phpstan-type FiberReturn (Response<mixed>|Error)
* @phpstan-type FiberResume (FiberReturn|null)
* @phpstan-type FiberSuspend (
* array{type: 'notification', notification: \Mcp\Schema\JsonRpc\Notification}|
* array{type: 'request', request: \Mcp\Schema\JsonRpc\Request, timeout?: int}
* )
* @phpstan-type FiberSuspend (\Mcp\Server\Suspension\NotificationSuspension|\Mcp\Server\Suspension\RequestSuspension)
* @phpstan-type McpFiber \Fiber<null, FiberReturn, FiberReturn, FiberSuspend>
*
* @author Christopher Hertel <mail@christopher-hertel.de>
Expand Down
8 changes: 4 additions & 4 deletions tests/Unit/Server/ClientGatewayTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
use Mcp\Schema\Result\ListRootsResult;
use Mcp\Server\ClientGateway;
use Mcp\Server\Session\SessionInterface;
use Mcp\Server\Suspension\RequestSuspension;
use PHPUnit\Framework\TestCase;
use Symfony\Component\Uid\Uuid;

Expand Down Expand Up @@ -217,11 +218,10 @@ private function runInFiber(\Closure $call, Response|Error $response, string $ex
$fiber = new \Fiber($call);
$suspend = $fiber->start();

$this->assertIsArray($suspend);
$this->assertSame('request', $suspend['type']);
$this->assertInstanceOf($expectedRequest, $suspend['request']);
$this->assertInstanceOf(RequestSuspension::class, $suspend);
$this->assertInstanceOf($expectedRequest, $suspend->request);

$request = $suspend['request'];
$request = $suspend->request;

$fiber->resume($response);

Expand Down
4 changes: 2 additions & 2 deletions tests/Unit/Server/InputRequiredShimTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
use Mcp\Server\Session\Session;
use Mcp\Server\Session\SessionInterface;
use Mcp\Server\Stateless\InputContext;
use Mcp\Server\Suspension\RequestSuspension;
use PHPUnit\Framework\Attributes\TestDox;
use PHPUnit\Framework\TestCase;

Expand Down Expand Up @@ -149,8 +150,7 @@ private function drive(
$suspended = $fiber->start();

while ($fiber->isSuspended()) {
$this->assertIsArray($suspended);
$this->assertSame('request', $suspended['type'], 'the shim only ever suspends to send a client request');
$this->assertInstanceOf(RequestSuspension::class, $suspended, 'the shim only ever suspends to send a client request');
$this->assertNotEmpty($answers, 'the shim sent more requests than the test queued answers for');

$suspended = $fiber->resume(array_shift($answers));
Expand Down
Loading
Loading