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
13 changes: 11 additions & 2 deletions packages/client/src/Client/Internal/AmphpHttpClient.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
Expand All @@ -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);
}
},
);
}

Expand Down
32 changes: 32 additions & 0 deletions packages/client/src/Client/Internal/CancellationError.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc\Client\Internal;

use Amp\CancelledException;
use Amp\TimeoutException;
use Thesis\Google\Rpc\Code;
use Thesis\Grpc\InvokeError;

/**
* Reports a call the caller stopped waiting for the way grpc-go does: an expired deadline
* (a {@see TimeoutException} behind the cancellation, e.g. from {@see \Amp\TimeoutCancellation})
* is DEADLINE_EXCEEDED, any other cancellation is CANCELLED. The original exception is kept
* as the previous one.
*
* @internal
*/
final class CancellationError
{
public static function from(CancelledException $cancelled): InvokeError
{
for ($cause = $cancelled->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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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]
Expand All @@ -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;
}
Expand Down
43 changes: 43 additions & 0 deletions tests/Client/Internal/CancellationErrorTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc\Client\Internal;

use Amp\CancelledException;
use Amp\TimeoutException;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\TestCase;
use Thesis\Google\Rpc\Code;

#[CoversClass(CancellationError::class)]
final class CancellationErrorTest extends TestCase
{
public function testExpiredDeadlineIsDeadlineExceeded(): void
{
$cancelled = new CancelledException(new TimeoutException('Too slow'));

$error = CancellationError::from($cancelled);

self::assertSame(Code::DEADLINE_EXCEEDED, $error->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);
}
}
193 changes: 193 additions & 0 deletions tests/ClientCancellationTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,193 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc;

use Amp\Cancellation;
use Amp\CancelledException;
use Amp\DeferredCancellation;
use Amp\TimeoutCancellation;
use Echos\Api\V1\EchoRequest;
use Echos\Api\V1\EchoResponse;
use Echos\Api\V1\EchoServiceClient;
use Echos\Api\V1\EchoServiceServer;
use Echos\Api\V1\EchoServiceServerRegistry;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\TestCase;
use Revolt\EventLoop;
use Thesis\Google\Protobuf\Timestamp;
use Thesis\Google\Rpc\Code;
use Thesis\Grpc\Client\Address;
use Thesis\Grpc\Client\Endpoint;
use Thesis\Grpc\Client\EndpointResolver;
use Thesis\Grpc\Client\EndpointResolverListener;
use Thesis\Grpc\Client\Internal\AmphpHttpClient;
use Thesis\Grpc\Client\Internal\CancellationError;
use Thesis\Grpc\Client\Internal\Http2\ConcurrentClientStream;
use Thesis\Grpc\Client\Resolution;
use Thesis\Grpc\Client\Target;
use Topic\Api\V1\Event;
use Topic\Api\V1\SubscribeRequest;
use Topic\Api\V1\TopicServiceClient;
use Topic\Api\V1\TopicServiceServer;
use Topic\Api\V1\TopicServiceServerRegistry;
use function Amp\delay;

#[CoversClass(AmphpHttpClient::class)]
#[CoversClass(ConcurrentClientStream::class)]
#[CoversClass(CancellationError::class)]
final class ClientCancellationTest extends TestCase
{
private const string ADDRESS = '127.0.0.1:50072';

private Server $server;

protected function setUp(): void
{
$this->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))]);
}
}
Loading