diff --git a/CHANGELOG.md b/CHANGELOG.md index 3018f03d..080efd8f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,9 @@ ### Fixed - Fixed wrong `@param` annotation in `Parameters` (#1614) +### Added +- Added retry support (#1616) + ## v3.3.0 (2025-11-24) ### Added - Added cluster support for `XADD`, `XDEL` and `XRANGE` (#1587) diff --git a/README.md b/README.md index 276baa54..b44bfaff 100644 --- a/README.md +++ b/README.md @@ -509,6 +509,43 @@ $client = new Predis\Client('tcp://127.0.0.1', [ For a more in-depth insight on how to create new connection backends you can refer to the actual implementation of the standard connection classes available in the `Predis\Connection` namespace. +### Retry exceptions + +You can enable automatic retry that is disabled by default, to be able to reduce the amount of +false-positives in case of network issues. By default, we're retrying on any connection, +timeout or socket initialization exception, but you can update the list of retry +exceptions. For now `EqualBackoff` and `ExponentialBackoff` strategies are available, +but you may provide your custom one. Retry may be configured with any type of communication +(standalone node, cluster, pipeline, transaction, replication). Here's an example of +configuration: + +```php +// Standalone client +$client = new Predis\Client([ + 'retry' => new \Predis\Retry\Retry( + new \Predis\Retry\Strategy\ExponentialBackoff(1000, 10000), // Base and cap configuration in microseconds + 3 // Number of retries + ), +]); + +// Cluster configuration +$options = [ + 'parameters' => [ + 'retry' => new \Predis\Retry\Retry(new \Predis\Retry\Strategy\ExponentialBackoff(1000, 10000), 3), + ], +]; + +$client = new Predis\Client(['tcp://host:port', 'tcp://host:port', 'tcp://host:port'], $options); + +$retry = new \Predis\Retry\Retry( + new \Predis\Retry\Strategy\ExponentialBackoff(1000, 10000), + 3 +); + +// Update a list of exceptions to catch +$retry->updateCatchableExceptions([Exception::class]); +``` + ## RESP3 ## ### Connection ### diff --git a/phpstan.dist.neon b/phpstan.dist.neon index 46821cf8..05a3c75c 100644 --- a/phpstan.dist.neon +++ b/phpstan.dist.neon @@ -29,9 +29,6 @@ parameters: - message: "#^Access to an undefined property Predis\\\\Connection\\\\ParametersInterface\\:\\:\\$weight\\.$#" count: 1 path: src/Connection/Cluster/PredisCluster.php - - message: "#^Variable \\$response might not be defined\\.$#" - count: 2 - path: src/Connection/Cluster/RedisCluster.php - message: "#^Access to an undefined property Predis\\\\Connection\\\\ParametersInterface\\:\\:\\$role\\.$#" count: 1 path: src/Connection/Replication/MasterSlaveReplication.php @@ -45,6 +42,3 @@ parameters: - message: "#^Variable \\$response might not be defined\\.$#" count: 1 path: src/Connection/Replication/MasterSlaveReplication.php - - message: "#^Variable \\$response might not be defined\\.$#" - count: 1 - path: src/Connection/Replication/SentinelReplication.php diff --git a/src/Client.php b/src/Client.php index 3b86083d..c0535c6f 100644 --- a/src/Client.php +++ b/src/Client.php @@ -22,6 +22,7 @@ use Predis\Command\RawCommand; use Predis\Command\ScriptCommand; use Predis\Configuration\Options; use Predis\Configuration\OptionsInterface; +use Predis\Connection\AggregateConnectionInterface; use Predis\Connection\ConnectionInterface; use Predis\Connection\Parameters; use Predis\Connection\ParametersInterface; @@ -41,6 +42,7 @@ use Predis\Response\ServerException; use Predis\Transaction\MultiExec as MultiExecTransaction; use ReturnTypeWillChange; use RuntimeException; +use Throwable; use Traversable; /** @@ -376,12 +378,25 @@ class Client implements ClientInterface, IteratorAggregate /** * {@inheritdoc} + * @throws Throwable */ public function executeCommand(CommandInterface $command) { - $response = $this->connection->executeCommand($command); $parameters = $this->connection->getParameters(); + if ($this->connection instanceof AggregateConnectionInterface || $this->connection instanceof RelayConnection) { + $response = $this->connection->executeCommand($command); + } else { + $response = $parameters->retry->callWithRetry( + function () use ($command) { + return $this->connection->executeCommand($command); + }, + function () { + $this->connection->disconnect(); + } + ); + } + if ($response instanceof ResponseInterface) { if ($response instanceof ErrorResponseInterface) { $response = $this->onErrorResponse($command, $response); diff --git a/src/Connection/AbstractConnection.php b/src/Connection/AbstractConnection.php index e04c8e72..777d66fa 100644 --- a/src/Connection/AbstractConnection.php +++ b/src/Connection/AbstractConnection.php @@ -19,6 +19,7 @@ use Predis\Connection\Resource\Exception\StreamInitException; use Predis\Protocol\Parser\ParserStrategyResolver; use Predis\Protocol\Parser\Strategy\ParserStrategyInterface; use Predis\Protocol\ProtocolException; +use Predis\TimeoutException; /** * Base class with the common logic used by connection classes to communicate @@ -158,6 +159,20 @@ abstract class AbstractConnection implements NodeConnectionInterface ); } + /** + * Helper method to handle timeout errors. + * + * @param int $code + * @return void + * @throws CommunicationException + */ + protected function onTimeoutError(int $code = 0): void + { + CommunicationException::handle( + new TimeoutException($this, $code) + ); + } + /** * Helper method to handle protocol errors. * diff --git a/src/Connection/Cluster/RedisCluster.php b/src/Connection/Cluster/RedisCluster.php index ba8ac5e2..bac1d533 100644 --- a/src/Connection/Cluster/RedisCluster.php +++ b/src/Connection/Cluster/RedisCluster.php @@ -28,10 +28,14 @@ use Predis\Connection\ConnectionException; use Predis\Connection\FactoryInterface; use Predis\Connection\NodeConnectionInterface; use Predis\Connection\ParametersInterface; +use Predis\Connection\RelayFactory; use Predis\NotSupportedException; use Predis\Response\Error as ErrorResponse; use Predis\Response\ErrorInterface as ErrorResponseInterface; use Predis\Response\ServerException; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; +use Predis\TimeoutException; use ReturnTypeWillChange; use Throwable; use Traversable; @@ -58,7 +62,7 @@ use Traversable; */ class RedisCluster extends AbstractAggregateConnection implements ClusterInterface, IteratorAggregate, Countable { - private $useClusterSlots = true; + public $useClusterSlots = true; /** * @var NodeConnectionInterface[] @@ -260,35 +264,31 @@ class RedisCluster extends AbstractAggregateConnection implements ClusterInterfa */ private function queryClusterNodeForSlotMap(NodeConnectionInterface $connection) { - $retries = 0; - $retryAfter = $this->retryInterval; + // Backward-compatible hardcoded retry + $retry = new Retry( + new ExponentialBackoff($this->retryInterval * 1000, -1), + $this->retryLimit, + [ConnectionException::class] + ); + $command = RawCommand::create('CLUSTER', 'SLOTS'); - while ($retries <= $this->retryLimit) { - try { - $response = $connection->executeCommand($command); - break; - } catch (ConnectionException $exception) { - $connection = $exception->getConnection(); - $connection->disconnect(); + $doCallback = function () use (&$connection, $command) { + return $connection->executeCommand($command); + }; - $this->remove($connection); + $failCallback = function (ConnectionException $exception) use (&$connection) { + $connection = $exception->getConnection(); + $connection->disconnect(); - if ($retries === $this->retryLimit) { - throw $exception; - } + $this->remove($connection); - if (!$connection = $this->getRandomConnection()) { - throw new ClientException('No connections left in the pool for `CLUSTER SLOTS`'); - } - - usleep($retryAfter * 1000); - $retryAfter *= 2; - ++$retries; + if (!$connection = $this->getRandomConnection()) { + throw new ClientException('No connections left in the pool for `CLUSTER SLOTS`'); } - } + }; - return $response; + return $retry->callWithRetry($doCallback, $failCallback); } /** @@ -545,51 +545,42 @@ class RedisCluster extends AbstractAggregateConnection implements ClusterInterfa * @param string $method Actual method. * * @return mixed + * @throws Throwable */ private function retryCommandOnFailure(CommandInterface $command, $method) { - $retries = 0; - $retryAfter = $this->retryInterval; - - while ($retries <= $this->retryLimit) { - try { - $response = $this->getConnectionByCommand($command)->$method($command); - - if ($response instanceof ErrorResponse) { - $message = $response->getMessage(); - - if (strpos($message, 'CLUSTERDOWN') !== false) { - throw new ServerException($message); - } - } - - break; - } catch (Throwable $exception) { - usleep($retryAfter * 1000); - $retryAfter *= 2; - - if ($exception instanceof ConnectionException) { - $connection = $exception->getConnection(); - - if ($connection) { - $connection->disconnect(); - $this->remove($connection); - } - } - - if ($retries === $this->retryLimit) { - throw $exception; - } - - if ($this->useClusterSlots) { - $this->askSlotMap(); - } - - ++$retries; - } + if ($this->connectionParameters->isDisabledRetry() || $this->connections instanceof RelayFactory) { + // Override default parameters, for backward-compatibility + // with current behaviour + $retry = new Retry( + new ExponentialBackoff($this->retryInterval * 1000, -1), + $this->retryLimit + ); + } else { + $retry = $this->connectionParameters->retry; } + $retry->updateCatchableExceptions([ServerException::class]); - return $response; + $doCallback = function () use ($command, $method) { + $response = $this->getConnectionByCommand($command)->$method($command); + + if ($response instanceof ErrorResponse) { + $message = $response->getMessage(); + + if (strpos($message, 'CLUSTERDOWN') !== false) { + throw new ServerException($message); + } + } + + return $response; + }; + + return $retry->callWithRetry( + $doCallback, + function (Throwable $e) { + $this->onFailCallback($e); + } + ); } /** @@ -740,4 +731,34 @@ class RedisCluster extends AbstractAggregateConnection implements ClusterInterfa usleep($this->readTimeout); } } + + /** + * Handle exceptions. + * + * @param Throwable $exception + * @return void + */ + private function onFailCallback(Throwable $exception) + { + if ($exception instanceof ConnectionException) { + $connection = $exception->getConnection(); + + if ($connection) { + $connection->disconnect(); + $this->remove($connection); + } + + if ($this->useClusterSlots) { + $this->askSlotMap(); + } + } + + if ($exception instanceof TimeoutException) { + $connection = $exception->getConnection(); + + if ($connection) { + $connection->disconnect(); + } + } + } } diff --git a/src/Connection/Parameters.php b/src/Connection/Parameters.php index 621ea655..2ee11c7d 100644 --- a/src/Connection/Parameters.php +++ b/src/Connection/Parameters.php @@ -13,6 +13,8 @@ namespace Predis\Connection; use InvalidArgumentException; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\NoBackoff; /** * Container for connection parameters used to initialize connections to Redis. @@ -36,11 +38,23 @@ class Parameters implements ParametersInterface */ protected $parameters; + /** + * @var bool + */ + private $disabledRetry = true; + /** * @param array $parameters Named array of connection parameters. */ public function __construct(array $parameters = []) { + if (!array_key_exists('retry', $parameters)) { + // Retries disabled by default + static::$defaults['retry'] = new Retry(new NoBackoff(), 0); + } else { + $this->disabledRetry = false; + } + $this->parameters = $this->filter($parameters + static::$defaults); } @@ -195,6 +209,16 @@ class Parameters implements ParametersInterface return "$this->scheme://$this->host:$this->port"; } + /** + * Returns if retries is disabled. + * + * @return bool + */ + public function isDisabledRetry(): bool + { + return $this->disabledRetry; + } + /** * {@inheritdoc} */ diff --git a/src/Connection/ParametersInterface.php b/src/Connection/ParametersInterface.php index 3514a9bf..de1a413c 100644 --- a/src/Connection/ParametersInterface.php +++ b/src/Connection/ParametersInterface.php @@ -12,6 +12,8 @@ namespace Predis\Connection; +use Predis\Retry\Retry; + /** * Interface defining a container for connection parameters. * @@ -35,9 +37,11 @@ namespace Predis\Connection; * @property bool $async_connect Performs the connect() operation asynchronously. * @property bool $tcp_nodelay Toggles the Nagle's algorithm for coalescing. * @property bool $client_info Whether to set LIB-NAME and LIB-VER when connecting. + * @property Retry $retry Retry configuration * @property bool $cache (Relay only) Whether to use in-memory caching. * @property string $serializer (Relay only) Serializer used for data serialization. * @property string $compression (Relay only) Algorithm used for data compression. + * @method bool isDisabledRetry() Specify if custom retry configuration was provided. */ interface ParametersInterface { diff --git a/src/Connection/RelayConnection.php b/src/Connection/RelayConnection.php index 65ad36cd..7a5bd55d 100644 --- a/src/Connection/RelayConnection.php +++ b/src/Connection/RelayConnection.php @@ -128,44 +128,6 @@ class RelayConnection extends AbstractConnection } } - /** - * Creates a new instance of the client. - * - * @return Relay - */ - private function createClient() - { - $client = new Relay(); - - // throw when errors occur and return `null` for non-existent keys - $client->setOption(Relay::OPT_PHPREDIS_COMPATIBILITY, false); - - // use reply literals - $client->setOption(Relay::OPT_REPLY_LITERAL, true); - - // disable Relay's command/connection retry - $client->setOption(Relay::OPT_MAX_RETRIES, 0); - - // whether to use in-memory caching - $client->setOption(Relay::OPT_USE_CACHE, $this->parameters->cache ?? true); - - // set data serializer - $client->setOption(Relay::OPT_SERIALIZER, constant(sprintf( - '%s::SERIALIZER_%s', - Relay::class, - strtoupper($this->parameters->serializer ?? 'none') - ))); - - // set data compression algorithm - $client->setOption(Relay::OPT_COMPRESSION, constant(sprintf( - '%s::COMPRESSION_%s', - Relay::class, - strtoupper($this->parameters->compression ?? 'none') - ))); - - return $client; - } - /** * Returns the underlying client. * diff --git a/src/Connection/RelayFactory.php b/src/Connection/RelayFactory.php index 57ca78c0..350519a4 100644 --- a/src/Connection/RelayFactory.php +++ b/src/Connection/RelayFactory.php @@ -15,6 +15,8 @@ namespace Predis\Connection; use InvalidArgumentException; use Predis\Command\RawCommand; use Predis\NotSupportedException; +use Predis\Retry\Strategy\EqualBackoff; +use Predis\Retry\Strategy\ExponentialBackoff; use Relay\Relay; class RelayFactory extends Factory @@ -64,7 +66,7 @@ class RelayFactory extends Factory } $initializer = $this->schemes[$scheme]; - $client = $this->createClient(); + $client = $this->createClient($parameters); $connection = new $initializer($parameters, $client); @@ -90,7 +92,7 @@ class RelayFactory extends Factory * * @return Relay */ - private function createClient() + private function createClient(ParametersInterface $parameters) { $client = new Relay(); @@ -100,26 +102,50 @@ class RelayFactory extends Factory // use reply literals $client->setOption(Relay::OPT_REPLY_LITERAL, true); - // disable Relay's command/connection retry - $client->setOption(Relay::OPT_MAX_RETRIES, 0); - // whether to use in-memory caching - $client->setOption(Relay::OPT_USE_CACHE, $this->parameters->cache ?? true); + $client->setOption(Relay::OPT_USE_CACHE, $parameters->cache ?? true); // set data serializer $client->setOption(Relay::OPT_SERIALIZER, constant(sprintf( '%s::SERIALIZER_%s', Relay::class, - strtoupper($this->parameters->serializer ?? 'none') + strtoupper($parameters->serializer ?? 'none') ))); // set data compression algorithm $client->setOption(Relay::OPT_COMPRESSION, constant(sprintf( '%s::COMPRESSION_%s', Relay::class, - strtoupper($this->parameters->compression ?? 'none') + strtoupper($parameters->compression ?? 'none') ))); + if ($parameters->isDisabledRetry()) { + $client->setOption(Relay::OPT_MAX_RETRIES, 0); + } else { + $client->setOption(Relay::OPT_MAX_RETRIES, $parameters->retry->getRetries()); + + $retryStrategy = $parameters->retry->getStrategy(); + + if ($retryStrategy instanceof ExponentialBackoff) { + $algorithm = Relay::BACKOFF_ALGORITHM_FULL_JITTER; + $base = $retryStrategy->getBase(); + $cap = $retryStrategy->getCap(); + } else { + $algorithm = Relay::BACKOFF_ALGORITHM_DEFAULT; + + if ($retryStrategy instanceof EqualBackoff) { + $base = $cap = $retryStrategy->compute(0); + } else { + $base = $retryStrategy::DEFAULT_BASE; + $cap = $retryStrategy::DEFAULT_CAP; + } + } + + $client->setOption(Relay::OPT_BACKOFF_ALGORITHM, $algorithm); + $client->setOption(Relay::OPT_BACKOFF_BASE, $base / 1000); + $client->setOption(Relay::OPT_BACKOFF_CAP, $cap / 1000); + } + return $client; } diff --git a/src/Connection/Replication/MasterSlaveReplication.php b/src/Connection/Replication/MasterSlaveReplication.php index ff5da336..9e186ee8 100644 --- a/src/Connection/Replication/MasterSlaveReplication.php +++ b/src/Connection/Replication/MasterSlaveReplication.php @@ -22,9 +22,12 @@ use Predis\Connection\ConnectionException; use Predis\Connection\FactoryInterface; use Predis\Connection\NodeConnectionInterface; use Predis\Connection\ParametersInterface; +use Predis\Connection\RelayFactory; use Predis\Replication\MissingMasterException; use Predis\Replication\ReplicationStrategy; use Predis\Response\ErrorInterface as ResponseErrorInterface; +use Predis\TimeoutException; +use Throwable; /** * Aggregate connection handling replication of Redis nodes configured in a @@ -476,9 +479,26 @@ class MasterSlaveReplication extends AbstractAggregateConnection implements Repl * @param string $method Actual method. * * @return mixed + * @throws Throwable */ private function retryCommandOnFailure(CommandInterface $command, $method) { + $parameters = $this->getParameters(); + + if (!$parameters->isDisabledRetry() && !$this->connectionFactory instanceof RelayFactory) { + $retry = $parameters->retry; + $retry->updateCatchableExceptions([MissingMasterException::class]); + + return $retry->callWithRetry( + function () use ($command, $method) { + return $this->executeCommandInternal($command, $method); + }, + function (Throwable $exception) { + $this->onFailCallback($exception); + } + ); + } + while (true) { try { $connection = $this->getConnectionByCommand($command); @@ -490,38 +510,35 @@ class MasterSlaveReplication extends AbstractAggregateConnection implements Repl break; } catch (ConnectionException $exception) { - $connection = $exception->getConnection(); - $connection->disconnect(); - - if ($connection === $this->master && !$this->autoDiscovery) { - // Throw immediately when master connection is failing, even - // when the command represents a read-only operation, unless - // automatic discovery has been enabled. - throw $exception; - } else { - // Otherwise remove the failing slave and attempt to execute - // the command again on one of the remaining slaves... - $this->remove($connection); - } - - // ... that is, unless we have no more connections to use. - if (!$this->slaves && !$this->master) { - throw $exception; - } elseif ($this->autoDiscovery) { - $this->discover(); - } + $this->onConnectionExceptionCallback($exception); } catch (MissingMasterException $exception) { - if ($this->autoDiscovery) { - $this->discover(); - } else { - throw $exception; - } + $this->onMissingMasterException($exception); } } return $response; } + /** + * Executes command against valid connection. + * + * @param CommandInterface $command + * @param string $method + * @return mixed + * @throws ConnectionException + */ + protected function executeCommandInternal(CommandInterface $command, string $method) + { + $connection = $this->getConnectionByCommand($command); + $response = $connection->$method($command); + + if ($response instanceof ResponseErrorInterface && $response->getErrorType() === 'LOADING') { + throw new ConnectionException($connection, "Redis is loading the dataset in memory [$connection]"); + } + + return $response; + } + /** * {@inheritdoc} */ @@ -571,4 +588,84 @@ class MasterSlaveReplication extends AbstractAggregateConnection implements Repl return null; } + + /** + * Handle connection exception. + * + * @param ConnectionException $exception + * @return void + * @throws ClientException|ConnectionException + */ + private function onConnectionExceptionCallback(ConnectionException $exception) + { + $connection = $exception->getConnection(); + $connection->disconnect(); + + if ($connection === $this->master && !$this->autoDiscovery) { + // Throw immediately when master connection is failing, even + // when the command represents a read-only operation, unless + // automatic discovery has been enabled. + throw $exception; + } else { + // Otherwise remove the failing slave and attempt to execute + // the command again on one of the remaining slaves... + $this->remove($connection); + } + + // ... that is, unless we have no more connections to use. + if (!$this->slaves && !$this->master) { + throw $exception; + } elseif ($this->autoDiscovery) { + $this->discover(); + } + } + + /** + * Exception handling callback. + * + * @param Throwable $exception + * @return void + * @throws Throwable + */ + private function onFailCallback(Throwable $exception) + { + if ($exception instanceof ConnectionException) { + $this->onConnectionExceptionCallback($exception); + + return; + } + + if ($exception instanceof MissingMasterException) { + $this->onMissingMasterException($exception); + + return; + } + + if ($exception instanceof TimeoutException) { + $connection = $exception->getConnection(); + + if ($connection) { + $connection->disconnect(); + + return; + } + } + + throw $exception; + } + + /** + * @param MissingMasterException $exception + * @return void + * @throws ClientException + * @throws MissingMasterException + */ + private function onMissingMasterException(MissingMasterException $exception) + { + if ($this->autoDiscovery) { + $this->discover(); + } else { + throw $exception; + } + } } diff --git a/src/Connection/Replication/SentinelReplication.php b/src/Connection/Replication/SentinelReplication.php index 1d0310f3..c742c436 100644 --- a/src/Connection/Replication/SentinelReplication.php +++ b/src/Connection/Replication/SentinelReplication.php @@ -16,17 +16,21 @@ use InvalidArgumentException; use Predis\Command\Command; use Predis\Command\CommandInterface; use Predis\Command\RawCommand; +use Predis\CommunicationException; use Predis\Connection\AbstractAggregateConnection; use Predis\Connection\ConnectionException; use Predis\Connection\FactoryInterface as ConnectionFactoryInterface; use Predis\Connection\NodeConnectionInterface; use Predis\Connection\Parameters; use Predis\Connection\ParametersInterface; +use Predis\Connection\RelayFactory; use Predis\Replication\ReplicationStrategy; use Predis\Replication\RoleException; use Predis\Response\Error; use Predis\Response\ErrorInterface as ErrorResponseInterface; use Predis\Response\ServerException; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; use Throwable; /** @@ -566,7 +570,10 @@ class SentinelReplication extends AbstractAggregateConnection implements Replica protected function assertConnectionRole(NodeConnectionInterface $connection, $role) { $role = strtolower($role); - $actualRole = $connection->executeCommand(RawCommand::create('ROLE')); + $retry = $connection->getParameters()->retry; + $actualRole = $retry->callWithRetry(function () use ($connection) { + return $connection->executeCommand(RawCommand::create('ROLE')); + }); if ($actualRole instanceof Error) { throw new ConnectionException($connection, $actualRole->getMessage()); @@ -710,33 +717,39 @@ class SentinelReplication extends AbstractAggregateConnection implements Replica */ private function retryCommandOnFailure(CommandInterface $command, $method) { - $retries = 0; + $parameters = $this->getParameters(); - while ($retries <= $this->retryLimit) { - try { - $response = $this->getConnectionByCommand($command)->$method($command); - if ($response instanceof Error && $response->getErrorType() === 'LOADING') { - throw new ConnectionException($this->current, $response->getMessage()); - } - break; - } catch (Throwable $exception) { - $this->wipeServerList(); - - if ($exception instanceof ConnectionException) { - $exception->getConnection()->disconnect(); - } - - if ($retries === $this->retryLimit) { - throw $exception; - } - - usleep($this->retryWait * 1000); - - ++$retries; - } + if ($parameters->isDisabledRetry() || $this->connectionFactory instanceof RelayFactory) { + // Override default parameters, for backward-compatibility + // with current behaviour + $retry = new Retry( + new ExponentialBackoff($this->retryWait * 1000, -1), + $this->retryLimit + ); + } else { + $retry = $parameters->retry; } + $retry->updateCatchableExceptions([Throwable::class]); - return $response; + $doCallback = function () use ($method, $command) { + $response = $this->getConnectionByCommand($command)->{$method}($command); + + if ($response instanceof Error && $response->getErrorType() === 'LOADING') { + throw new ConnectionException($this->current, $response->getMessage()); + } + + return $response; + }; + + $failCallback = function (Throwable $exception) { + $this->wipeServerList(); + + if ($exception instanceof CommunicationException) { + $exception->getConnection()->disconnect(); + } + }; + + return $retry->callWithRetry($doCallback, $failCallback); } /** diff --git a/src/Connection/Resource/Stream.php b/src/Connection/Resource/Stream.php index 271a1607..d9309c7d 100644 --- a/src/Connection/Resource/Stream.php +++ b/src/Connection/Resource/Stream.php @@ -207,7 +207,27 @@ class Stream implements StreamInterface $result = fwrite($this->stream, $string); - if ($result === false) { + if ($result === false || $result === 0) { + $metadata = $this->getMetadata(); + + if ($this->eof()) { + throw new RuntimeException('Connection closed by peer during write', 1); + } + + if (!is_resource($this->stream)) { + throw new RuntimeException( + 'Stream resource is no longer valid', + 1 + ); + } + + if (array_key_exists('timed_out', $metadata) && $metadata['timed_out']) { + throw new RuntimeException( + 'Stream has been timed out', + 2 + ); + } + throw new RuntimeException('Unable to write to stream', 1); } @@ -252,6 +272,26 @@ class Stream implements StreamInterface } if (false === $string) { + $metadata = $this->getMetadata(); + + if ($this->eof()) { + throw new RuntimeException('Connection closed by peer during read', 1); + } + + if (!is_resource($this->stream)) { + throw new RuntimeException( + 'Stream resource is no longer valid', + 1 + ); + } + + if (array_key_exists('timed_out', $metadata) && $metadata['timed_out']) { + throw new RuntimeException( + 'Stream has been timed out', + 2 + ); + } + throw new RuntimeException('Unable to read from stream', 1); } diff --git a/src/Connection/StreamConnection.php b/src/Connection/StreamConnection.php index 667dede7..67ed8ab4 100644 --- a/src/Connection/StreamConnection.php +++ b/src/Connection/StreamConnection.php @@ -340,11 +340,14 @@ class StreamConnection extends AbstractConnection * @param string|null $message * @throws RuntimeException|CommunicationException */ - protected function onStreamError(RuntimeException $e, ?string $message = null) + protected function onStreamError($e, ?string $message = null) { - // Code = 1 represents issues related to read/write operation. + // Code = 1 represents issues related to read/write operation, connection broken. if ($e->getCode() === 1) { $this->onConnectionError($message); + } elseif ($e->getCode() === 2) { + // Operation has been timed out, connection not necessarily broken. + $this->onTimeoutError(); } throw $e; diff --git a/src/Pipeline/Atomic.php b/src/Pipeline/Atomic.php index 369c1fd8..66fa989b 100644 --- a/src/Pipeline/Atomic.php +++ b/src/Pipeline/Atomic.php @@ -14,13 +14,16 @@ namespace Predis\Pipeline; use Predis\ClientException; use Predis\ClientInterface; -use Predis\Connection\AggregateConnectionInterface; +use Predis\Command\Command; +use Predis\Command\CommandInterface; +use Predis\CommunicationException; use Predis\Connection\ConnectionInterface; use Predis\Connection\NodeConnectionInterface; use Predis\Response\ErrorInterface as ErrorResponseInterface; use Predis\Response\ResponseInterface; use Predis\Response\ServerException; use SplQueue; +use Throwable; /** * Command pipeline wrapped into a MULTI / EXEC transaction. @@ -63,24 +66,18 @@ class Atomic extends Pipeline protected function executePipeline(ConnectionInterface $connection, SplQueue $commands) { $commandFactory = $this->getClient()->getCommandFactory(); - $connection->executeCommand($commandFactory->create('multi')); + $retry = $connection->getParameters()->retry; + $this->executeCommandWithRetry($connection, $commandFactory->create('multi')); - if ($connection instanceof AggregateConnectionInterface) { - $this->writeToMultiNode($connection, $commands); - } else { - $this->writeToSingleNode($connection, $commands); - } - - foreach ($commands as $command) { - $response = $connection->readResponse($command); - - if ($response instanceof ErrorResponseInterface) { - $connection->executeCommand($commandFactory->create('discard')); - throw new ServerException($response->getMessage()); + $retry->callWithRetry(function () use ($connection, $commands) { + $this->queuePipeline($connection, $commands); + }, function (Throwable $exception) { + if ($exception instanceof CommunicationException) { + $exception->getConnection()->disconnect(); } - } + }); - $executed = $connection->executeCommand($commandFactory->create('exec')); + $executed = $this->executeCommandWithRetry($connection, $commandFactory->create('exec')); if (!isset($executed)) { throw new ClientException( @@ -123,4 +120,44 @@ class Atomic extends Pipeline return $responses; } + + /** + * @param ConnectionInterface $connection + * @param SplQueue $commands + * @return void + * @throws Throwable + */ + protected function queuePipeline(ConnectionInterface $connection, SplQueue $commands) + { + $commandFactory = $this->getClient()->getCommandFactory(); + $this->writeToSingleNode($connection, $commands); + + foreach ($commands as $command) { + $response = $connection->readResponse($command); + + if ($response instanceof ErrorResponseInterface) { + $this->executeCommandWithRetry($connection, $commandFactory->create('discard')); + throw new ServerException($response->getMessage()); + } + } + } + + /** + * @param ConnectionInterface $connection + * @param Command $command + * @return mixed + * @throws Throwable + */ + protected function executeCommandWithRetry(ConnectionInterface $connection, CommandInterface $command) + { + $retry = $connection->getParameters()->retry; + + return $retry->callWithRetry(function () use ($connection, $command) { + return $connection->executeCommand($command); + }, function (Throwable $e) { + if ($e instanceof CommunicationException) { + $e->getConnection()->disconnect(); + } + }); + } } diff --git a/src/Pipeline/FireAndForget.php b/src/Pipeline/FireAndForget.php index aa5bb1b0..77c3b360 100644 --- a/src/Pipeline/FireAndForget.php +++ b/src/Pipeline/FireAndForget.php @@ -12,9 +12,11 @@ namespace Predis\Pipeline; +use Predis\CommunicationException; use Predis\Connection\AggregateConnectionInterface; use Predis\Connection\ConnectionInterface; use SplQueue; +use Throwable; /** * Command pipeline that writes commands to the servers but discards responses. @@ -26,11 +28,19 @@ class FireAndForget extends Pipeline */ protected function executePipeline(ConnectionInterface $connection, SplQueue $commands) { - if ($connection instanceof AggregateConnectionInterface) { - $this->writeToMultiNode($connection, $commands); - } else { - $this->writeToSingleNode($connection, $commands); - } + $retry = $connection->getParameters()->retry; + + $retry->callWithRetry(function () use ($connection, $commands) { + if ($connection instanceof AggregateConnectionInterface) { + $this->writeToMultiNode($connection, $commands); + } else { + $this->writeToSingleNode($connection, $commands); + } + }, function (Throwable $e) { + if ($e instanceof CommunicationException) { + $e->getConnection()->disconnect(); + } + }); $connection->disconnect(); diff --git a/src/Pipeline/Pipeline.php b/src/Pipeline/Pipeline.php index 9ca6705f..7e398dd3 100644 --- a/src/Pipeline/Pipeline.php +++ b/src/Pipeline/Pipeline.php @@ -18,13 +18,18 @@ use Predis\ClientContextInterface; use Predis\ClientException; use Predis\ClientInterface; use Predis\Command\CommandInterface; +use Predis\CommunicationException; use Predis\Connection\AggregateConnectionInterface; +use Predis\Connection\Cluster\RedisCluster; +use Predis\Connection\ConnectionException; use Predis\Connection\ConnectionInterface; use Predis\Connection\Replication\ReplicationInterface; use Predis\Response\ErrorInterface as ErrorResponseInterface; use Predis\Response\ResponseInterface; use Predis\Response\ServerException; +use Predis\TimeoutException; use SplQueue; +use Throwable; /** * Implementation of a command pipeline in which write and read operations of @@ -129,22 +134,65 @@ class Pipeline implements ClientContextInterface * @param SplQueue $commands Queued commands. * * @return array + * @throws Throwable */ protected function executePipeline(ConnectionInterface $connection, SplQueue $commands) { + $retry = $connection->getParameters()->retry; + $backupQueue = $this->createDeepCloneQueue($commands); + + return $retry->callWithRetry( + function () use ($connection, &$commands) { + return $this->executePipelineInternal($connection, $commands); + }, + function (Throwable $e) use (&$commands, $backupQueue, $connection) { + if (!$e instanceof CommunicationException) { + throw $e; + } + + if ($connection instanceof AggregateConnectionInterface) { + $this->onAggregateConnectionFailCallback($connection, $e); + } else { + $connection = $e->getConnection(); + $connection->disconnect(); + } + + // In case of error whole pipeline should be retried + // So we need to write all original commands again + $commands = $this->createDeepCloneQueue($backupQueue); + } + ); + } + + /** + * @param ConnectionInterface $connection + * @param SplQueue $commands + * @return array + * @throws ServerException + * @throws Throwable + */ + protected function executePipelineInternal( + ConnectionInterface $connection, + SplQueue $commands + ): array { + $responses = []; + $exceptions = $this->throwServerExceptions(); + $protocolVersion = (int) $connection->getParameters()->protocol; + if ($connection instanceof AggregateConnectionInterface) { $this->writeToMultiNode($connection, $commands); } else { $this->writeToSingleNode($connection, $commands); } - $responses = []; - $exceptions = $this->throwServerExceptions(); - $protocolVersion = (int) $connection->getParameters()->protocol; - while (!$commands->isEmpty()) { $command = $commands->dequeue(); - $response = $connection->readResponse($command); + + if ($connection instanceof AggregateConnectionInterface) { + $response = $connection->getConnectionByCommand($command)->readResponse($command); + } else { + $response = $connection->readResponse($command); + } if (!$response instanceof ResponseInterface) { if ($protocolVersion === 2) { @@ -162,12 +210,30 @@ class Pipeline implements ClientContextInterface return $responses; } + /** + * Creates a deep copy of commands queue for backup. + * + * @param SplQueue $queue + * @return SplQueue + */ + private function createDeepCloneQueue(SplQueue $queue): SplQueue + { + $new = new SplQueue(); + + foreach ($queue as $command) { + $new->enqueue(clone $command); + } + + return $new; + } + /** * Writes pipelined commands to single node connection. * * @param ConnectionInterface $connection * @param SplQueue $commands * @return void + * @throws Throwable */ protected function writeToSingleNode(ConnectionInterface $connection, SplQueue $commands) { @@ -186,9 +252,12 @@ class Pipeline implements ClientContextInterface * @param AggregateConnectionInterface $connection * @param SplQueue $commands * @return void + * @throws Throwable */ protected function writeToMultiNode(AggregateConnectionInterface $connection, SplQueue $commands) { + $retry = $connection->getParameters()->retry; + foreach ($commands as $command) { $nodeConnection = $connection->getConnectionByCommand($command); $nodeConnection->write($command->serializeCommand()); @@ -286,4 +355,37 @@ class Pipeline implements ClientContextInterface { return $this->client; } + + /** + * Handle aggregate connection exception. + * + * @param AggregateConnectionInterface $connection + * @param CommunicationException $e + * @return void + */ + private function onAggregateConnectionFailCallback(AggregateConnectionInterface $connection, Throwable $e) + { + if ($e instanceof ConnectionException) { + $nodeConnection = $e->getConnection(); + + if ($nodeConnection) { + $nodeConnection->disconnect(); + $connection->remove($nodeConnection); + } + + if ($connection instanceof RedisCluster) { + if ($connection->useClusterSlots) { + $connection->askSlotMap(); + } + } + } + + if ($e instanceof TimeoutException) { + $nodeConnection = $e->getConnection(); + + if ($nodeConnection) { + $nodeConnection->disconnect(); + } + } + } } diff --git a/src/Retry/Retry.php b/src/Retry/Retry.php new file mode 100644 index 00000000..a6b8d50f --- /dev/null +++ b/src/Retry/Retry.php @@ -0,0 +1,143 @@ +backoffStrategy = $backoffStrategy; + $this->retries = $retries; + + if (null !== $catchableExceptions) { + $this->catchableExceptions = $catchableExceptions; + } + } + + /** + * Update the retry count. + * + * @param int $retries + * @return void + */ + public function updateRetriesCount(int $retries): void + { + $this->retries = $retries; + } + + /** + * Extend catchable exceptions list. + * + * @param array $catchableExceptions + * @return void + */ + public function updateCatchableExceptions(array $catchableExceptions): void + { + $this->catchableExceptions = array_merge($this->catchableExceptions, $catchableExceptions); + } + + /** + * @return int + */ + public function getRetries(): int + { + return $this->retries; + } + + /** + * @return RetryStrategyInterface + */ + public function getStrategy(): RetryStrategyInterface + { + return $this->backoffStrategy; + } + + /** + * @param callable(): mixed $do + * @param callable(Throwable): void|null $fail + * @return mixed + * @throws Throwable + */ + public function callWithRetry(callable $do, ?callable $fail = null) + { + $failures = 0; + + while (true) { + try { + return $do(); + } catch (Throwable $e) { + if (null !== $this->catchableExceptions) { + $isCatchable = false; + foreach ($this->catchableExceptions as $catchableException) { + if ($e instanceof $catchableException) { + $isCatchable = true; + } + } + + if (!$isCatchable) { + throw $e; + } + } + + $backoff = $this->backoffStrategy->compute($failures); + ++$failures; + + if ($this->retries >= 0 && $failures > $this->retries) { + throw $e; + } + + if ($fail !== null) { + $fail($e); + } + + if ($backoff > 0) { + usleep($backoff); + } + } + } + } +} diff --git a/src/Retry/Strategy/EqualBackoff.php b/src/Retry/Strategy/EqualBackoff.php new file mode 100644 index 00000000..6235d8c6 --- /dev/null +++ b/src/Retry/Strategy/EqualBackoff.php @@ -0,0 +1,37 @@ +backoff = $backoff; + } + + public function compute(int $failures): int + { + return $this->backoff; + } +} diff --git a/src/Retry/Strategy/ExponentialBackoff.php b/src/Retry/Strategy/ExponentialBackoff.php new file mode 100644 index 00000000..c16c146e --- /dev/null +++ b/src/Retry/Strategy/ExponentialBackoff.php @@ -0,0 +1,75 @@ +base = $base; + $this->cap = $cap; + $this->withJitter = $withJitter; + } + + /** + * {@inheritDoc} + */ + public function compute(int $failures): int + { + if ($this->withJitter) { + return min($this->cap, (mt_rand(0, mt_getrandmax() - 1) / mt_getrandmax()) * ($this->base * 2 ** $failures)); + } + + if ($this->cap > 0) { + return min($this->cap, $this->base * 2 ** $failures); + } + + return $this->base * 2 ** $failures; + } + + /** + * @return int + */ + public function getBase(): int + { + return $this->base; + } + + /** + * @return int + */ + public function getCap(): int + { + return $this->cap; + } +} diff --git a/src/Retry/Strategy/NoBackoff.php b/src/Retry/Strategy/NoBackoff.php new file mode 100644 index 00000000..fa29ea59 --- /dev/null +++ b/src/Retry/Strategy/NoBackoff.php @@ -0,0 +1,24 @@ +executeBypassingTransaction($command); } - return $this->connection->executeCommand($command); + $retry = $this->connection->getParameters()->retry; + + return $retry->callWithRetry( + function () use ($command) { + return $this->connection->executeCommand($command); + }, function (CommunicationException $e) { + $this->onFailCallback($e); + } + ); } /** @@ -99,10 +112,19 @@ abstract class NonClusterConnectionStrategy implements StrategyInterface /** * {@inheritDoc} + * @throws Throwable */ public function unwatch() { - return $this->connection->executeCommand(new UNWATCH()); + $retry = $this->connection->getParameters()->retry; + + return $retry->callWithRetry( + function () { + return $this->connection->executeCommand(new UNWATCH()); + }, function (CommunicationException $e) { + $this->onFailCallback($e); + } + ); } /** @@ -118,12 +140,20 @@ abstract class NonClusterConnectionStrategy implements StrategyInterface * * @param CommandInterface $command * @return BypassTransactionResponse - * @throws ServerException + * @throws ServerException|Throwable */ protected function executeBypassingTransaction(CommandInterface $command): BypassTransactionResponse { + $retry = $this->connection->getParameters()->retry; + try { - $response = $this->connection->executeCommand($command); + $response = $retry->callWithRetry( + function () use ($command) { + return $this->connection->executeCommand($command); + }, function (CommunicationException $e) { + $this->onFailCallback($e); + } + ); } catch (ServerException $exception) { if (!$this->connection instanceof RelayConnection) { throw $exception; @@ -146,4 +176,38 @@ abstract class NonClusterConnectionStrategy implements StrategyInterface return new BypassTransactionResponse($response); } + + /** + * Handle communication exception. + * + * @param CommunicationException $e + * @return void + */ + private function onFailCallback(CommunicationException $e) + { + $connection = $e->getConnection(); + + if ($connection instanceof NodeConnectionInterface) { + $connection->disconnect(); + + return; + } + + if ($e instanceof ConnectionException) { + $nodeConnection = $e->getConnection(); + + if ($nodeConnection) { + $nodeConnection->disconnect(); + $this->connection->remove($nodeConnection); + } + } + + if ($e instanceof TimeoutException) { + $nodeConnection = $e->getConnection(); + + if ($nodeConnection) { + $nodeConnection->disconnect(); + } + } + } } diff --git a/tests/PHPUnit/PredisConnectionTestCase.php b/tests/PHPUnit/PredisConnectionTestCase.php index fc81535f..9cfd2134 100644 --- a/tests/PHPUnit/PredisConnectionTestCase.php +++ b/tests/PHPUnit/PredisConnectionTestCase.php @@ -15,6 +15,7 @@ namespace Predis\Connection; use PHPUnit\Framework\MockObject\MockObject; use Predis\Command\CommandInterface; use Predis\Command\RawCommand; +use Predis\TimeoutException; use PredisTestCase; /** @@ -459,7 +460,7 @@ abstract class PredisConnectionTestCase extends PredisTestCase */ public function testThrowsExceptionOnReadWriteTimeout(): void { - $this->expectException('Predis\Connection\ConnectionException'); + $this->expectException(TimeoutException::class); $commands = $this->getCommandFactory(); diff --git a/tests/PHPUnit/PredisTestCase.php b/tests/PHPUnit/PredisTestCase.php index 77115211..ad983187 100644 --- a/tests/PHPUnit/PredisTestCase.php +++ b/tests/PHPUnit/PredisTestCase.php @@ -256,7 +256,7 @@ abstract class PredisTestCase extends PHPUnit\Framework\TestCase * * @return Client */ - protected function createClient(?array $parameters = null, ?array $options = null, ?bool $flushdb = true): Client + public function createClient(?array $parameters = null, ?array $options = null, ?bool $flushdb = true): Client { $parameters = array_merge( $this->getDefaultParametersArray(), diff --git a/tests/Predis/ClientTest.php b/tests/Predis/ClientTest.php index 9ee51700..979f372e 100644 --- a/tests/Predis/ClientTest.php +++ b/tests/Predis/ClientTest.php @@ -12,16 +12,26 @@ namespace Predis; +use Exception; use Iterator; use PHPUnit\Framework\MockObject\MockObject; use Predis\Command\Factory as CommandFactory; use Predis\Command\Processor\KeyPrefixProcessor; +use Predis\Command\RawCommand; +use Predis\Connection\Cluster\RedisCluster; +use Predis\Connection\Factory; use Predis\Connection\NodeConnectionInterface; use Predis\Connection\Parameters; use Predis\Connection\ParametersInterface; use Predis\Connection\Replication\MasterSlaveReplication; +use Predis\Connection\Resource\StreamFactoryInterface; +use Predis\Connection\StreamConnection; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; use PredisTestCase; +use Psr\Http\Message\StreamInterface; use ReflectionProperty; +use RuntimeException; use stdClass; class ClientTest extends PredisTestCase @@ -188,7 +198,7 @@ class ClientTest extends PredisTestCase */ public function testConstructorWithConnectionArgument(): void { - $factory = new Connection\Factory(); + $factory = new Factory(); $connection = $factory->create('tcp://localhost:7000'); $client = new Client($connection); @@ -211,7 +221,7 @@ class ClientTest extends PredisTestCase { $cluster = new Connection\Cluster\PredisCluster(new Parameters()); - $factory = new Connection\Factory(); + $factory = new Factory(); $cluster->add($factory->create('tcp://localhost:7000')); $cluster->add($factory->create('tcp://localhost:7001')); @@ -228,7 +238,7 @@ class ClientTest extends PredisTestCase { $replication = new MasterSlaveReplication(); - $factory = new Connection\Factory(); + $factory = new Factory(); $replication->add($factory->create('tcp://host1?alias=master')); $replication->add($factory->create('tcp://host2?alias=slave')); @@ -610,6 +620,10 @@ class ClientTest extends PredisTestCase ->expects($this->once()) ->method('executeCommand') ->willReturn($expectedResponse); + $connection + ->expects($this->once()) + ->method('getParameters') + ->willReturn(new Parameters()); $client = new Client($connection); $client->executeCommand($ping); @@ -628,6 +642,10 @@ class ClientTest extends PredisTestCase ->expects($this->once()) ->method('executeCommand') ->willReturn($expectedResponse); + $connection + ->expects($this->once()) + ->method('getParameters') + ->willReturn(new Parameters()); $client = new Client($connection, ['exceptions' => false]); $response = $client->executeCommand($ping); @@ -689,6 +707,11 @@ class ClientTest extends PredisTestCase ->with($this->isRedisCommand('PING')) ->willReturn($expectedResponse); + $connection + ->expects($this->once()) + ->method('getParameters') + ->willReturn(new Parameters()); + $client = new Client($connection); $client->ping(); } @@ -706,6 +729,10 @@ class ClientTest extends PredisTestCase ->method('executeCommand') ->with($this->isRedisCommand('PING')) ->willReturn($expectedResponse); + $connection + ->expects($this->once()) + ->method('getParameters') + ->willReturn(new Parameters()); $client = new Client($connection, ['exceptions' => false]); $response = $client->ping(); @@ -957,7 +984,7 @@ class ClientTest extends PredisTestCase */ public function testGetClientByMethodSupportsSelectingConnectionByCommand(): void { - $command = Command\RawCommand::create('GET', 'key'); + $command = RawCommand::create('GET', 'key'); $connection = $this->getMockBuilder('Predis\Connection\ConnectionInterface')->getMock(); $aggregate = $this->getMockBuilder('Predis\Connection\AggregateConnectionInterface') @@ -1185,6 +1212,10 @@ class ClientTest extends PredisTestCase ->expects($this->once()) ->method('executeCommand') ->willReturn(new Response\Status('QUEUED')); + $connection + ->expects($this->any()) + ->method('getParameters') + ->willReturn(new Parameters()); $callable = $this->getMockBuilder('stdClass') ->addMethods(['__invoke']) @@ -1308,6 +1339,50 @@ class ClientTest extends PredisTestCase $this->assertSame('127.0.0.1:6381', $iterator->key()); } + /** + * @group disconnected + */ + public function testExecuteCommandRetryCommandOnRetryableException() + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + + $mockStream + ->expects($this->atLeast(3)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(4)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + 1000 + ); + + $mockStream + ->expects($this->once()) + ->method('read') + ->withAnyParameters() + ->willReturn("+PONG\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $connection = new StreamConnection($parameters, $mockStreamFactory); + $client = new Client($connection); + $this->assertEquals('PONG', $client->ping()); + } + /** * @group connected * @group relay-incompatible @@ -1439,6 +1514,172 @@ class ClientTest extends PredisTestCase $this->assertEquals(1, $clientTestUser->acl->delUser('test_user')); } + /** + * @group connected + * @return void + * @requiresRedisVersion >= 7.0.0 + */ + public function testStandaloneNodeRetryCommandExecutionOnTimeoutException(): void + { + $retries = 0; + $mockDisconnect = function () use (&$retries) { + $streamConnection = new StreamConnection(new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ])); + $disconnectFunc = [$streamConnection, 'disconnect']; + ++$retries; + $disconnectFunc(); + }; + + $stubConnection = $this->getMockBuilder(StreamConnection::class) + ->setConstructorArgs([new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(100, 1000), 3), + 'read_write_timeout' => 0.1, + ])]) + ->onlyMethods(['disconnect']) + ->getMock(); + + $stubConnection + ->expects($this->exactly(7)) + ->method('disconnect') + ->willReturnCallback($mockDisconnect); + + $stubConnection->addConnectCommand(new RawCommand('auth', ['foobar'])); + + $client = new Client($stubConnection); + + $this->expectException(TimeoutException::class); + + $client->blmpop(3, ['random_key']); + $this->assertEquals(3, $retries); + } + + /** + * @group connected + * @group relay-incompatible + * @return void + * @requiresRedisVersion >= 7.0.0 + */ + public function testStandaloneNodeRetryCommandExecutionOnTimeoutExceptionIntegration(): void + { + // Retry used to wrap callback around, so we can count retries + $retry = new Retry(new ExponentialBackoff(100, 1000), 3); + $retriesCount = 0; + $retryWrapperFunc = function (callable $do, ?callable $fail = null) use ($retry, &$retriesCount) { + $failWrapperFunc = function (Exception $e) use (&$retriesCount, $fail) { + ++$retriesCount; + $fail($e); + }; + + return $retry->callWithRetry($do, $failWrapperFunc); + }; + + $mockRetry = $this->getMockBuilder(Retry::class) + ->setConstructorArgs([new ExponentialBackoff(100, 1000), 3]) + ->onlyMethods(['callWithRetry']) + ->getMock(); + + $mockRetry + ->expects($this->any()) + ->method('callWithRetry') + ->willReturnCallback($retryWrapperFunc); + + // Create a real connection with mocked retry and short read_write_timeout + $client = $this->createClient([ + 'retry' => $mockRetry, + 'read_write_timeout' => 0.1, + ]); + + $this->expectException(TimeoutException::class); + + try { + // blmpop with 3 second timeout will exceed the 0.1 second read_write_timeout + // causing TimeoutException to be thrown and retried 3 times before failing + $client->blmpop(3, ['random_key_that_does_not_exist']); + } finally { + $this->assertGreaterThanOrEqual(3, $retriesCount); + } + } + + /** + * @group connected + * @group cluster + * @return void + * @requiresRedisVersion >= 2.0.0 + */ + public function testClusterRetryCommandExecutionOnTimeoutException(): void + { + $defaultParams = $this->getDefaultParametersArray(); + $parsedParams = []; + + foreach ($defaultParams as $param) { + $parsedParam = Parameters::parse($param); + $parsedParam['retry'] = new Retry(new ExponentialBackoff(1000, 10000), 3); + $parsedParam['read_write_timeout'] = 0.1; + $parsedParams[] = Parameters::create($parsedParam); + } + + $retries = 0; + $mockDisconnect = function () use (&$retries) { + $streamConnection = new StreamConnection(new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ])); + $disconnectFunc = [$streamConnection, 'disconnect']; + ++$retries; + $disconnectFunc(); + }; + + $stubConnection1 = $this->getMockBuilder(StreamConnection::class) + ->setConstructorArgs([$parsedParams[0]]) + ->onlyMethods(['disconnect']) + ->getMock(); + + $stubConnection1 + ->expects($this->any()) + ->method('disconnect') + ->willReturnCallback($mockDisconnect); + + $stubConnection1->addConnectCommand(new RawCommand('auth', [$parsedParams[0]->password])); + + $stubConnection2 = $this->getMockBuilder(StreamConnection::class) + ->setConstructorArgs([$parsedParams[1]]) + ->onlyMethods(['disconnect']) + ->getMock(); + + $stubConnection2 + ->expects($this->any()) + ->method('disconnect') + ->willReturnCallback($mockDisconnect); + + $stubConnection2->addConnectCommand(new RawCommand('auth', [$parsedParams[1]->password])); + + $stubConnection3 = $this->getMockBuilder(StreamConnection::class) + ->setConstructorArgs([$parsedParams[2]]) + ->onlyMethods(['disconnect']) + ->getMock(); + + $stubConnection3 + ->expects($this->any()) + ->method('disconnect') + ->willReturnCallback($mockDisconnect); + + $stubConnection3->addConnectCommand(new RawCommand('auth', [$parsedParams[2]->password])); + + $mockFactory = $this->getMockBuilder(Factory::class)->getMock(); + $clusterConnection = new RedisCluster($mockFactory, $parsedParams[0]); + + $clusterConnection->add($stubConnection1); + $clusterConnection->add($stubConnection2); + $clusterConnection->add($stubConnection3); + + $client = new Client($clusterConnection); + + $this->expectException(TimeoutException::class); + + $client->blpop(['random_key'], 3); + $this->assertEquals(3, $retries); + } + // ******************************************************************** // // ---- HELPER METHODS ------------------------------------------------ // // ******************************************************************** // diff --git a/tests/Predis/Connection/Cluster/RedisClusterTest.php b/tests/Predis/Connection/Cluster/RedisClusterTest.php index 925c5de1..92fa61b4 100644 --- a/tests/Predis/Connection/Cluster/RedisClusterTest.php +++ b/tests/Predis/Connection/Cluster/RedisClusterTest.php @@ -19,8 +19,14 @@ use Predis\Command; use Predis\Connection; use Predis\Connection\FactoryInterface; use Predis\Connection\Parameters; +use Predis\Connection\Resource\StreamFactoryInterface; +use Predis\Connection\StreamConnection; use Predis\Response; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; use PredisTestCase; +use Psr\Http\Message\StreamInterface; +use RuntimeException; class RedisClusterTest extends PredisTestCase { @@ -1373,6 +1379,59 @@ class RedisClusterTest extends PredisTestCase $cluster->executeCommand($command); } + /** + * @medium + * @group disconnected + * @group slow + */ + public function testRetryCommandFailureOnCustomRetryConfiguration() + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + + $mockStream + ->expects($this->exactly(4)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(4)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + 1000, + 1000 + ); + + $mockStream + ->expects($this->once()) + ->method('read') + ->withAnyParameters() + ->willReturn("+OK\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $connection = new StreamConnection($parameters, $mockStreamFactory); + $cluster = new RedisCluster(new Connection\Factory(), $parameters); + $cluster->useClusterSlots(false); + $cluster->add($connection); + + $this->assertEquals( + 'OK', + $cluster->executeCommand(Command\RawCommand::create('SET', 1001)) + ); + } + /** * @medium * @group disconnected diff --git a/tests/Predis/Connection/ParametersTest.php b/tests/Predis/Connection/ParametersTest.php index 78bc4ee6..84c53731 100644 --- a/tests/Predis/Connection/ParametersTest.php +++ b/tests/Predis/Connection/ParametersTest.php @@ -12,6 +12,8 @@ namespace Predis\Connection; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\NoBackoff; use PredisTestCase; class ParametersTest extends PredisTestCase @@ -430,6 +432,7 @@ class ParametersTest extends PredisTestCase 'host' => '127.0.0.1', 'port' => 6379, 'protocol' => 2, + 'retry' => new Retry(new NoBackoff(), 0), ]; } diff --git a/tests/Predis/Connection/Replication/MasterSlaveReplicationTest.php b/tests/Predis/Connection/Replication/MasterSlaveReplicationTest.php index ee92f38e..c77e1c5c 100644 --- a/tests/Predis/Connection/Replication/MasterSlaveReplicationTest.php +++ b/tests/Predis/Connection/Replication/MasterSlaveReplicationTest.php @@ -15,9 +15,16 @@ namespace Predis\Connection\Replication; use PHPUnit\Framework\MockObject\MockObject; use Predis\Command; use Predis\Connection; +use Predis\Connection\Parameters; +use Predis\Connection\Resource\StreamFactoryInterface; +use Predis\Connection\StreamConnection; use Predis\Replication\ReplicationStrategy; use Predis\Response; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; use PredisTestCase; +use Psr\Http\Message\StreamInterface; +use RuntimeException; class MasterSlaveReplicationTest extends PredisTestCase { @@ -1489,6 +1496,58 @@ repl_backlog_histlen:12978 $replication->write($command1->serializeCommand() . $command2->serializeCommand() . $command3->serializeCommand()); } + /** + * @medium + * @group disconnected + * @group slow + */ + public function testRetryCommandFailureOnCustomRetryConfiguration() + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + 'role' => 'master', + ]); + + $mockStream + ->expects($this->exactly(4)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(4)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + 1000 + ); + + $mockStream + ->expects($this->once()) + ->method('read') + ->withAnyParameters() + ->willReturn("+OK\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $connection = new StreamConnection($parameters, $mockStreamFactory); + $replication = new MasterSlaveReplication(); + $replication->add($connection); + + $this->assertEquals( + 'OK', + $replication->executeCommand(Command\RawCommand::create('SET', 1001)) + ); + } + public function connectionsProvider(): array { return [ diff --git a/tests/Predis/Connection/Replication/SentinelReplicationTest.php b/tests/Predis/Connection/Replication/SentinelReplicationTest.php index 49b9171a..f0528ead 100644 --- a/tests/Predis/Connection/Replication/SentinelReplicationTest.php +++ b/tests/Predis/Connection/Replication/SentinelReplicationTest.php @@ -16,10 +16,17 @@ use Exception; use PHPUnit\Framework\MockObject\MockObject; use Predis\Command; use Predis\Connection; +use Predis\Connection\Parameters; +use Predis\Connection\Resource\StreamFactoryInterface; +use Predis\Connection\StreamConnection; use Predis\Replication; use Predis\Response; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; use PredisTestCase; +use Psr\Http\Message\StreamInterface; use ReflectionProperty; +use RuntimeException; class SentinelReplicationTest extends PredisTestCase { @@ -103,7 +110,7 @@ class SentinelReplicationTest extends PredisTestCase */ public function testConnectionParametersInstanceForSentinelConnectionIsNotModified(): void { - $originalParameters = Connection\Parameters::create( + $originalParameters = Parameters::create( 'tcp://127.0.0.1:5381?role=sentinel&database=1&password=secret' ); @@ -123,8 +130,8 @@ class SentinelReplicationTest extends PredisTestCase */ public function testConnectionParametersInstanceForSentinelConnectionIsNotModifiedEmptyPassword(): void { - $sentinel1 = Connection\Parameters::create('tcp://127.0.0.1:5381?role=sentinel&database=1&password='); - $sentinel2 = Connection\Parameters::create('tcp://127.0.0.1:5381?role=sentinel&database=1'); + $sentinel1 = Parameters::create('tcp://127.0.0.1:5381?role=sentinel&database=1&password='); + $sentinel2 = Parameters::create('tcp://127.0.0.1:5381?role=sentinel&database=1'); $replication1 = $this->getReplicationConnection('svc', [$sentinel1]); $replication2 = $this->getReplicationConnection('svc', [$sentinel2]); @@ -1774,6 +1781,73 @@ class SentinelReplicationTest extends PredisTestCase $this->assertSame($slave2, $replication->getCurrent()); } + /** + * @medium + * @group disconnected + * @group slow + */ + public function testRetryCommandFailureOnCustomRetryConfiguration() + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + 'role' => 'master', + ]); + + $mockStream + ->expects($this->exactly(3)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(6)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + 1000, + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + 1000, + 1000 + ); + + $mockStream + ->expects($this->exactly(7)) + ->method('read') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls("*1\r\n", "$5\r\n", "master\r\n", "*1\r\n", "$5\r\n", "master\r\n", "+OK\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $connection = new StreamConnection($parameters, $mockStreamFactory); + + $mockFactory = $this->getMockBuilder(Connection\Factory::class)->getMock(); + $mockFactory + ->expects($this->any()) + ->method('create') + ->willReturn($connection); + + $sentinel = $this->getMockSentinelConnection(); + $sentinel + ->expects($this->any()) + ->method('executeCommand') + ->willReturn(['127.0.0.1', '6381']); + + $replication = new SentinelReplication('src', [$sentinel], $mockFactory); + $replication->add($connection); + + $this->assertEquals( + 'OK', + $replication->executeCommand(Command\RawCommand::create('SET', 1001)) + ); + } + // ******************************************************************** // // ---- HELPER METHODS ------------------------------------------------ // // ******************************************************************** // diff --git a/tests/Predis/Connection/Resource/StreamTest.php b/tests/Predis/Connection/Resource/StreamTest.php index 0e18bdc5..f61f42ce 100644 --- a/tests/Predis/Connection/Resource/StreamTest.php +++ b/tests/Predis/Connection/Resource/StreamTest.php @@ -419,6 +419,18 @@ class StreamTest extends TestCase fclose($handle); } + /** + * @return void + */ + public function testWriteEmptyData(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + + $this->expectException(RuntimeException::class); + $stream->write(''); + } + public function writableModeProvider(): array { return [ diff --git a/tests/Predis/Pipeline/AtomicTest.php b/tests/Predis/Pipeline/AtomicTest.php index d99c6483..60855a1c 100644 --- a/tests/Predis/Pipeline/AtomicTest.php +++ b/tests/Predis/Pipeline/AtomicTest.php @@ -12,11 +12,15 @@ namespace Predis\Pipeline; +use Exception; use Predis\Client; use Predis\ClientInterface; use Predis\Command\Redis\PING; use Predis\Connection\Parameters; use Predis\Response; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; +use Predis\TimeoutException; use PredisTestCase; class AtomicTest extends PredisTestCase @@ -56,7 +60,7 @@ class AtomicTest extends PredisTestCase ); $connection - ->expects($this->once()) + ->expects($this->exactly(4)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -104,6 +108,10 @@ class AtomicTest extends PredisTestCase $queued, $queued ); + $connection + ->expects($this->exactly(3)) + ->method('getParameters') + ->willReturn(new Parameters(['protocol' => 2])); $pipeline = new Atomic(new Client($connection)); @@ -150,6 +158,10 @@ class AtomicTest extends PredisTestCase $queued, $error ); + $connection + ->expects($this->exactly(3)) + ->method('getParameters') + ->willReturn(new Parameters(['protocol' => 2])); $pipeline = new Atomic(new Client($connection)); @@ -191,6 +203,10 @@ class AtomicTest extends PredisTestCase ->willReturn( new Response\Error('ERR Test error') ); + $connection + ->expects($this->exactly(3)) + ->method('getParameters') + ->willReturn(new Parameters(['protocol' => 2])); $pipeline = new Atomic(new Client($connection)); @@ -236,7 +252,7 @@ class AtomicTest extends PredisTestCase ); $connection - ->expects($this->once()) + ->expects($this->exactly(4)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -265,6 +281,71 @@ class AtomicTest extends PredisTestCase $pipeline->execute(); } + /** + * @group disconnected + * @throws Exception + */ + public function testRetryStandalonePipelineOnRetryableErrors(): void + { + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + $mockConnection = $this->getMockConnection(); + + $mockConnection + ->expects($this->exactly(2)) + ->method('executeCommand') + ->withConsecutive( + [$this->isRedisCommand('MULTI')], + [$this->isRedisCommand('EXEC')] + ) + ->willReturnOnConsecutiveCalls( + new Response\Status('OK'), + ['PONG', 'PONG', 'PONG'] + ); + + $mockConnection + ->expects($this->exactly(4)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + 1000 + ); + + $mockConnection + ->expects($this->exactly(3)) + ->method('readResponse') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + "+QUEUED\r\n", + "+QUEUED\r\n", + "+QUEUED\r\n" + ); + + $mockConnection + ->expects($this->atLeast(3)) + ->method('disconnect') + ->withAnyParameters(); + + $mockConnection + ->expects($this->exactly(4)) + ->method('getParameters') + ->willReturn($parameters); + + $pipeline = new Atomic(new Client($mockConnection)); + + $responses = $pipeline->execute(function (Pipeline $pipe) { + $pipe->ping(); + $pipe->ping(); + $pipe->ping(); + }); + + $this->assertEquals(['PONG', 'PONG', 'PONG'], $responses); + } + /** * @group connected * @group relay-incompatible @@ -273,7 +354,7 @@ class AtomicTest extends PredisTestCase { $parameters = $this->getDefaultParametersArray(); - $client = $this->getClient( + $client = new Client( ["tcp://{$parameters['host']}:{$parameters['port']}?role=master&database={$parameters['database']}&password={$parameters['password']}"], ['replication' => 'predis'] ); diff --git a/tests/Predis/Pipeline/FireAndForgetTest.php b/tests/Predis/Pipeline/FireAndForgetTest.php index f903389d..88ceef17 100644 --- a/tests/Predis/Pipeline/FireAndForgetTest.php +++ b/tests/Predis/Pipeline/FireAndForgetTest.php @@ -12,10 +12,17 @@ namespace Predis\Pipeline; +use Exception; use Predis\Client; use Predis\ClientInterface; use Predis\Command\Redis\PING; +use Predis\Connection\Cluster\RedisCluster; +use Predis\Connection\Parameters; +use Predis\Connection\Replication\MasterSlaveReplication; use Predis\Response; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; +use Predis\TimeoutException; use PredisTestCase; class FireAndForgetTest extends PredisTestCase @@ -34,6 +41,10 @@ class FireAndForgetTest extends PredisTestCase $connection ->expects($this->never()) ->method('readResponse'); + $connection + ->expects($this->exactly(1)) + ->method('getParameters') + ->willReturn(new Parameters(['protocol' => 2])); $pipeline = new FireAndForget(new Client($connection)); @@ -68,6 +79,11 @@ class FireAndForgetTest extends PredisTestCase ->expects($this->never()) ->method('readResponse'); + $connection + ->expects($this->exactly(2)) + ->method('getParameters') + ->willReturn(new Parameters(['protocol' => 2])); + $pipeline = new FireAndForget(new Client($connection)); $pipeline->ping(); @@ -77,6 +93,147 @@ class FireAndForgetTest extends PredisTestCase $this->assertEmpty($pipeline->execute()); } + /** + * @group disconnected + * @throws Exception + */ + public function testRetryStandalonePipelineOnRetryableErrors(): void + { + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + $mockConnection = $this->getMockConnection(); + + $mockConnection + ->expects($this->exactly(4)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + 1000 + ); + + $mockConnection + ->expects($this->atLeast(3)) + ->method('disconnect') + ->withAnyParameters(); + + $mockConnection + ->expects($this->exactly(1)) + ->method('getParameters') + ->willReturn($parameters); + + $pipeline = new FireAndForget(new Client($mockConnection)); + + $pipeline->execute(function (Pipeline $pipe) { + $pipe->ping(); + $pipe->ping(); + $pipe->ping(); + }); + } + + /** + * @group disconnected + * @throws Exception + */ + public function testRetryClusterPipelineOnRetryableErrors(): void + { + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + $mockConnection = $this->getMockConnection(); + $mockClusterConnection = $this->getMockBuilder(RedisCluster::class) + ->disableOriginalConstructor()->getMock(); + + $mockConnection + ->expects($this->exactly(6)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + 1000, + 1000, + 1000 + ); + + $mockConnection + ->expects($this->atLeast(3)) + ->method('disconnect') + ->withAnyParameters(); + + $mockClusterConnection + ->expects($this->exactly(5)) + ->method('getParameters') + ->willReturn($parameters); + + $mockClusterConnection + ->expects($this->exactly(6)) + ->method('getConnectionByCommand') + ->willReturn($mockConnection); + + $pipeline = new FireAndForget(new Client($mockClusterConnection)); + + $pipeline->execute(function (Pipeline $pipe) { + $pipe->ping(); + $pipe->ping(); + $pipe->ping(); + }); + } + + /** + * @group disconnected + * @throws Exception + */ + public function testRetryReplicationPipelineOnRetryableErrors(): void + { + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + $mockConnection = $this->getMockConnection(); + $mockReplicationConnection = $this->getMockBuilder(MasterSlaveReplication::class) + ->disableOriginalConstructor()->getMock(); + + $mockConnection + ->expects($this->exactly(6)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + 1000, + 1000, + 1000 + ); + + $mockConnection + ->expects($this->atLeast(3)) + ->method('disconnect') + ->withAnyParameters(); + + $mockReplicationConnection + ->expects($this->exactly(5)) + ->method('getParameters') + ->willReturn($parameters); + + $mockReplicationConnection + ->expects($this->exactly(6)) + ->method('getConnectionByCommand') + ->willReturn($mockConnection); + + $pipeline = new FireAndForget(new Client($mockReplicationConnection)); + + $pipeline->execute(function (Pipeline $pipe) { + $pipe->ping(); + $pipe->ping(); + $pipe->ping(); + }); + } + /** * @group connected * @group cluster @@ -105,7 +262,7 @@ class FireAndForgetTest extends PredisTestCase { $parameters = $this->getDefaultParametersArray(); - $client = $this->getClient( + $client = new Client( ["tcp://{$parameters['host']}:{$parameters['port']}?role=master&database={$parameters['database']}&password={$parameters['password']}"], ['replication' => 'predis'] ); diff --git a/tests/Predis/Pipeline/PipelineTest.php b/tests/Predis/Pipeline/PipelineTest.php index e69a538e..61870bfe 100644 --- a/tests/Predis/Pipeline/PipelineTest.php +++ b/tests/Predis/Pipeline/PipelineTest.php @@ -20,9 +20,20 @@ use Predis\ClientInterface; use Predis\Command\CommandInterface; use Predis\Command\Redis\ECHO_; use Predis\Command\Redis\PING; +use Predis\Connection\Cluster\RedisCluster; +use Predis\Connection\Factory; use Predis\Connection\Parameters; +use Predis\Connection\Replication\MasterSlaveReplication; +use Predis\Connection\Resource\StreamFactoryInterface; +use Predis\Connection\StreamConnection; +use Predis\Replication\ReplicationStrategy; use Predis\Response; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; +use Predis\TimeoutException; use PredisTestCase; +use Psr\Http\Message\StreamInterface; +use RuntimeException; use stdClass; class PipelineTest extends PredisTestCase @@ -91,7 +102,7 @@ class PipelineTest extends PredisTestCase ->willReturn($object); $connection - ->expects($this->once()) + ->expects($this->exactly(2)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -119,7 +130,7 @@ class PipelineTest extends PredisTestCase ->willReturn($error); $connection - ->expects($this->once()) + ->expects($this->exactly(2)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -145,7 +156,7 @@ class PipelineTest extends PredisTestCase ->willReturn($error); $connection - ->expects($this->once()) + ->expects($this->exactly(2)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -238,7 +249,7 @@ class PipelineTest extends PredisTestCase ->willReturnCallback($this->getReadCallback()); $connection - ->expects($this->once()) + ->expects($this->exactly(2)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -301,7 +312,7 @@ class PipelineTest extends PredisTestCase ->willReturnCallback($this->getReadCallback()); $connection - ->expects($this->exactly(2)) + ->expects($this->exactly(4)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -333,15 +344,15 @@ class PipelineTest extends PredisTestCase ->expects($this->once()) ->method('switchToMaster'); $connection - ->expects($this->exactly(3)) + ->expects($this->exactly(6)) ->method('getConnectionByCommand') ->willReturn($nodeConnection); - $connection + $nodeConnection ->expects($this->exactly(3)) ->method('readResponse') ->willReturn($pong); $connection - ->expects($this->once()) + ->expects($this->exactly(3)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -429,7 +440,7 @@ class PipelineTest extends PredisTestCase ->method('readResponse') ->willReturnCallback($this->getReadCallback()); $connection - ->expects($this->once()) + ->expects($this->exactly(2)) ->method('getParameters') ->willReturn(new Parameters(['protocol' => 2])); @@ -478,6 +489,327 @@ class PipelineTest extends PredisTestCase $this->assertNull($responses); } + /** + * @group disconnected + * @throws Exception + */ + public function testRetryStandalonePipelineOnRetryableErrors(): void + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + + $mockStream + ->expects($this->atLeast(3)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(4)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + 1000 + ); + + $mockStream + ->expects($this->exactly(3)) + ->method('read') + ->withAnyParameters() + ->willReturn("+PONG\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $connection = new StreamConnection($parameters, $mockStreamFactory); + $pipeline = new Pipeline(new Client($connection)); + + $responses = $pipeline->execute(function (Pipeline $pipe) { + $pipe->ping(); + $pipe->ping(); + $pipe->ping(); + }); + + $this->assertEquals(['PONG', 'PONG', 'PONG'], $responses); + } + + /** + * @group disconnected + * @throws Exception + */ + public function testRetryClusterPipelineOnRetryableErrors(): void + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $mockConnectionFactory = $this->getMockBuilder(Factory::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + + $mockStream + ->expects($this->atLeast(3)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(6)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + 1000, + 1000, + 1000 + ); + + $mockStream + ->expects($this->exactly(3)) + ->method('read') + ->withAnyParameters() + ->willReturn("+OK\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $streamConnection = new StreamConnection($parameters, $mockStreamFactory); + $connection = new RedisCluster($mockConnectionFactory, $parameters); + $connection->add($streamConnection); + + $pipeline = new Pipeline(new Client($connection)); + + $responses = $pipeline->execute(function (Pipeline $pipe) { + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + }); + + $this->assertEquals(['OK', 'OK', 'OK'], $responses); + } + + /** + * @group disconnected + * @throws Exception + */ + public function testExecutePipelineInvokesOnAggregateConnectionFailCallbackOnConnectionException(): void + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $mockConnectionFactory = $this->getMockBuilder(Factory::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + + $streamConnection = new StreamConnection($parameters, $mockStreamFactory); + $connection = $this->getMockBuilder(RedisCluster::class) + ->setConstructorArgs([$mockConnectionFactory, $parameters]) + ->onlyMethods(['getConnectionByCommand', 'remove', 'askSlotMap']) + ->getMock(); + + // Disable useClusterSlots to avoid askSlotMap() calls + $connection->useClusterSlots = false; + + $mockStream + ->expects($this->atLeast(3)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(6)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new \Predis\Connection\ConnectionException($streamConnection, 'Connection failed')), + $this->throwException(new \Predis\Connection\ConnectionException($streamConnection, 'Connection failed')), + $this->throwException(new \Predis\Connection\ConnectionException($streamConnection, 'Connection failed')), + 1000, + 1000, + 1000 + ); + + $mockStream + ->expects($this->exactly(3)) + ->method('read') + ->withAnyParameters() + ->willReturn("+OK\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $connection + ->expects($this->exactly(9)) + ->method('getConnectionByCommand') + ->willReturn($streamConnection); + + // Verify that remove() is called on the aggregate connection during retry + $connection + ->expects($this->exactly(3)) + ->method('remove') + ->with($streamConnection); + + // Verify that askSlotMap() is NOT called since useClusterSlots is false + $connection + ->expects($this->never()) + ->method('askSlotMap'); + + $pipeline = new Pipeline(new Client($connection)); + + $responses = $pipeline->execute(function (Pipeline $pipe) { + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + }); + + $this->assertEquals(['OK', 'OK', 'OK'], $responses); + } + + /** + * @group disconnected + * @throws Exception + */ + public function testExecutePipelineInvokesOnAggregateConnectionFailCallbackOnTimeoutException(): void + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $mockConnectionFactory = $this->getMockBuilder(Factory::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + + $streamConnection = new StreamConnection($parameters, $mockStreamFactory); + $connection = $this->getMockBuilder(RedisCluster::class) + ->setConstructorArgs([$mockConnectionFactory, $parameters]) + ->onlyMethods(['getConnectionByCommand', 'remove']) + ->getMock(); + + // Disable useClusterSlots to avoid askSlotMap() calls + $connection->useClusterSlots = false; + + $mockStream + ->expects($this->atLeast(3)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(6)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new TimeoutException($streamConnection, 0)), + $this->throwException(new TimeoutException($streamConnection, 0)), + $this->throwException(new TimeoutException($streamConnection, 0)), + 1000, + 1000, + 1000 + ); + + $mockStream + ->expects($this->exactly(3)) + ->method('read') + ->withAnyParameters() + ->willReturn("+OK\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $connection + ->expects($this->exactly(9)) + ->method('getConnectionByCommand') + ->willReturn($streamConnection); + + // Verify that remove() is NOT called for TimeoutException + $connection + ->expects($this->never()) + ->method('remove'); + + $pipeline = new Pipeline(new Client($connection)); + + $responses = $pipeline->execute(function (Pipeline $pipe) { + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + }); + + $this->assertEquals(['OK', 'OK', 'OK'], $responses); + } + + /** + * @group disconnected + * @throws Exception + */ + public function testRetryReplicationPipelineOnRetryableErrors(): void + { + $mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + $mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + 'role' => 'master', + ]); + + $mockStream + ->expects($this->atLeast(3)) + ->method('close') + ->withAnyParameters(); + + $mockStream + ->expects($this->exactly(6)) + ->method('write') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + $this->throwException(new RuntimeException('', 2)), + 1000, + 1000, + 1000 + ); + + $mockStream + ->expects($this->exactly(3)) + ->method('read') + ->withAnyParameters() + ->willReturn("+OK\r\n"); + + $mockStreamFactory + ->expects($this->exactly(4)) + ->method('createStream') + ->withAnyParameters() + ->willReturn($mockStream); + + $streamConnection = new StreamConnection($parameters, $mockStreamFactory); + + $connection = new MasterSlaveReplication(new ReplicationStrategy()); + $connection->add($streamConnection); + + $pipeline = new Pipeline(new Client($connection)); + + $responses = $pipeline->execute(function (Pipeline $pipe) { + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + $pipe->set('key', 'value'); + }); + + $this->assertEquals(['OK', 'OK', 'OK'], $responses); + } + // ******************************************************************** // // ---- INTEGRATION TESTS --------------------------------------------- // // ******************************************************************** // @@ -648,6 +980,54 @@ class PipelineTest extends PredisTestCase $this->assertSameValues($expectedResults, $results); } + /** + * @group connected + * @group relay-incompatible + * @requiresRedisVersion >= 6.2.0 + * @return void + */ + public function testStandaloneRetryPipelineOnTimeoutException(): void + { + $client = $this->getClient([ + 'retry' => new Retry(new ExponentialBackoff(100, 1000), 3), + 'read_write_timeout' => 0.1, + ]); + + $this->expectException(TimeoutException::class); + + $client->pipeline(function (Pipeline $pipe) use (&$retries) { + $pipe->incr('test_key'); + $pipe->blpop('foo', 3); + }); + $this->assertEquals(3, $client->get('test_key')); + } + + /** + * @group connected + * @group cluster + * @group relay-incompatible + * @requiresRedisVersion >= 6.2.0 + * @return void + */ + public function testClusterRetryPipelineOnTimeoutException(): void + { + $retries = 0; + $client = $this->getClient([], [ + 'parameters' => [ + 'retry' => new Retry(new ExponentialBackoff(100, 1000), 3), + 'read_write_timeout' => 0.1, + ], + ]); + + $this->expectException(TimeoutException::class); + + $client->pipeline(function (Pipeline $pipe) use (&$retries) { + ++$retries; + $pipe->blpop('foo', 3); + }); + $this->assertEquals(3, $retries); + } + /** * @group connected * @group relay-incompatible @@ -656,7 +1036,7 @@ class PipelineTest extends PredisTestCase { $parameters = $this->getDefaultParametersArray(); - $client = $this->getClient( + $client = new Client( ["tcp://{$parameters['host']}:{$parameters['port']}?role=master&database={$parameters['database']}&password={$parameters['password']}"], ['replication' => 'predis'] ); diff --git a/tests/Predis/Retry/RetryTest.php b/tests/Predis/Retry/RetryTest.php new file mode 100644 index 00000000..ea4f24dd --- /dev/null +++ b/tests/Predis/Retry/RetryTest.php @@ -0,0 +1,154 @@ += $retries) { + return; + } + + ++$retriesCount; + throw new StreamInitException(); + }; + + $startTime = microtime(true); + $retry->callWithRetry($callable); + $executionTime = microtime(true) - $startTime; + + $this->assertEquals($retriesCount, $retries); + $this->assertEqualsWithDelta($expectedExecutionTime, $executionTime, $delta); + + $retry->updateRetriesCount(10); + $this->assertEquals(10, $retry->getRetries()); + } + + /** + * @group disconnected + * @return void + */ + public function testNoRetriesOnExcludedRetryableExceptions() + { + $retry = new Retry(new NoBackoff(), 3, [ConnectionException::class]); + $retriesCount = 0; + $callCount = 0; + + $doCallable = function () use (&$callCount) { + ++$callCount; + + if ($callCount <= 3) { + throw new RuntimeException(); + } elseif ($callCount <= 7) { + throw new ConnectionException( + $this->getMockBuilder(NodeConnectionInterface::class)->getMock() + ); + } else { + throw new StreamInitException(); + } + }; + + $failCallable = function () use (&$retriesCount) { + ++$retriesCount; + }; + + // Ensures that no retries happens on excluded exception. + while ($callCount < 3) { + try { + $retry->callWithRetry($doCallable, $failCallable); + } catch (Throwable $e) { + $this->assertInstanceOf(RuntimeException::class, $e); + $this->assertEquals(0, $retriesCount); + } + } + + // Ensures that retries happens on specified exception. + try { + $retry->callWithRetry($doCallable, $failCallable); + } catch (Throwable $e) { + $this->assertInstanceOf(ConnectionException::class, $e); + $this->assertEquals(3, $retriesCount); + } + + $retry->updateCatchableExceptions([StreamInitException::class]); + + // Ensures that retries happens on updated catchable exceptions. + try { + $retry->callWithRetry($doCallable, $failCallable); + } catch (Throwable $e) { + $this->assertInstanceOf(StreamInitException::class, $e); + $this->assertEquals(6, $retriesCount); + } + + $this->assertEquals(11, $callCount); + } + + public function strategyProvider(): array + { + return [ + 'NoBackoff' => [ + new NoBackoff(), + 3, + 1, + 1, + ], + 'EqualBackoff' => [ + new EqualBackoff(0.3 * 1000000), + 3, + 0.9, + 0.1, + ], + 'ExponentialBackoff - no jitter' => [ + new ExponentialBackoff(), + 3, + 0.112, + 0.08, + ], + 'ExponentialBackoff - with jitter' => [ + new ExponentialBackoff( + RetryStrategyInterface::DEFAULT_BASE, + RetryStrategyInterface::DEFAULT_CAP, + true + ), + 3, + 0.112, + 0.112, // Theoretically, jitter==0 might happen sequentially 3 times + ], + ]; + } +} diff --git a/tests/Predis/Retry/Strategy/EqualBackoffTest.php b/tests/Predis/Retry/Strategy/EqualBackoffTest.php new file mode 100644 index 00000000..ab1cd1c1 --- /dev/null +++ b/tests/Predis/Retry/Strategy/EqualBackoffTest.php @@ -0,0 +1,28 @@ +assertEquals(1, $backoff->compute(1)); + } +} diff --git a/tests/Predis/Retry/Strategy/ExponentialBackoffTest.php b/tests/Predis/Retry/Strategy/ExponentialBackoffTest.php new file mode 100644 index 00000000..49e1e888 --- /dev/null +++ b/tests/Predis/Retry/Strategy/ExponentialBackoffTest.php @@ -0,0 +1,70 @@ +assertLessThanOrEqual(RetryStrategyInterface::DEFAULT_CAP, $backoff->compute(100)); + + // Test default base + $this->assertGreaterThanOrEqual(RetryStrategyInterface::DEFAULT_BASE, $backoff->compute(0)); + + $interval = $backoff->compute(2); + + // Test between + $this->assertGreaterThanOrEqual(RetryStrategyInterface::DEFAULT_BASE, $interval); + $this->assertLessThanOrEqual(RetryStrategyInterface::DEFAULT_CAP, $interval); + + $backoff = new ExponentialBackoff(1000000, 10000000); + + // Test adjusted cap + $this->assertLessThanOrEqual(10000000, $backoff->compute(100)); + + // Test adjusted base + $this->assertGreaterThanOrEqual(1000000, $backoff->compute(0)); + + $backoff = new ExponentialBackoff(RetryStrategyInterface::DEFAULT_BASE, -1); + + // Test with no cap + $this->assertEquals(RetryStrategyInterface::DEFAULT_BASE * 2, $backoff->compute(1)); + + $backoff = new ExponentialBackoff( + RetryStrategyInterface::DEFAULT_BASE, + RetryStrategyInterface::DEFAULT_CAP, + true + ); + + $interval = $backoff->compute(0); + + // Test with jitter - default base + $this->assertGreaterThanOrEqual(0, $interval); + $this->assertLessThanOrEqual(RetryStrategyInterface::DEFAULT_BASE, $interval); + + $interval = $backoff->compute(6); + + // Test with jitter - default cap + $this->assertGreaterThanOrEqual(0, $interval); + $this->assertLessThanOrEqual(RetryStrategyInterface::DEFAULT_CAP, $interval); + } +} diff --git a/tests/Predis/Retry/Strategy/NoBackoffTest.php b/tests/Predis/Retry/Strategy/NoBackoffTest.php new file mode 100644 index 00000000..db59eda3 --- /dev/null +++ b/tests/Predis/Retry/Strategy/NoBackoffTest.php @@ -0,0 +1,28 @@ +assertEquals(0, $backoff->compute(1)); + } +} diff --git a/tests/Predis/Transaction/MultiExecTest.php b/tests/Predis/Transaction/MultiExecTest.php index d284c16a..eda76e18 100644 --- a/tests/Predis/Transaction/MultiExecTest.php +++ b/tests/Predis/Transaction/MultiExecTest.php @@ -20,6 +20,9 @@ use Predis\Command\CommandInterface; use Predis\Connection\NodeConnectionInterface; use Predis\Connection\Parameters; use Predis\Response; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; +use Predis\TimeoutException; use Predis\Transaction\Exception\TransactionException; use PredisTestCase; use RuntimeException; @@ -675,6 +678,55 @@ class MultiExecTest extends PredisTestCase $tx->multi()->echo('test')->exec(); } + /** + * @group disconnected + * @throws Exception + */ + public function testRetryReplicationPipelineOnRetryableErrors(): void + { + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + 'role' => 'master', + ]); + + $mockConnection = $this->getMockConnection(); + + $mockConnection + ->expects($this->any()) + ->method('getParameters') + ->willReturn($parameters); + + $mockConnection + ->expects($this->atLeast(3)) + ->method('disconnect') + ->withAnyParameters(); + + $mockConnection + ->expects($this->exactly(8)) + ->method('executeCommand') + ->withAnyParameters() + ->willReturnOnConsecutiveCalls( + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + $this->throwException(new TimeoutException($mockConnection)), + new Response\Status('OK'), + new Response\Status('QUEUED'), + new Response\Status('QUEUED'), + new Response\Status('QUEUED'), + ['OK', 'OK', 'OK'] + ); + + $tx = new MultiExec(new Client($mockConnection)); + + $responses = $tx->execute(function (MultiExec $tx) { + $tx->set('key', 'value'); + $tx->set('key', 'value'); + $tx->set('key', 'value'); + }); + + $this->assertEquals(['OK', 'OK', 'OK'], $responses); + } + // ******************************************************************** // // ---- INTEGRATION TESTS --------------------------------------------- // // ******************************************************************** // diff --git a/tests/Predis/Transaction/Strategy/NodeConnectionStrategyTest.php b/tests/Predis/Transaction/Strategy/NodeConnectionStrategyTest.php index 386cbec3..f7e637b2 100644 --- a/tests/Predis/Transaction/Strategy/NodeConnectionStrategyTest.php +++ b/tests/Predis/Transaction/Strategy/NodeConnectionStrategyTest.php @@ -15,6 +15,10 @@ namespace Predis\Transaction\Strategy; use PHPUnit\Framework\TestCase; use Predis\Command\CommandInterface; use Predis\Connection\NodeConnectionInterface; +use Predis\Connection\Parameters; +use Predis\Retry\Retry; +use Predis\Retry\Strategy\ExponentialBackoff; +use Predis\TimeoutException; use Predis\Transaction\MultiExecState; class NodeConnectionStrategyTest extends TestCase @@ -46,8 +50,72 @@ class NodeConnectionStrategyTest extends TestCase ->with($this->mockCommand) ->willReturn('OK'); + $this->mockConnection + ->expects($this->any()) + ->method('getParameters') + ->willReturn(new Parameters()); + $strategy = new NodeConnectionStrategy($this->mockConnection, new MultiExecState()); $this->assertEquals('OK', $strategy->executeCommand($this->mockCommand)); } + + /** + * @return void + */ + public function testUnwatch(): void + { + $this->mockConnection + ->expects($this->once()) + ->method('executeCommand') + ->with($this->callback(function ($command) { + return $command->getId() === 'UNWATCH'; + })) + ->willReturn('OK'); + + $this->mockConnection + ->expects($this->any()) + ->method('getParameters') + ->willReturn(new Parameters()); + + $strategy = new NodeConnectionStrategy($this->mockConnection, new MultiExecState()); + + $this->assertEquals('OK', $strategy->unwatch()); + } + + /** + * @return void + */ + public function testUnwatchWithRetries(): void + { + $parameters = new Parameters([ + 'retry' => new Retry(new ExponentialBackoff(1000, 10000), 3), + ]); + + $this->mockConnection + ->expects($this->exactly(4)) + ->method('executeCommand') + ->with($this->callback(function ($command) { + return $command->getId() === 'UNWATCH'; + })) + ->willReturnOnConsecutiveCalls( + $this->throwException(new TimeoutException($this->mockConnection)), + $this->throwException(new TimeoutException($this->mockConnection)), + $this->throwException(new TimeoutException($this->mockConnection)), + 'OK' + ); + + $this->mockConnection + ->expects($this->any()) + ->method('getParameters') + ->willReturn($parameters); + + $this->mockConnection + ->expects($this->exactly(3)) + ->method('disconnect'); + + $strategy = new NodeConnectionStrategy($this->mockConnection, new MultiExecState()); + + $this->assertEquals('OK', $strategy->unwatch()); + } } diff --git a/tests/Predis/Transaction/Strategy/ReplicationConnectionStrategyTest.php b/tests/Predis/Transaction/Strategy/ReplicationConnectionStrategyTest.php index db4575fe..6eefac65 100644 --- a/tests/Predis/Transaction/Strategy/ReplicationConnectionStrategyTest.php +++ b/tests/Predis/Transaction/Strategy/ReplicationConnectionStrategyTest.php @@ -14,6 +14,7 @@ namespace Predis\Transaction\Strategy; use PHPUnit\Framework\TestCase; use Predis\Command\CommandInterface; +use Predis\Connection\Parameters; use Predis\Connection\Replication\ReplicationInterface; use Predis\Transaction\MultiExecState; @@ -46,6 +47,11 @@ class ReplicationConnectionStrategyTest extends TestCase ->with($this->mockCommand) ->willReturn('OK'); + $this->mockConnection + ->expects($this->any()) + ->method('getParameters') + ->willReturn(new Parameters()); + $strategy = new ReplicationConnectionStrategy($this->mockConnection, new MultiExecState()); $this->assertEquals('OK', $strategy->executeCommand($this->mockCommand));