diff --git a/packages/client/src/Client/Internal/Connection/LazyConnection.php b/packages/client/src/Client/Internal/Connection/LazyConnection.php index 8428691..7c89a4a 100644 --- a/packages/client/src/Client/Internal/Connection/LazyConnection.php +++ b/packages/client/src/Client/Internal/Connection/LazyConnection.php @@ -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; + } } } diff --git a/tests/Client/Internal/Connection/LazyConnectionTest.php b/tests/Client/Internal/Connection/LazyConnectionTest.php new file mode 100644 index 0000000..81fd44a --- /dev/null +++ b/tests/Client/Internal/Connection/LazyConnectionTest.php @@ -0,0 +1,158 @@ +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()), + ); + } +} diff --git a/tests/Client/Internal/Connection/StubConnection.php b/tests/Client/Internal/Connection/StubConnection.php new file mode 100644 index 0000000..8023f81 --- /dev/null +++ b/tests/Client/Internal/Connection/StubConnection.php @@ -0,0 +1,37 @@ + $invoke + * @return StubStream + */ + #[\Override] + public function createStream(Invoke $invoke, Metadata $md, Cancellation $cancellation, PickContext $pick): ClientStream + { + /** @var StubStream */ + return new StubStream(); + } + + #[\Override] + public function close(Cancellation $cancellation = new NullCancellation()): void + { + $this->closed = true; + } +} diff --git a/tests/Client/Internal/Connection/StubStream.php b/tests/Client/Internal/Connection/StubStream.php new file mode 100644 index 0000000..3013cb2 --- /dev/null +++ b/tests/Client/Internal/Connection/StubStream.php @@ -0,0 +1,50 @@ + + */ +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 []; + } +}