From 90c777e6657961251384a284f912e01409f8c9f1 Mon Sep 17 00:00:00 2001 From: kafkiansky Date: Tue, 29 Sep 2026 14:48:08 +0300 Subject: [PATCH] Fix grpc deadline --- .../src/Client/Internal/AmphpHttpClient.php | 13 +- .../src/Client/Internal/CancellationError.php | 32 +++ .../Internal/Http2/ConcurrentClientStream.php | 30 ++- .../Client/Internal/CancellationErrorTest.php | 43 ++++ tests/ClientCancellationTest.php | 193 ++++++++++++++++++ 5 files changed, 302 insertions(+), 9 deletions(-) create mode 100644 packages/client/src/Client/Internal/CancellationError.php create mode 100644 tests/Client/Internal/CancellationErrorTest.php create mode 100644 tests/ClientCancellationTest.php diff --git a/packages/client/src/Client/Internal/AmphpHttpClient.php b/packages/client/src/Client/Internal/AmphpHttpClient.php index 751bf5c..6433ef8 100644 --- a/packages/client/src/Client/Internal/AmphpHttpClient.php +++ b/packages/client/src/Client/Internal/AmphpHttpClient.php @@ -5,6 +5,7 @@ namespace Thesis\Grpc\Client\Internal; use Amp\Cancellation; +use Amp\CancelledException; use Amp\NullCancellation; use Thesis\Google\Rpc\Code; use Thesis\Grpc\Client; @@ -53,6 +54,8 @@ function ( return $stream->receive(); } catch (GrpcException $e) { throw $e; + } catch (CancelledException $e) { + throw CancellationError::from($e); } catch (\Throwable $e) { // Transport-level failures (e.g. a refused connection) map to UNAVAILABLE, // so interceptors above see a gRPC status rather than a raw amphp exception. @@ -74,11 +77,17 @@ public function createStream( $invoke, $md, $cancellation, - fn( + function ( Client\Invoke $invoke, Metadata $md, Cancellation $cancellation, - ): ClientStream => $this->connection->createStream($invoke, $md, $cancellation, $pick), + ) use ($pick): ClientStream { + try { + return $this->connection->createStream($invoke, $md, $cancellation, $pick); + } catch (CancelledException $e) { + throw CancellationError::from($e); + } + }, ); } diff --git a/packages/client/src/Client/Internal/CancellationError.php b/packages/client/src/Client/Internal/CancellationError.php new file mode 100644 index 0000000..d2ee424 --- /dev/null +++ b/packages/client/src/Client/Internal/CancellationError.php @@ -0,0 +1,32 @@ +getPrevious(); $cause !== null; $cause = $cause->getPrevious()) { + if ($cause instanceof TimeoutException) { + return new InvokeError(Code::DEADLINE_EXCEEDED, $cause->getMessage(), previous: $cancelled); + } + } + + return new InvokeError(Code::CANCELLED, $cancelled->getMessage(), previous: $cancelled); + } +} diff --git a/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php b/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php index 1c30643..7bd4e99 100644 --- a/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php +++ b/packages/client/src/Client/Internal/Http2/ConcurrentClientStream.php @@ -5,11 +5,13 @@ namespace Thesis\Grpc\Client\Internal\Http2; use Amp\Cancellation; +use Amp\CancelledException; use Amp\Future; use Amp\Http\Client\Response; use Amp\NullCancellation; use Amp\Pipeline; use Thesis\Google\Rpc\Code; +use Thesis\Grpc\Client\Internal\CancellationError; use Thesis\Grpc\ClientStream; use Thesis\Grpc\Exception\ClientStreamIsClosed; use Thesis\Grpc\InvokeError; @@ -58,17 +60,25 @@ public function send(object $message): void #[\Override] public function receive(): object { - if (!$this->recv->continue()) { - throw $this->errors->obtain($this) ?? new InvokeError(Code::UNKNOWN); - } + try { + if (!$this->recv->continue()) { + throw $this->errors->obtain($this) ?? new InvokeError(Code::UNKNOWN); + } - return $this->recv->getValue(); + return $this->recv->getValue(); + } catch (CancelledException $e) { + throw CancellationError::from($e); + } } #[\Override] public function headers(): Metadata { - return new Metadata($this->response->getHeaders()); + try { + return new Metadata($this->response->getHeaders()); + } catch (CancelledException $e) { + throw CancellationError::from($e); + } } #[\Override] @@ -91,9 +101,15 @@ public function close(): void #[\Override] public function getIterator(): \Traversable { - yield from $this->recv; + try { + yield from $this->recv; + + // A cancelled call ends the message stream quietly and surfaces here, while reading the trailers. + $error = $this->errors->obtain($this); + } catch (CancelledException $e) { + throw CancellationError::from($e); + } - $error = $this->errors->obtain($this); if ($error !== null) { throw $error; } diff --git a/tests/Client/Internal/CancellationErrorTest.php b/tests/Client/Internal/CancellationErrorTest.php new file mode 100644 index 0000000..9124c4a --- /dev/null +++ b/tests/Client/Internal/CancellationErrorTest.php @@ -0,0 +1,43 @@ +statusCode); + self::assertSame('Too slow', $error->statusMessage); + self::assertSame($cancelled, $error->getPrevious()); + } + + public function testOtherCancellationIsCancelled(): void + { + $cancelled = new CancelledException(); + + $error = CancellationError::from($cancelled); + + self::assertSame(Code::CANCELLED, $error->statusCode); + self::assertSame($cancelled, $error->getPrevious()); + } + + public function testFindsTheDeadlineDeeperInTheChain(): void + { + $cancelled = new CancelledException(new CancelledException(new TimeoutException())); + + self::assertSame(Code::DEADLINE_EXCEEDED, CancellationError::from($cancelled)->statusCode); + } +} diff --git a/tests/ClientCancellationTest.php b/tests/ClientCancellationTest.php new file mode 100644 index 0000000..0520bee --- /dev/null +++ b/tests/ClientCancellationTest.php @@ -0,0 +1,193 @@ +server = new Server\Builder() + ->withAddresses(self::ADDRESS) + ->withServices( + new EchoServiceServerRegistry(new SlowEchoServer()), + new TopicServiceServerRegistry(new SlowTopicServer()), + ) + ->build(); + + $this->server->start(); + } + + protected function tearDown(): void + { + $this->server->stop(); + } + + public function testUnaryDeadlineIsDeadlineExceeded(): void + { + $error = self::catch(static fn() => self::echoClient()->echo(new EchoRequest('ping'), cancellation: new TimeoutCancellation(0.1))); + + self::assertSame(Code::DEADLINE_EXCEEDED, $error->statusCode); + self::assertInstanceOf(CancelledException::class, $error->getPrevious()); + } + + public function testUnaryCancellationIsCancelled(): void + { + $cancellation = new DeferredCancellation(); + EventLoop::delay(0.1, static fn() => $cancellation->cancel()); + + $error = self::catch(static fn() => self::echoClient()->echo(new EchoRequest('ping'), cancellation: $cancellation->getCancellation())); + + self::assertSame(Code::CANCELLED, $error->statusCode); + } + + public function testDeadlineBeforeTheConnectionIsEstablished(): void + { + $client = new EchoServiceClient( + new Client\Builder() + ->withHost('slow:///echo') + ->withEndpointResolver('slow', new SlowResolver(self::ADDRESS)) + ->build(), + ); + + $error = self::catch(static fn() => $client->echo(new EchoRequest('ping'), cancellation: new TimeoutCancellation(0.05))); + + self::assertSame(Code::DEADLINE_EXCEEDED, $error->statusCode); + } + + public function testStreamDeadlineIsDeadlineExceeded(): void + { + $stream = self::topicClient()->subscribe(new SubscribeRequest('events'), cancellation: new TimeoutCancellation(0.2)); + $stream->receive(); + + $error = self::catch(static fn() => $stream->receive()); + + self::assertSame(Code::DEADLINE_EXCEEDED, $error->statusCode); + } + + public function testStreamCancellationIsCancelled(): void + { + $cancellation = new DeferredCancellation(); + $stream = self::topicClient()->subscribe(new SubscribeRequest('events'), cancellation: $cancellation->getCancellation()); + $stream->receive(); + EventLoop::delay(0.1, static fn() => $cancellation->cancel()); + + $error = self::catch(static fn() => $stream->receive()); + + self::assertSame(Code::CANCELLED, $error->statusCode); + } + + public function testStreamIterationDeadlineIsDeadlineExceeded(): void + { + $stream = self::topicClient()->subscribe(new SubscribeRequest('events'), cancellation: new TimeoutCancellation(0.2)); + + $error = self::catch(static function () use ($stream): void { + foreach ($stream as $_); + + }); + + self::assertSame(Code::DEADLINE_EXCEEDED, $error->statusCode); + } + + /** + * @param \Closure(): mixed $call + */ + private static function catch(\Closure $call): InvokeError + { + try { + $call(); + } catch (InvokeError $e) { + return $e; + } + + self::fail('Expected an InvokeError.'); + } + + private static function echoClient(): EchoServiceClient + { + return new EchoServiceClient(new Client\Builder()->withHost(self::ADDRESS)->build()); + } + + private static function topicClient(): TopicServiceClient + { + return new TopicServiceClient(new Client\Builder()->withHost(self::ADDRESS)->build()); + } +} + +final readonly class SlowEchoServer implements EchoServiceServer +{ + #[\Override] + public function echo(EchoRequest $request, Metadata $md, Cancellation $cancellation): EchoResponse + { + delay(1, cancellation: $cancellation); + + return new EchoResponse($request->sentence); + } +} + +final readonly class SlowTopicServer implements TopicServiceServer +{ + #[\Override] + public function subscribe(SubscribeRequest $request, Metadata $md, Cancellation $cancellation): iterable + { + yield new Event('first', '', new Timestamp(1)); + delay(1, cancellation: $cancellation); + yield new Event('second', '', new Timestamp(2)); + } +} + +final readonly class SlowResolver implements EndpointResolver +{ + /** + * @param non-empty-string $address + */ + public function __construct( + private string $address, + ) {} + + #[\Override] + public function resolve(Target $target, EndpointResolverListener $listener, Cancellation $cancellation): Resolution + { + delay(1, cancellation: $cancellation); + + return new Resolution([new Endpoint(new Address($this->address))]); + } +}