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
31 changes: 29 additions & 2 deletions packages/client/src/Client/Internal/Connection/LazyConnection.php
Original file line number Diff line number Diff line change
Expand Up @@ -47,11 +47,38 @@ public function close(Cancellation $cancellation = new NullCancellation()): void
$future = $this->future;
$this->future = null;

$future?->await($cancellation)->close($cancellation);
if ($future === null) {
return;
}

try {
$connection = $future->await($cancellation);
} catch (\Throwable $e) {
if ($future->isComplete()) {
// The connection was never established: there is nothing to close.
return;
}

throw $e;
}

$connection->close($cancellation);
}

private function createConnection(Cancellation $cancellation): Connection
{
return ($this->future ??= async($this->factory))->await($cancellation);
$future = $this->future ??= async($this->factory);

try {
return $future->await($cancellation);
} catch (\Throwable $e) {
// Forget a failed attempt so the next call retries it. A caller that merely
// stopped waiting (cancellation) leaves the attempt running for the others.
if ($future->isComplete() && $this->future === $future) {
$this->future = null;
}

throw $e;
}
}
}
158 changes: 158 additions & 0 deletions tests/Client/Internal/Connection/LazyConnectionTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,158 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc\Client\Internal\Connection;

use Amp\Cancellation;
use Amp\CancelledException;
use Amp\DeferredCancellation;
use Amp\NullCancellation;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\TestCase;
use Thesis\Grpc\Client\Internal\Connection;
use Thesis\Grpc\Client\Invoke;
use Thesis\Grpc\Client\PickContext;
use Thesis\Grpc\Metadata;
use Thesis\Grpc\RpcType;
use function Amp\async;
use function Amp\delay;
use function Amp\Future\await;

#[CoversClass(LazyConnection::class)]
final class LazyConnectionTest extends TestCase
{
public function testCreatesTheConnectionOnce(): void
{
$created = 0;
$connection = new LazyConnection(static function () use (&$created): Connection {
++$created;

return new StubConnection();
});

self::createStream($connection);
self::createStream($connection);

self::assertSame(1, $created);
}

public function testConcurrentCallersShareOneAttempt(): void
{
$created = 0;
$connection = new LazyConnection(static function () use (&$created): Connection {
++$created;
delay(0.01);

return new StubConnection();
});

await([
async(self::createStream(...), $connection),
async(self::createStream(...), $connection),
]);

self::assertSame(1, $created);
}

public function testRetriesAfterAFailedAttempt(): void
{
$attempts = 0;
$connection = new LazyConnection(static function () use (&$attempts): Connection {
if (++$attempts === 1) {
throw new \RuntimeException('Resolution failed.');
}

return new StubConnection();
});

try {
self::createStream($connection);
self::fail('Expected the first attempt to fail.');
} catch (\RuntimeException $e) {
self::assertSame('Resolution failed.', $e->getMessage());
}

self::createStream($connection);

self::assertSame(2, $attempts);
}

public function testCancelledWaitDoesNotAbandonTheAttempt(): void
{
$attempts = 0;
$connection = new LazyConnection(static function () use (&$attempts): Connection {
++$attempts;
delay(0.02);

return new StubConnection();
});

$cancellation = new DeferredCancellation();
$waiter = async(self::createStream(...), $connection, $cancellation->getCancellation());
delay(0.005);
$cancellation->cancel();

try {
$waiter->await();
self::fail('Expected the wait to be cancelled.');
} catch (CancelledException) {
}

self::createStream($connection);

self::assertSame(1, $attempts);
}

public function testClosesTheCreatedConnection(): void
{
$stub = new StubConnection();
$connection = new LazyConnection(static fn(): Connection => $stub);

self::createStream($connection);
$connection->close();

self::assertTrue($stub->closed);
}

public function testCloseIgnoresAFailedAttempt(): void
{
$connection = new LazyConnection(static function (): Connection {
delay(0.01);

throw new \RuntimeException('Resolution failed.');
});

$caller = async(self::createStream(...), $connection);
delay(0.001);

$connection->close();

$this->expectException(\RuntimeException::class);
$caller->await();
}

public function testCloseBeforeFirstUseIsANoOp(): void
{
$created = 0;
$connection = new LazyConnection(static function () use (&$created): Connection {
++$created;

return new StubConnection();
});

$connection->close();

self::assertSame(0, $created);
}

private static function createStream(LazyConnection $connection, Cancellation $cancellation = new NullCancellation()): void
{
$connection->createStream(
new Invoke('/svc/Method', \stdClass::class, RpcType::Unary),
new Metadata(),
$cancellation,
new PickContext('/svc/Method', new Metadata()),
);
}
}
37 changes: 37 additions & 0 deletions tests/Client/Internal/Connection/StubConnection.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc\Client\Internal\Connection;

use Amp\Cancellation;
use Amp\NullCancellation;
use Thesis\Grpc\Client\Internal\Connection;
use Thesis\Grpc\Client\Invoke;
use Thesis\Grpc\Client\PickContext;
use Thesis\Grpc\ClientStream;
use Thesis\Grpc\Metadata;

final class StubConnection implements Connection
{
public bool $closed = false;

/**
* @template In of object
* @template Out of object
* @param Invoke<In, Out> $invoke
* @return StubStream<In, Out>
*/
#[\Override]
public function createStream(Invoke $invoke, Metadata $md, Cancellation $cancellation, PickContext $pick): ClientStream
{
/** @var StubStream<In, Out> */
return new StubStream();
}

#[\Override]
public function close(Cancellation $cancellation = new NullCancellation()): void
{
$this->closed = true;
}
}
50 changes: 50 additions & 0 deletions tests/Client/Internal/Connection/StubStream.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
<?php

declare(strict_types=1);

namespace Thesis\Grpc\Client\Internal\Connection;

use Amp\Cancellation;
use Amp\NullCancellation;
use Thesis\Grpc\ClientStream;
use Thesis\Grpc\Metadata;

/**
* A stream that never carries messages: the connection tests only need one to exist.
*
* @template In of object
* @template-covariant Out of object
* @implements ClientStream<In, Out>
*/
final readonly class StubStream implements ClientStream
{
#[\Override]
public function send(object $message): void {}

#[\Override]
public function receive(): object
{
throw new \LogicException('The stub stream carries no messages.');
}

#[\Override]
public function headers(): Metadata
{
return new Metadata();
}

#[\Override]
public function trailers(Cancellation $cancellation = new NullCancellation()): Metadata
{
return new Metadata();
}

#[\Override]
public function close(): void {}

#[\Override]
public function getIterator(): \Traversable
{
yield from [];
}
}
Loading