diff --git a/composer.json b/composer.json index 9e2662ed..9f7740bb 100644 --- a/composer.json +++ b/composer.json @@ -22,7 +22,8 @@ } ], "require": { - "php": "^7.2 || ^8.0" + "php": "^7.2 || ^8.0", + "psr/http-message": "^1.0|^2.0" }, "require-dev": { "friendsofphp/php-cs-fixer": "^3.3", diff --git a/examples/sharded_dispatcher_loop.php b/examples/sharded_dispatcher_loop.php index 49b2c1d0..5e6b4c7a 100644 --- a/examples/sharded_dispatcher_loop.php +++ b/examples/sharded_dispatcher_loop.php @@ -35,8 +35,8 @@ $client = new Client( 'tcp://127.0.0.1:6373?read_write_timeout=0', 'tcp://127.0.0.1:6374?read_write_timeout=0', ], [ - 'cluster' => 'redis', -]); + 'cluster' => 'redis', + ]); // 2. Run pub/sub loop. $pubSub = $client->pubSubLoop(); diff --git a/examples/sharded_pubsub_consumer.php b/examples/sharded_pubsub_consumer.php index a31a6454..f883b025 100644 --- a/examples/sharded_pubsub_consumer.php +++ b/examples/sharded_pubsub_consumer.php @@ -21,8 +21,8 @@ $client = new Client( 'tcp://127.0.0.1:6373?read_write_timeout=0', 'tcp://127.0.0.1:6374?read_write_timeout=0', ], [ - 'cluster' => 'redis', -]); + 'cluster' => 'redis', + ]); // 2. Run pub/sub loop. Sharded channels belongs to different shards. $pubSub = $client->pubSubLoop(); diff --git a/examples/transaction_using_cas.php b/examples/transaction_using_cas.php index edc85faa..2ad8e780 100644 --- a/examples/transaction_using_cas.php +++ b/examples/transaction_using_cas.php @@ -32,7 +32,7 @@ function zpop($client, $key) 'cas' => true, // Initialize with support for CAS operations 'watch' => $key, // Key that needs to be WATCHed to detect changes 'retry' => 3, // Number of retries on aborted transactions, after - // which the client bails out with an exception. + // which the client bails out with an exception. ]; $client->transaction($options, function ($tx) use ($key, &$element) { diff --git a/src/Connection/AbstractConnection.php b/src/Connection/AbstractConnection.php index 238414ac..7edd3dc8 100644 --- a/src/Connection/AbstractConnection.php +++ b/src/Connection/AbstractConnection.php @@ -12,10 +12,10 @@ namespace Predis\Connection; -use InvalidArgumentException; use Predis\Command\CommandInterface; use Predis\Command\RawCommand; use Predis\CommunicationException; +use Predis\Connection\Resource\Exception\StreamInitException; use Predis\Protocol\Parser\ParserStrategyResolver; use Predis\Protocol\Parser\Strategy\ParserStrategyInterface; use Predis\Protocol\ProtocolException; @@ -36,7 +36,7 @@ abstract class AbstractConnection implements NodeConnectionInterface */ protected $clientId; - private $resource; + protected $resource; private $cachedId; protected $parameters; @@ -51,7 +51,7 @@ abstract class AbstractConnection implements NodeConnectionInterface */ public function __construct(ParametersInterface $parameters) { - $this->parameters = $this->assertParameters($parameters); + $this->parameters = $parameters; $this->setParserStrategy(); } @@ -64,23 +64,6 @@ abstract class AbstractConnection implements NodeConnectionInterface $this->disconnect(); } - /** - * Checks some of the parameters used to initialize the connection. - * - * @param ParametersInterface $parameters Initialization parameters for the connection. - * - * @return ParametersInterface - * @throws InvalidArgumentException - */ - abstract protected function assertParameters(ParametersInterface $parameters); - - /** - * Creates the underlying resource used to communicate with Redis. - * - * @return mixed - */ - abstract protected function createResource(); - /** * {@inheritdoc} */ @@ -97,6 +80,14 @@ abstract class AbstractConnection implements NodeConnectionInterface return true; } + /** + * Creates a stream resource to communicate with Redis. + * + * @return mixed + * @throws StreamInitException + */ + abstract protected function createResource(); + /** * {@inheritdoc} */ @@ -156,10 +147,11 @@ abstract class AbstractConnection implements NodeConnectionInterface /** * Helper method to handle connection errors. * - * @param string $message Error message. - * @param int $code Error code. + * @param string $message Error message. + * @param int $code Error code. + * @throws CommunicationException */ - protected function onConnectionError($message, $code = 0) + protected function onConnectionError($message, $code = 0): void { CommunicationException::handle( new ConnectionException($this, "$message [{$this->getParameters()}]", $code) @@ -169,7 +161,8 @@ abstract class AbstractConnection implements NodeConnectionInterface /** * Helper method to handle protocol errors. * - * @param string $message Error message. + * @param string $message Error message. + * @throws CommunicationException */ protected function onProtocolError($message) { diff --git a/src/Connection/CompositeStreamConnection.php b/src/Connection/CompositeStreamConnection.php index 431663ca..12891b42 100644 --- a/src/Connection/CompositeStreamConnection.php +++ b/src/Connection/CompositeStreamConnection.php @@ -16,24 +16,28 @@ use InvalidArgumentException; use Predis\Command\CommandInterface; use Predis\Protocol\ProtocolProcessorInterface; use Predis\Protocol\Text\ProtocolProcessor as TextProtocolProcessor; +use Psr\Http\Message\StreamInterface; +use RuntimeException; /** * Connection abstraction to Redis servers based on PHP's stream that uses an * external protocol processor defining the protocol used for the communication. + * + * @method StreamInterface getResource() */ class CompositeStreamConnection extends StreamConnection implements CompositeConnectionInterface { protected $protocol; /** - * @param ParametersInterface $parameters Initialization parameters for the connection. - * @param ProtocolProcessorInterface $protocol Protocol processor. + * @param ParametersInterface $parameters Initialization parameters for the connection. + * @param ProtocolProcessorInterface|null $protocol Protocol processor. */ public function __construct( ParametersInterface $parameters, ?ProtocolProcessorInterface $protocol = null ) { - $this->parameters = $this->assertParameters($parameters); + parent::__construct($parameters); $this->protocol = $protocol ?: new TextProtocolProcessor(); } @@ -63,17 +67,21 @@ class CompositeStreamConnection extends StreamConnection implements CompositeCon } $value = ''; - $socket = $this->getResource(); + $stream = $this->getResource(); + + if ($stream->eof()) { + $this->onStreamError(new RuntimeException('Stream is already at the end'), ''); + } do { - $chunk = fread($socket, $length); - - if ($chunk === false || $chunk === '') { - $this->onConnectionError('Error while reading bytes from the server.'); + try { + $chunk = $stream->read($length); + } catch (RuntimeException $e) { + $this->onStreamError($e, 'Error while reading bytes from the server.'); } - $value .= $chunk; - } while (($length -= strlen($chunk)) > 0); + $value .= $chunk; // @phpstan-ignore-line + } while (($length -= strlen($chunk)) > 0); // @phpstan-ignore-line return $value; } @@ -84,16 +92,20 @@ class CompositeStreamConnection extends StreamConnection implements CompositeCon public function readLine() { $value = ''; - $socket = $this->getResource(); + $stream = $this->getResource(); + + if ($stream->eof()) { + $this->onStreamError(new RuntimeException('Stream is already at the end'), ''); + } do { - $chunk = fgets($socket); - - if ($chunk === false || $chunk === '') { - $this->onConnectionError('Error while reading line from the server.'); + try { + $chunk = $stream->read(-1); + } catch (RuntimeException $e) { + $this->onStreamError($e, 'Error while reading bytes from the server.'); } - $value .= $chunk; + $value .= $chunk; // @phpstan-ignore-line } while (substr($value, -2) !== "\r\n"); return substr($value, 0, -2); diff --git a/src/Connection/Resource/Exception/StreamInitException.php b/src/Connection/Resource/Exception/StreamInitException.php new file mode 100644 index 00000000..a42df853 --- /dev/null +++ b/src/Connection/Resource/Exception/StreamInitException.php @@ -0,0 +1,19 @@ +stream = $stream; + $metadata = stream_get_meta_data($this->stream); + $this->seekable = $metadata['seekable']; + $this->readable = (bool) preg_match(self::READABLE_MODES, $metadata['mode']); + $this->writable = (bool) preg_match(self::WRITABLE_MODES, $metadata['mode']); + } + + /** + * Closes the stream on garbage collection cleaning. + */ + public function __destruct() + { + $this->close(); + } + + /** + * {@inheritDoc} + */ + public function __toString(): string + { + if ($this->isSeekable()) { + $this->seek(0); + } + + return $this->getContents(); + } + + /** + * {@inheritDoc} + */ + public function close(): void + { + if (isset($this->stream)) { + fclose($this->stream); + } + + $this->detach(); + } + + /** + * {@inheritDoc} + */ + public function detach() + { + if (!isset($this->stream)) { + return null; + } + + $result = $this->stream; + unset($this->stream); + $this->readable = $this->writable = $this->seekable = false; + + return $result; + } + + /** + * {@inheritDoc} + */ + public function getSize(): ?int + { + if (!isset($this->stream)) { + return null; + } + + $stats = fstat($this->stream); + if (is_array($stats) && isset($stats['size'])) { + return $stats['size']; + } + + return null; + } + + /** + * {@inheritDoc} + */ + public function tell(): int + { + if (!isset($this->stream)) { + throw new RuntimeException('Stream is detached'); + } + + $result = ftell($this->stream); + + if ($result === false) { + throw new RuntimeException('Unable to determine stream position'); + } + + return $result; + } + + /** + * {@inheritDoc} + */ + public function eof(): bool + { + if (!isset($this->stream)) { + throw new RuntimeException('Stream is detached'); + } + + return feof($this->stream); + } + + /** + * {@inheritDoc} + */ + public function isSeekable(): bool + { + return $this->seekable; + } + + /** + * {@inheritDoc} + */ + public function seek(int $offset, int $whence = SEEK_SET): void + { + if (!isset($this->stream)) { + throw new RuntimeException('Stream is detached'); + } + + if (!$this->isSeekable()) { + throw new RuntimeException('Stream is not seekable'); + } + + if (fseek($this->stream, $offset, $whence) === -1) { + throw new RuntimeException("Unable to seek stream from offset {$offset} to whence {$whence}"); + } + } + + /** + * {@inheritDoc} + */ + public function rewind(): void + { + $this->seek(0); + } + + /** + * {@inheritDoc} + */ + public function isWritable(): bool + { + return $this->writable; + } + + /** + * {@inheritDoc} + * @throws RuntimeException + */ + public function write(string $string): int + { + if (!isset($this->stream)) { + throw new RuntimeException('Stream is detached'); + } + + if (!$this->isWritable()) { + throw new RuntimeException('Cannot write to a non-writable stream'); + } + + $result = fwrite($this->stream, $string); + + if ($result === false) { + throw new RuntimeException('Unable to write to stream', 1); + } + + return $result; + } + + /** + * {@inheritDoc} + */ + public function isReadable(): bool + { + return $this->readable; + } + + /** + * {@inheritDoc} + * @param int $length If length = -1, reads a stream line by line (e.g fgets()) + * @throws RuntimeException + */ + public function read(int $length): string + { + if (!isset($this->stream)) { + throw new RuntimeException('Stream is detached'); + } + + if (!$this->isReadable()) { + throw new RuntimeException('Cannot read from non-readable stream'); + } + + if ($length < -1) { + throw new RuntimeException('Length parameter cannot be negative'); + } + + if (0 === $length) { + return ''; + } + + if ($length === -1) { + $string = fgets($this->stream); + } else { + $string = fread($this->stream, $length); + } + + if (false === $string) { + throw new RuntimeException('Unable to read from stream', 1); + } + + return $string; + } + + /** + * {@inheritDoc} + */ + public function getContents(): string + { + if (!isset($this->stream)) { + throw new RuntimeException('Stream is detached'); + } + + if (!$this->isReadable()) { + throw new RuntimeException('Cannot read from non-readable stream'); + } + + return stream_get_contents($this->stream); + } + + /** + * {@inheritDoc} + */ + public function getMetadata(?string $key = null) + { + if (!isset($this->stream)) { + return null; + } + + if (!$key) { + return stream_get_meta_data($this->stream); + } + + $metadata = stream_get_meta_data($this->stream); + + return $metadata[$key] ?? null; + } +} diff --git a/src/Connection/Resource/StreamFactory.php b/src/Connection/Resource/StreamFactory.php new file mode 100644 index 00000000..c3fb2335 --- /dev/null +++ b/src/Connection/Resource/StreamFactory.php @@ -0,0 +1,227 @@ +assertParameters($parameters); + + switch ($parameters->scheme) { + case 'tcp': + case 'redis': + $stream = $this->tcpStreamInitializer($parameters); + break; + + case 'unix': + $stream = $this->unixStreamInitializer($parameters); + break; + + case 'tls': + case 'rediss': + $stream = $this->tlsStreamInitializer($parameters); + break; + + default: + throw new InvalidArgumentException("Invalid scheme: '{$parameters->scheme}'."); + } + + return new Stream($stream); + } + + /** + * Checks some parameters used to initialize the connection. + * + * @param ParametersInterface $parameters Initialization parameters for the connection. + * + * @return ParametersInterface + * @throws InvalidArgumentException + */ + protected function assertParameters(ParametersInterface $parameters): ParametersInterface + { + switch ($parameters->scheme) { + case 'tcp': + case 'redis': + case 'unix': + case 'tls': + case 'rediss': + break; + + default: + throw new InvalidArgumentException("Invalid scheme: '$parameters->scheme'."); + } + + return $parameters; + } + + /** + * Initializes a TCP stream resource. + * + * @param ParametersInterface $parameters Initialization parameters for the connection. + * + * @return resource + * @throws StreamInitException + */ + protected function tcpStreamInitializer(ParametersInterface $parameters) + { + if (!filter_var($parameters->host, FILTER_VALIDATE_IP, FILTER_FLAG_IPV6)) { + $address = "tcp://$parameters->host:$parameters->port"; + } else { + $address = "tcp://[$parameters->host]:$parameters->port"; + } + + $flags = STREAM_CLIENT_CONNECT; + + if (isset($parameters->async_connect) && $parameters->async_connect) { + $flags |= STREAM_CLIENT_ASYNC_CONNECT; + } + + if (isset($parameters->persistent)) { + if (false !== $persistent = filter_var($parameters->persistent, FILTER_VALIDATE_BOOLEAN, FILTER_NULL_ON_FAILURE)) { + $flags |= STREAM_CLIENT_PERSISTENT; + + if ($persistent === null) { + $address = "{$address}/{$parameters->persistent}"; + } + } + } + + return $this->createStreamSocket($parameters, $address, $flags); + } + + /** + * Initializes a UNIX stream resource. + * + * @param ParametersInterface $parameters Initialization parameters for the connection. + * + * @return resource + * @throws StreamInitException + */ + protected function unixStreamInitializer(ParametersInterface $parameters) + { + if (!isset($parameters->path)) { + throw new InvalidArgumentException('Missing UNIX domain socket path.'); + } + + $flags = STREAM_CLIENT_CONNECT; + + if (isset($parameters->persistent)) { + if (false !== $persistent = filter_var($parameters->persistent, FILTER_VALIDATE_BOOLEAN, FILTER_NULL_ON_FAILURE)) { + $flags |= STREAM_CLIENT_PERSISTENT; + + if ($persistent === null) { + throw new InvalidArgumentException( + 'Persistent connection IDs are not supported when using UNIX domain sockets.' + ); + } + } + } + + return $this->createStreamSocket($parameters, "unix://{$parameters->path}", $flags); + } + + /** + * Initializes a SSL-encrypted TCP stream resource. + * + * @param ParametersInterface $parameters Initialization parameters for the connection. + * + * @return resource + * @throws StreamInitException + */ + protected function tlsStreamInitializer(ParametersInterface $parameters) + { + $resource = $this->tcpStreamInitializer($parameters); + $metadata = stream_get_meta_data($resource); + + // Detect if crypto mode is already enabled for this stream (PHP >= 7.0.0). + if (isset($metadata['crypto'])) { + return $resource; + } + + if (isset($parameters->ssl) && is_array($parameters->ssl)) { + $options = $parameters->ssl; + } else { + $options = []; + } + + if (!isset($options['crypto_type'])) { + $options['crypto_type'] = STREAM_CRYPTO_METHOD_TLS_CLIENT; + } + + if (!stream_context_set_option($resource, ['ssl' => $options])) { + $this->onInitializationError($resource, $parameters, 'Error while setting SSL context options'); + } + + if (!stream_socket_enable_crypto($resource, true, $options['crypto_type'])) { + $this->onInitializationError($resource, $parameters, 'Error while switching to encrypted communication'); + } + + return $resource; + } + + /** + * Creates a connected stream socket resource. + * + * @param ParametersInterface $parameters Connection parameters. + * @param string $address Address for stream_socket_client(). + * @param int $flags Flags for stream_socket_client(). + * + * @return resource + * @throws StreamInitException + */ + protected function createStreamSocket(ParametersInterface $parameters, $address, $flags) + { + $timeout = (isset($parameters->timeout) ? (float) $parameters->timeout : 5.0); + $context = stream_context_create(['socket' => ['tcp_nodelay' => (bool) $parameters->tcp_nodelay]]); + + if (!$resource = @stream_socket_client($address, $errno, $errstr, $timeout, $flags, $context)) { + $this->onInitializationError($resource, $parameters, trim($errstr), $errno); + } + + if (isset($parameters->read_write_timeout)) { + $rwtimeout = (float) $parameters->read_write_timeout; + $rwtimeout = $rwtimeout > 0 ? $rwtimeout : -1; + $timeoutSeconds = floor($rwtimeout); + $timeoutUSeconds = ($rwtimeout - $timeoutSeconds) * 1000000; + stream_set_timeout($resource, $timeoutSeconds, $timeoutUSeconds); + } + + return $resource; + } + + /** + * Helper method to handle connection errors. + * + * @param string $message Error message. + * @param int $code Error code. + * @throws StreamInitException + */ + protected function onInitializationError($stream, ParametersInterface $parameters, string $message, int $code = 0): void + { + if (is_resource($stream)) { + fclose($stream); + } + + throw new StreamInitException("$message [{$parameters}]", $code); + } +} diff --git a/src/Connection/Resource/StreamFactoryInterface.php b/src/Connection/Resource/StreamFactoryInterface.php new file mode 100644 index 00000000..ef5286d6 --- /dev/null +++ b/src/Connection/Resource/StreamFactoryInterface.php @@ -0,0 +1,27 @@ +streamFactory = $factory ?? new StreamFactory(); + } + /** * Disconnects from the server and destroys the underlying resource when the * garbage collector kicks in only if the connection has not been marked as @@ -57,174 +80,9 @@ class StreamConnection extends AbstractConnection /** * {@inheritdoc} */ - protected function assertParameters(ParametersInterface $parameters) + protected function createResource(): StreamInterface { - switch ($parameters->scheme) { - case 'tcp': - case 'redis': - case 'unix': - case 'tls': - case 'rediss': - break; - - default: - throw new InvalidArgumentException("Invalid scheme: '$parameters->scheme'."); - } - - return $parameters; - } - - /** - * {@inheritdoc} - */ - protected function createResource() - { - switch ($this->parameters->scheme) { - case 'tcp': - case 'redis': - return $this->tcpStreamInitializer($this->parameters); - - case 'unix': - return $this->unixStreamInitializer($this->parameters); - - case 'tls': - case 'rediss': - return $this->tlsStreamInitializer($this->parameters); - - default: - throw new InvalidArgumentException("Invalid scheme: '{$this->parameters->scheme}'."); - } - } - - /** - * Creates a connected stream socket resource. - * - * @param ParametersInterface $parameters Connection parameters. - * @param string $address Address for stream_socket_client(). - * @param int $flags Flags for stream_socket_client(). - * - * @return resource - */ - protected function createStreamSocket(ParametersInterface $parameters, $address, $flags) - { - $timeout = (isset($parameters->timeout) ? (float) $parameters->timeout : 5.0); - $context = stream_context_create(['socket' => ['tcp_nodelay' => (bool) $parameters->tcp_nodelay]]); - - if (!$resource = @stream_socket_client($address, $errno, $errstr, $timeout, $flags, $context)) { - $this->onConnectionError(trim($errstr), $errno); - } - - if (isset($parameters->read_write_timeout)) { - $rwtimeout = (float) $parameters->read_write_timeout; - $rwtimeout = $rwtimeout > 0 ? $rwtimeout : -1; - $timeoutSeconds = floor($rwtimeout); - $timeoutUSeconds = ($rwtimeout - $timeoutSeconds) * 1000000; - stream_set_timeout($resource, $timeoutSeconds, $timeoutUSeconds); - } - - return $resource; - } - - /** - * Initializes a TCP stream resource. - * - * @param ParametersInterface $parameters Initialization parameters for the connection. - * - * @return resource - */ - protected function tcpStreamInitializer(ParametersInterface $parameters) - { - if (!filter_var($parameters->host, FILTER_VALIDATE_IP, FILTER_FLAG_IPV6)) { - $address = "tcp://$parameters->host:$parameters->port"; - } else { - $address = "tcp://[$parameters->host]:$parameters->port"; - } - - $flags = STREAM_CLIENT_CONNECT; - - if (isset($parameters->async_connect) && $parameters->async_connect) { - $flags |= STREAM_CLIENT_ASYNC_CONNECT; - } - - if (isset($parameters->persistent)) { - if (false !== $persistent = filter_var($parameters->persistent, FILTER_VALIDATE_BOOLEAN, FILTER_NULL_ON_FAILURE)) { - $flags |= STREAM_CLIENT_PERSISTENT; - - if ($persistent === null) { - $address = "{$address}/{$parameters->persistent}"; - } - } - } - - return $this->createStreamSocket($parameters, $address, $flags); - } - - /** - * Initializes a UNIX stream resource. - * - * @param ParametersInterface $parameters Initialization parameters for the connection. - * - * @return resource - */ - protected function unixStreamInitializer(ParametersInterface $parameters) - { - if (!isset($parameters->path)) { - throw new InvalidArgumentException('Missing UNIX domain socket path.'); - } - - $flags = STREAM_CLIENT_CONNECT; - - if (isset($parameters->persistent)) { - if (false !== $persistent = filter_var($parameters->persistent, FILTER_VALIDATE_BOOLEAN, FILTER_NULL_ON_FAILURE)) { - $flags |= STREAM_CLIENT_PERSISTENT; - - if ($persistent === null) { - throw new InvalidArgumentException( - 'Persistent connection IDs are not supported when using UNIX domain sockets.' - ); - } - } - } - - return $this->createStreamSocket($parameters, "unix://{$parameters->path}", $flags); - } - - /** - * Initializes a SSL-encrypted TCP stream resource. - * - * @param ParametersInterface $parameters Initialization parameters for the connection. - * - * @return resource - */ - protected function tlsStreamInitializer(ParametersInterface $parameters) - { - $resource = $this->tcpStreamInitializer($parameters); - $metadata = stream_get_meta_data($resource); - - // Detect if crypto mode is already enabled for this stream (PHP >= 7.0.0). - if (isset($metadata['crypto'])) { - return $resource; - } - - if (isset($parameters->ssl) && is_array($parameters->ssl)) { - $options = $parameters->ssl; - } else { - $options = []; - } - - if (!isset($options['crypto_type'])) { - $options['crypto_type'] = STREAM_CRYPTO_METHOD_TLS_CLIENT; - } - - if (!stream_context_set_option($resource, ['ssl' => $options])) { - $this->onConnectionError('Error while setting SSL context options'); - } - - if (!stream_socket_enable_crypto($resource, true, $options['crypto_type'])) { - $this->onConnectionError('Error while switching to encrypted communication'); - } - - return $resource; + return $this->streamFactory->createStream($this->parameters); } /** @@ -247,51 +105,56 @@ class StreamConnection extends AbstractConnection public function disconnect() { if ($this->isConnected()) { - $resource = $this->getResource(); - if (is_resource($resource)) { - fclose($resource); - } + $this->getResource()->close(); + parent::disconnect(); } } /** * {@inheritDoc} + * @throws CommunicationException */ public function write(string $buffer): void { - $socket = $this->getResource(); + $stream = $this->getResource(); while (($length = strlen($buffer)) > 0) { - $written = is_resource($socket) ? @fwrite($socket, $buffer) : false; + try { + $written = $stream->write($buffer); + } catch (RuntimeException $e) { + $this->onStreamError($e, 'Error while writing bytes to the server.'); + } - if ($length === $written) { + if ($length === $written) { // @phpstan-ignore-line return; } - if ($written === false || $written === 0) { - $this->onConnectionError('Error while writing bytes to the server.'); - } - - $buffer = substr($buffer, $written); + $buffer = substr($buffer, $written); // @phpstan-ignore-line } } /** * {@inheritdoc} * @throws PushNotificationException + * @throws StreamInitException|CommunicationException */ public function read() { - $socket = $this->getResource(); - $chunk = fgets($socket); + $stream = $this->getResource(); - if ($chunk === false || $chunk === '') { - $this->onConnectionError('Error while reading line from the server.'); + if ($stream->eof()) { + $this->onStreamError(new RuntimeException('Stream is already at the end'), ''); } try { - $parsedData = $this->parserStrategy->parseData($chunk); + $chunk = $stream->read(-1); + } catch (RuntimeException $e) { + $this->onStreamError($e, 'Error while reading line from the server.'); + } + + try { + $parsedData = $this->parserStrategy->parseData($chunk); // @phpstan-ignore-line } catch (UnexpectedTypeException $e) { $this->onProtocolError("Unknown response prefix: '{$e->getType()}'."); @@ -321,17 +184,17 @@ class StreamConnection extends AbstractConnection return $data; case Resp2Strategy::TYPE_BULK_STRING: - $bulkData = $this->readByChunks($socket, $parsedData['value']); + $bulkData = $this->readByChunks($stream, $parsedData['value']); return substr($bulkData, 0, -2); case Resp3Strategy::TYPE_VERBATIM_STRING: - $bulkData = $this->readByChunks($socket, $parsedData['value']); + $bulkData = $this->readByChunks($stream, $parsedData['value']); return substr($bulkData, $parsedData['offset'], -2); case Resp3Strategy::TYPE_BLOB_ERROR: - $errorMessage = $this->readByChunks($socket, $parsedData['value']); + $errorMessage = $this->readByChunks($stream, $parsedData['value']); return new Error(substr($errorMessage, 0, -2)); @@ -376,40 +239,30 @@ class StreamConnection extends AbstractConnection */ public function hasDataToRead(): bool { - $resource = $this->getResource(); - - if ($resource) { - $resourceArray = [$resource]; - $write = null; - $except = null; - $num = stream_select($resourceArray, $write, $except, 0); - - return $num > 0; - } - - return false; + return !$this->getResource()->eof(); } /** * Reads given resource split on chunks with given size. * - * @param $resource - * @param int $chunkSize + * @param StreamInterface $stream + * @param int $chunkSize * @return string + * @throws CommunicationException */ - private function readByChunks($resource, int $chunkSize): string + private function readByChunks(StreamInterface $stream, int $chunkSize): string { $string = ''; $bytesLeft = ($chunkSize += 2); do { - $chunk = is_resource($resource) ? fread($resource, min($bytesLeft, 4096)) : false; - - if ($chunk === false || $chunk === '') { - $this->onConnectionError('Error while reading bytes from the server.'); + try { + $chunk = $stream->read(min($bytesLeft, 4096)); + } catch (RuntimeException $e) { + $this->onStreamError($e, 'Error while reading bytes from the server.'); } - $string .= $chunk; + $string .= $chunk; // @phpstan-ignore-line $bytesLeft = $chunkSize - strlen($string); } while ($bytesLeft > 0); @@ -419,9 +272,10 @@ class StreamConnection extends AbstractConnection /** * Handle response from on-connect command. * - * @param $response - * @param CommandInterface $command + * @param $response + * @param CommandInterface $command * @return void + * @throws CommunicationException */ private function handleOnConnectResponse($response, CommandInterface $command): void { @@ -448,6 +302,7 @@ class StreamConnection extends AbstractConnection * @param ErrorResponseInterface $error * @param CommandInterface $failedCommand * @return void + * @throws CommunicationException */ private function handleError(ErrorResponseInterface $error, CommandInterface $failedCommand): void { @@ -474,4 +329,21 @@ class StreamConnection extends AbstractConnection $this->onConnectionError("Failed: {$error->getMessage()}"); } + + /** + * Handles stream-related exceptions. + * + * @param RuntimeException $e + * @param string|null $message + * @throws RuntimeException|CommunicationException + */ + protected function onStreamError(RuntimeException $e, ?string $message = null) + { + // Code = 1 represents issues related to read/write operation. + if ($e->getCode() === 1) { + $this->onConnectionError($message); + } + + throw $e; + } } diff --git a/tests/PHPUnit/PredisConnectionTestCase.php b/tests/PHPUnit/PredisConnectionTestCase.php index d5c7f30d..6bd03640 100644 --- a/tests/PHPUnit/PredisConnectionTestCase.php +++ b/tests/PHPUnit/PredisConnectionTestCase.php @@ -81,17 +81,6 @@ abstract class PredisConnectionTestCase extends PredisTestCase $this->assertInstanceOf('Predis\Connection\NodeConnectionInterface', $connection); } - /** - * @group disconnected - */ - public function testThrowsExceptionOnInvalidScheme(): void - { - $this->expectException('InvalidArgumentException'); - $this->expectExceptionMessage("Invalid scheme: 'udp'"); - - $this->createConnectionWithParams(['scheme' => 'udp']); - } - /** * @group disconnected */ @@ -201,7 +190,8 @@ abstract class PredisConnectionTestCase extends PredisTestCase $connection = $this->createConnection(); $this->assertFalse($connection->isConnected()); - $this->assertIsResource($connection->getResource()); + $connection->connect(); + $this->assertTrue($connection->isConnected()); } @@ -472,57 +462,6 @@ abstract class PredisConnectionTestCase extends PredisTestCase $this->assertSame(['foo', 'hoge', 'lol'], $connection->read()); } - /** - * @group connected - * @group slow - */ - public function testThrowsExceptionOnConnectionTimeout(): void - { - $this->expectException('Predis\Connection\ConnectionException'); - $this->expectExceptionMessageMatches('/.* \[tcp:\/\/169.254.10.10:6379\]/'); - - $connection = $this->createConnectionWithParams([ - 'host' => '169.254.10.10', - 'timeout' => 0.1, - ], false); - - $connection->connect(); - } - - /** - * @group connected - * @group slow - */ - public function testThrowsExceptionOnConnectionTimeoutIPv6(): void - { - $this->expectException('Predis\Connection\ConnectionException'); - $this->expectExceptionMessageMatches('/.* \[tcp:\/\/\[0:0:0:0:0:ffff:a9fe:a0a\]:6379\]/'); - - $connection = $this->createConnectionWithParams([ - 'host' => '0:0:0:0:0:ffff:a9fe:a0a', - 'timeout' => 0.1, - ], false); - - $connection->connect(); - } - - /** - * @group connected - * @group slow - */ - public function testThrowsExceptionOnUnixDomainSocketNotFound(): void - { - $this->expectException('Predis\Connection\ConnectionException'); - $this->expectExceptionMessageMatches('/.* \[unix:\/tmp\/nonexistent\/redis\.sock]/'); - - $connection = $this->createConnectionWithParams([ - 'scheme' => 'unix', - 'path' => '/tmp/nonexistent/redis.sock', - ], false); - - $connection->connect(); - } - /** * @group connected * @group slow @@ -540,23 +479,6 @@ abstract class PredisConnectionTestCase extends PredisTestCase $connection->executeCommand($commands->create('brpop', ['foo', 3])); } - /** - * @medium - * @group connected - */ - public function testThrowsExceptionOnProtocolDesynchronizationErrors(): void - { - $this->expectException('Predis\Protocol\ProtocolException'); - - $connection = $this->createConnection(); - $stream = $connection->getResource(); - - $connection->writeRequest($this->getCommandFactory()->create('ping')); - fread($stream, 1); - - $connection->read(); - } - // ******************************************************************** // // ---- HELPER METHODS ------------------------------------------------ // // ******************************************************************** // @@ -590,11 +512,11 @@ abstract class PredisConnectionTestCase extends PredisTestCase * This assertion will trigger a connect() operation if the connection has * not been open yet. * - * @param NodeConnectionInterface $connection Connection instance + * @param resource $resource */ - protected function assertPersistentConnection(NodeConnectionInterface $connection): void + protected function assertPersistentConnection($resource): void { - $this->assertSame('persistent stream', get_resource_type($connection->getResource())); + $this->assertSame('persistent stream', get_resource_type($resource)); } /** @@ -603,11 +525,11 @@ abstract class PredisConnectionTestCase extends PredisTestCase * This assertion will trigger a connect() operation if the connection has * not been open yet. * - * @param NodeConnectionInterface $connection Connection instance + * @param resource $resource */ - protected function assertNonPersistentConnection(NodeConnectionInterface $connection): void + protected function assertNonPersistentConnection($resource): void { - $this->assertSame('stream', get_resource_type($connection->getResource())); + $this->assertSame('stream', get_resource_type($resource)); } /** diff --git a/tests/Predis/Collection/Iterator/HashKeyTest.php b/tests/Predis/Collection/Iterator/HashKeyTest.php index 4a8c0d5e..b4f4eb71 100644 --- a/tests/Predis/Collection/Iterator/HashKeyTest.php +++ b/tests/Predis/Collection/Iterator/HashKeyTest.php @@ -63,7 +63,7 @@ class HashKeyTest extends PredisTestCase ->with('key:hash', 0, []) ->willReturn( [0, [], - ]); + ]); $iterator = new HashKey($client, 'key:hash'); diff --git a/tests/Predis/Command/Redis/GEOSEARCH_Test.php b/tests/Predis/Command/Redis/GEOSEARCH_Test.php index 281a83ea..2582f0ed 100644 --- a/tests/Predis/Command/Redis/GEOSEARCH_Test.php +++ b/tests/Predis/Command/Redis/GEOSEARCH_Test.php @@ -209,7 +209,7 @@ class GEOSEARCH_Test extends PredisCommandTestCase [ 'member1' => ['lng' => 1.1, 'lat' => 2.2], 'member2' => ['lng' => 2.2, 'lat' => 3.3], - 'member3' => ['lng' => 3.3, 'lat' => 4.4]], + 'member3' => ['lng' => 3.3, 'lat' => 4.4], ], ], 'with WITHDIST modifier' => [ [['member1', '111.111'], ['member2', '222.222'], ['member3', '333.333']], diff --git a/tests/Predis/Command/Redis/TimeSeries/TSINFO_Test.php b/tests/Predis/Command/Redis/TimeSeries/TSINFO_Test.php index f4a5ff41..775a6aff 100644 --- a/tests/Predis/Command/Redis/TimeSeries/TSINFO_Test.php +++ b/tests/Predis/Command/Redis/TimeSeries/TSINFO_Test.php @@ -70,7 +70,7 @@ class TSINFO_Test extends PredisCommandTestCase $redis = $this->getClient(); $expectedResponse = ['totalSamples', 0, 'memoryUsage', 4239, 'firstTimestamp', 0, 'lastTimestamp', 0, 'retentionTime', 60000, 'chunkCount', 1, 'chunkSize', 4096, 'chunkType', 'compressed', 'duplicatePolicy', - 'max', 'labels', [['sensor_id', '2'], ['area_id', '32']], 'sourceKey', null, 'rules', []]; + 'max', 'labels', [['sensor_id', '2'], ['area_id', '32']], 'sourceKey', null, 'rules', [], ]; $arguments = (new CreateArguments()) ->retentionMsecs(60000) @@ -95,7 +95,7 @@ class TSINFO_Test extends PredisCommandTestCase $redis = $this->getResp3Client(); $expectedResponse = ['totalSamples' => 0, 'memoryUsage' => 4239, 'firstTimestamp' => 0, 'lastTimestamp' => 0, 'retentionTime' => 60000, 'chunkCount' => 1, 'chunkSize' => 4096, 'chunkType' => 'compressed', - 'duplicatePolicy' => 'max', 'labels' => ['sensor_id' => '2', 'area_id' => '32'], 'sourceKey' => null, 'rules' => []]; + 'duplicatePolicy' => 'max', 'labels' => ['sensor_id' => '2', 'area_id' => '32'], 'sourceKey' => null, 'rules' => [], ]; $arguments = (new CreateArguments()) ->retentionMsecs(60000) diff --git a/tests/Predis/Command/Redis/TimeSeries/TSMRANGE_Test.php b/tests/Predis/Command/Redis/TimeSeries/TSMRANGE_Test.php index 03ce7afe..1fe1418a 100644 --- a/tests/Predis/Command/Redis/TimeSeries/TSMRANGE_Test.php +++ b/tests/Predis/Command/Redis/TimeSeries/TSMRANGE_Test.php @@ -117,17 +117,17 @@ class TSMRANGE_Test extends PredisCommandTestCase { $redis = $this->getResp3Client(); $expectedResponse = [ - 'type=stock' => [ - ['type' => 'stock'], - ['reducers' => ['max']], - ['sources' => ['stock:A', 'stock:B']], - [ - [1000, 120], - [1010, 110], - [1020, 120], - ], + 'type=stock' => [ + ['type' => 'stock'], + ['reducers' => ['max']], + ['sources' => ['stock:A', 'stock:B']], + [ + [1000, 120], + [1010, 110], + [1020, 120], ], - ]; + ], + ]; $this->assertEquals( 'OK', diff --git a/tests/Predis/Command/Redis/TimeSeries/TSMREVRANGE_Test.php b/tests/Predis/Command/Redis/TimeSeries/TSMREVRANGE_Test.php index 988d9f01..26b1e464 100644 --- a/tests/Predis/Command/Redis/TimeSeries/TSMREVRANGE_Test.php +++ b/tests/Predis/Command/Redis/TimeSeries/TSMREVRANGE_Test.php @@ -118,16 +118,16 @@ class TSMREVRANGE_Test extends PredisCommandTestCase $redis = $this->getResp3Client(); $expectedResponse = [ 'type=stock' => [ - ['type' => 'stock'], - ['reducers' => ['max']], - ['sources' => ['stock:A', 'stock:B']], - [ - [1020, 120], - [1010, 110], - [1000, 120], - ], + ['type' => 'stock'], + ['reducers' => ['max']], + ['sources' => ['stock:A', 'stock:B']], + [ + [1020, 120], + [1010, 110], + [1000, 120], ], - ]; + ], + ]; $this->assertEquals( 'OK', diff --git a/tests/Predis/Connection/CompositeStreamConnectionTest.php b/tests/Predis/Connection/CompositeStreamConnectionTest.php index 8bb8461a..fa54eb72 100644 --- a/tests/Predis/Connection/CompositeStreamConnectionTest.php +++ b/tests/Predis/Connection/CompositeStreamConnectionTest.php @@ -86,16 +86,16 @@ class CompositeStreamConnectionTest extends PredisConnectionTestCase public function testPersistentParameterWithFalseLikeValues(): void { $connection1 = $this->createConnectionWithParams(['persistent' => 0]); - $this->assertNonPersistentConnection($connection1); + $this->assertNonPersistentConnection($connection1->getResource()->detach()); $connection2 = $this->createConnectionWithParams(['persistent' => false]); - $this->assertNonPersistentConnection($connection2); + $this->assertNonPersistentConnection($connection2->getResource()->detach()); $connection3 = $this->createConnectionWithParams(['persistent' => '0']); - $this->assertNonPersistentConnection($connection3); + $this->assertNonPersistentConnection($connection3->getResource()->detach()); $connection4 = $this->createConnectionWithParams(['persistent' => 'false']); - $this->assertNonPersistentConnection($connection4); + $this->assertNonPersistentConnection($connection4->getResource()->detach()); } /** @@ -105,16 +105,16 @@ class CompositeStreamConnectionTest extends PredisConnectionTestCase public function testPersistentParameterWithTrueLikeValues(): void { $connection1 = $this->createConnectionWithParams(['persistent' => 1]); - $this->assertPersistentConnection($connection1); + $this->assertPersistentConnection($connection1->getResource()->detach()); $connection2 = $this->createConnectionWithParams(['persistent' => true]); - $this->assertPersistentConnection($connection2); + $this->assertPersistentConnection($connection2->getResource()->detach()); $connection3 = $this->createConnectionWithParams(['persistent' => '1']); - $this->assertPersistentConnection($connection3); + $this->assertPersistentConnection($connection3->getResource()->detach()); $connection4 = $this->createConnectionWithParams(['persistent' => 'true']); - $this->assertPersistentConnection($connection4); + $this->assertPersistentConnection($connection4->getResource()->detach()); $connection1->disconnect(); } @@ -128,10 +128,10 @@ class CompositeStreamConnectionTest extends PredisConnectionTestCase $connection1 = $this->createConnectionWithParams(['persistent' => true]); $connection2 = $this->createConnectionWithParams(['persistent' => true]); - $this->assertPersistentConnection($connection1); - $this->assertPersistentConnection($connection2); + $this->assertPersistentConnection($connection1->getResource()->detach()); + $this->assertPersistentConnection($connection2->getResource()->detach()); - $this->assertSame($connection1->getResource(), $connection2->getResource()); + $this->assertEquals($connection1->getResource(), $connection2->getResource()); $connection1->disconnect(); } @@ -145,8 +145,8 @@ class CompositeStreamConnectionTest extends PredisConnectionTestCase $connection1 = $this->createConnectionWithParams(['persistent' => 'conn1']); $connection2 = $this->createConnectionWithParams(['persistent' => 'conn2']); - $this->assertPersistentConnection($connection1); - $this->assertPersistentConnection($connection2); + $this->assertPersistentConnection($connection1->getResource()->detach()); + $this->assertPersistentConnection($connection2->getResource()->detach()); $this->assertNotSame($connection1->getResource(), $connection2->getResource()); } diff --git a/tests/Predis/Connection/Resource/StreamTest.php b/tests/Predis/Connection/Resource/StreamTest.php new file mode 100644 index 00000000..ebb896d2 --- /dev/null +++ b/tests/Predis/Connection/Resource/StreamTest.php @@ -0,0 +1,439 @@ +expectException(InvalidArgumentException::class); + $this->expectExceptionMessage('Given stream is not a valid resource'); + + new Stream(false); + } + + /** + * @return void + */ + public function testToStringReturnsAllRemainingContent(): void + { + $handle = fopen('php://temp', 'rb+'); + fwrite($handle, 'data'); + $stream = new Stream($handle); + + $this->assertSame('data', (string) $stream); + } + + /** + * @return void + */ + public function testClosesStream(): void + { + $handle = fopen('php://temp', 'rb+'); + fwrite($handle, 'data'); + $stream = new Stream($handle); + + $stream->close(); + $this->assertTrue(true); + } + + /** + * @return void + */ + public function testDetachReturnsStreamAndDetachItFromObject(): void + { + $handle = fopen('php://temp', 'rb+'); + fwrite($handle, 'data'); + $stream = new Stream($handle); + $detachedStream = $stream->detach(); + fseek($detachedStream, 0); + + $this->assertSame('data', stream_get_contents($detachedStream)); + $this->assertNull($stream->detach()); + $this->assertFalse($stream->isReadable()); + $this->assertFalse($stream->isWritable()); + $this->assertFalse($stream->isSeekable()); + } + + /** + * @return void + */ + public function testGetSize(): void + { + $handle = fopen('php://temp', 'rb+'); + fwrite($handle, 'data'); + $stream = new Stream($handle); + + $this->assertSame(4, $stream->getSize()); + $stream->detach(); + $this->assertNull($stream->getSize()); + } + + /** + * @return void + */ + public function testTellThrowsExceptionOnDetachedStream(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->detach(); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Stream is detached'); + + $stream->tell(); + } + + /** + * @return void + */ + public function testTellReturnsCurrentPosition(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + + $this->assertSame(0, $stream->tell()); + } + + /** + * @return void + */ + public function testEofChecksIfPointerAtTheEndOfTheStream(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->read(1); + + $this->assertTrue($stream->eof()); + } + + /** + * @return void + */ + public function testEofThrowsExceptionOnDetachedStream(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->detach(); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Stream is detached'); + + $stream->eof(); + } + + /** + * @return void + */ + public function testIsSeekable(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + + $this->assertTrue($stream->isSeekable()); + } + + /** + * @return void + */ + public function testSeekThrowsExceptionOnDetachedStream(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->detach(); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Stream is detached'); + + $stream->seek(0, 1); + } + + /** + * @return void + */ + public function testSeekThrowsExceptionOnIncorrectOffset(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Unable to seek stream from offset 10 to whence 1'); + + $stream->seek(10, 1); + } + + /** + * @return void + */ + public function testRewind(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + + $stream->rewind(); + $this->assertTrue(true); + } + + /** + * @dataProvider writableModeProvider + * @param string $mode + * @return void + */ + public function testIsWritable(string $mode): void + { + $handle = fopen('php://temp', $mode); + $stream = new Stream($handle); + + $this->assertTrue($stream->isWritable()); + } + + /** + * @return void + */ + public function testWrite(): void + { + $handle = fopen('php://temp', 'wb+'); + $stream = new Stream($handle); + + $this->assertSame(4, $stream->write('data')); + } + + /** + * @return void + */ + public function testWriteThrowsExceptionOnDetachedStream(): void + { + $handle = fopen('php://temp', 'wb+'); + $stream = new Stream($handle); + $stream->detach(); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Stream is detached'); + + $stream->write('data'); + } + + /** + * @return void + */ + public function testWriteThrowsExceptionOnReadOnlyStream(): void + { + $handle = fopen('php://temp', 'rb'); + $stream = new Stream($handle); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Cannot write to a non-writable stream'); + + $stream->write('data'); + } + + /** + * @dataProvider readableModeProvider + * @param string $mode + * @return void + */ + public function testIsReadable(string $mode): void + { + $handle = fopen('php://temp', $mode); + $stream = new Stream($handle); + + $this->assertTrue($stream->isReadable()); + } + + /** + * @return void + */ + public function testRead(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->write('data'); + $stream->rewind(); + + $this->assertSame('data', $stream->read(4)); + } + + /** + * @return void + */ + public function testReadReturnsEmptyStringOnZeroLength(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->write('data'); + $stream->rewind(); + + $this->assertSame('', $stream->read(0)); + } + + /** + * @return void + */ + public function testReadThrowsExceptionOnDetachedStream(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->detach(); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Stream is detached'); + + $stream->read(4); + } + + /** + * @return void + */ + public function testReadThrowsExceptionOnWriteOnlyStream(): void + { + $handle = fopen('php://output', 'wb'); + $stream = new Stream($handle); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Cannot read from non-readable stream'); + + $stream->read(4); + } + + /** + * @return void + */ + public function testReadThrowsExceptionOnNegativeLength(): void + { + $handle = fopen('php://temp', 'wb'); + $stream = new Stream($handle); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Length parameter cannot be negative'); + + $stream->read(-2); + } + + /** + * @return void + */ + public function testGetContents(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->write('data'); + $stream->rewind(); + + $this->assertSame('data', $stream->getContents()); + } + + /** + * @return void + */ + public function testGetContentsThrowsExceptionOnWriteOnlyStream(): void + { + $handle = fopen('php://output', 'wb'); + $stream = new Stream($handle); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Cannot read from non-readable stream'); + + $stream->getContents(); + } + + /** + * @return void + */ + public function testGetContentsThrowsExceptionOnDetachedStream(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $stream->detach(); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Stream is detached'); + + $stream->getContents(); + } + + /** + * @return void + */ + public function testGetMetadata(): void + { + $handle = fopen('php://temp', 'rb+'); + $stream = new Stream($handle); + $metadata = $stream->getMetadata(); + + $this->assertArrayHasKey('wrapper_type', $metadata); + $this->assertArrayHasKey('stream_type', $metadata); + $this->assertArrayHasKey('mode', $metadata); + $this->assertArrayHasKey('unread_bytes', $metadata); + $this->assertArrayHasKey('seekable', $metadata); + $this->assertArrayHasKey('uri', $metadata); + $this->assertSame('php://temp', $stream->getMetadata('uri')); + + $stream->detach(); + + $this->assertNull($stream->getMetadata()); + } + + public function writableModeProvider(): array + { + return [ + ['w'], + ['w+'], + ['rw'], + ['r+'], + ['x+'], + ['c+'], + ['wb'], + ['w+b'], + ['r+b'], + ['rb+'], + ['x+b'], + ['c+b'], + ['w+t'], + ['r+t'], + ['x+t'], + ['c+t'], + ['a'], + ['a+'], + ]; + } + + public function readableModeProvider(): iterable + { + return [ + ['r'], + ['w+'], + ['r+'], + ['x+'], + ['c+'], + ['rb'], + ['w+b'], + ['r+b'], + ['x+b'], + ['c+b'], + ['rt'], + ['w+t'], + ['r+t'], + ['x+t'], + ['c+t'], + ['a+'], + ['rb+'], + ]; + } +} diff --git a/tests/Predis/Connection/StreamConnectionTest.php b/tests/Predis/Connection/StreamConnectionTest.php index 88bc3b27..464e6a8e 100644 --- a/tests/Predis/Connection/StreamConnectionTest.php +++ b/tests/Predis/Connection/StreamConnectionTest.php @@ -15,10 +15,38 @@ namespace Predis\Connection; use PHPUnit\Framework\MockObject\MockObject; use Predis\Client; use Predis\Command\RawCommand; +use Predis\Connection\Resource\Exception\StreamInitException; +use Predis\Connection\Resource\StreamFactoryInterface; +use Predis\Consumer\Push\PushResponse; +use Predis\Protocol\ProtocolException; use Predis\Response\Error as ErrorResponse; +use Psr\Http\Message\StreamInterface; +use RuntimeException; +/** + * @method StreamConnection createConnection(bool $initialize = false) + * @method StreamConnection createConnectionWithParams($parameters, $initialize = false) + */ class StreamConnectionTest extends PredisConnectionTestCase { + /** + * @var StreamFactoryInterface + */ + private $mockStreamFactory; + + /** + * @var StreamInterface + */ + private $mockStream; + + protected function setUp(): void + { + parent::setUp(); + + $this->mockStreamFactory = $this->getMockBuilder(StreamFactoryInterface::class)->getMock(); + $this->mockStream = $this->getMockBuilder(StreamInterface::class)->getMock(); + } + /** * {@inheritDoc} */ @@ -40,7 +68,7 @@ class StreamConnectionTest extends PredisConnectionTestCase /** @var NodeConnectionInterface|MockObject */ $connection = $this ->getMockBuilder($this->getConnectionClass()) - ->onlyMethods(['write', 'read', 'createResource']) + ->onlyMethods(['write', 'read']) ->setConstructorArgs([new Parameters()]) ->getMock(); $connection @@ -52,8 +80,6 @@ class StreamConnectionTest extends PredisConnectionTestCase new ErrorResponse('ERR invalid DB index') ); - $connection->method('createResource'); - $connection->addConnectCommand($cmdSelect); $connection->connect(); } @@ -61,24 +87,538 @@ class StreamConnectionTest extends PredisConnectionTestCase /** * @group disconnected */ - public function testDoesntThrowErrorOnInvalidResource(): void + public function testConnectWithNoConnectCommands(): void { - $this->expectException('Predis\Connection\ConnectionException'); + $parameters = new Parameters(); - $cmdSelect = RawCommand::create('SELECT', '1000'); - $invalidResource = null; + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); - /** @var NodeConnectionInterface|MockObject */ - $connection = $this - ->getMockBuilder($this->getConnectionClass()) - ->onlyMethods(['getResource']) - ->setConstructorArgs([new Parameters()]) - ->getMock(); - $connection - ->method('getResource') - ->willReturn($invalidResource); + $this->mockStream + ->expects($this->never()) + ->method('write') + ->withAnyParameters(); - $connection->writeRequest($cmdSelect); + $this->mockStream + ->expects($this->never()) + ->method('read') + ->withAnyParameters(); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $connection->connect(); + } + + /** + * @group disconnected + */ + public function testConnectWithConnectCommands(): void + { + $parameters = new Parameters(); + $command1 = new RawCommand('AUTH', [12345, 12345]); + $command2 = new RawCommand('SELECT', [10]); + $command3 = new RawCommand('CLIENT', ['SETNAME', 'predis']); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->exactly(3)) + ->method('write') + ->withConsecutive( + [$command1->serializeCommand()], + [$command2->serializeCommand()], + [$command3->serializeCommand()] + ) + ->willReturnOnConsecutiveCalls( + strlen($command1->serializeCommand()), + strlen($command2->serializeCommand()), + strlen($command3->serializeCommand()) + ); + + $this->mockStream + ->expects($this->exactly(3)) + ->method('read') + ->with(-1) + ->willReturn('+OK\r\n'); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $connection->addConnectCommand($command1); + $connection->addConnectCommand($command2); + $connection->addConnectCommand($command3); + + $connection->connect(); + } + + /** + * @group disconnected + */ + public function testDisconnectDoNothingOnAlreadyDisconnectedConnection(): void + { + $parameters = new Parameters(); + + $this->mockStreamFactory + ->expects($this->never()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->never()) + ->method('close') + ->withAnyParameters(); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $connection->disconnect(); + } + + /** + * @group disconnected + */ + public function testDisconnectOnAlreadyConnectedConnection(): void + { + $parameters = new Parameters(); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('close') + ->withAnyParameters(); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $connection->connect(); + $connection->disconnect(); + } + + /** + * @group disconnected + */ + public function testWriteWholeBufferAtOnce(): void + { + $parameters = new Parameters(); + $command = new RawCommand('PING'); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('write') + ->with($command->serializeCommand()) + ->willReturn(strlen($command->serializeCommand())); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + $connection->write($command->serializeCommand()); + } + + /** + * @group disconnected + */ + public function testWriteBufferByChunks(): void + { + $parameters = new Parameters(); + $command = new RawCommand('PING'); + $firstChunk = substr($command->serializeCommand(), 0, strlen($command->serializeCommand()) / 2); + $secondChunk = substr($command->serializeCommand(), strlen($command->serializeCommand()) / 2); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->exactly(2)) + ->method('write') + ->withConsecutive([$command->serializeCommand()], [$secondChunk]) + ->willReturnOnConsecutiveCalls(strlen($firstChunk), strlen($secondChunk)); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + $connection->write($command->serializeCommand()); + } + + /** + * @group disconnected + */ + public function testWriteThrowsExceptionOnNonSuccessfulStreamWrite(): void + { + $parameters = new Parameters(); + $command = new RawCommand('PING'); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('write') + ->with($command->serializeCommand()) + ->willThrowException(new RuntimeException('Error from stream during write', 1)); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->expectException(ConnectionException::class); + $this->expectExceptionMessage('Error while writing bytes to the server. [tcp://127.0.0.1:6379]'); + + $connection->write($command->serializeCommand()); + } + + /** + * @dataProvider simpleDataTypesProvider + * @group disconnected + */ + public function testReadSimpleDataTypes(string $payload, $expectedResponse): void + { + $parameters = new Parameters(['protocol' => 3]); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with($parameters) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('read') + ->with(-1) + ->willReturn($payload); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->assertEquals($expectedResponse, $connection->read()); + } + + /** + * @dataProvider aggregateDataTypesProvider + * @group disconnected + */ + public function testReadAggregateDataTypes(int $count, array $lengths, array $processedLines, $expectedResponse): void + { + $parameters = new Parameters(['protocol' => 3]); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with($parameters) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->exactly($count)) + ->method('read') + ->withConsecutive(...$lengths) + ->willReturnOnConsecutiveCalls(...$processedLines); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->assertEquals($expectedResponse, $connection->read()); + } + + /** + * @group disconnected + */ + public function testReadAggregateDataTypesThrowsExceptionOnBrokenChunk(): void + { + $parameters = new Parameters(['protocol' => 3]); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with($parameters) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->exactly(2)) + ->method('read') + ->withConsecutive([-1], [8]) + ->willReturnOnConsecutiveCalls( + "$6\r\nfoobar\r\n", + $this->throwException(new RuntimeException('Error while reading', 1)) + ); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->expectException(ConnectionException::class); + $this->expectExceptionMessage('Error while reading bytes from the server. [tcp://127.0.0.1:6379]'); + + $connection->read(); + } + + /** + * @group disconnected + */ + public function testReadThrowsExceptionOnUnknownDataType(): void + { + $parameters = new Parameters(['protocol' => 3]); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with($parameters) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('read') + ->with(-1) + ->willReturn("@wrongtype\r\n"); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->expectException(ProtocolException::class); + $this->expectExceptionMessage("Unknown response prefix: '@'. [tcp://127.0.0.1:6379]"); + + $connection->read(); + } + + /** + * @group disconnected + */ + public function testReadThrowsExceptionOnInvalidBrokenStreamTransport(): void + { + $parameters = new Parameters(['protocol' => 3]); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with($parameters) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('read') + ->with(-1) + ->willThrowException(new RuntimeException('Error on read', 1)); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->expectException(ConnectionException::class); + $this->expectExceptionMessage('Error while reading line from the server. [tcp://127.0.0.1:6379]'); + + $connection->read(); + } + + /** + * @group disconnected + */ + public function testReadThrowsExceptionOnNonReadableStream(): void + { + $parameters = new Parameters(['protocol' => 3]); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with($parameters) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('read') + ->with(-1) + ->willThrowException(new RuntimeException('Cannot read from non-readable stream')); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Cannot read from non-readable stream'); + + $connection->read(); + } + + /** + * @group disconnected + */ + public function testReadThrowsExceptionOnEOF(): void + { + $parameters = new Parameters(['protocol' => 3]); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with($parameters) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('eof') + ->willReturn(true); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $this->expectException(RuntimeException::class); + $this->expectExceptionMessage('Stream is already at the end'); + + $connection->read(); + } + + /** + * @group disconnected + */ + public function testWriteRequest(): void + { + $parameters = new Parameters(); + $command = new RawCommand('PING'); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('write') + ->with($command->serializeCommand()) + ->willReturn(strlen($command->serializeCommand())); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + $connection->writeRequest($command); + } + + /** + * @group disconnected + */ + public function testHasDataToRead(): void + { + $parameters = new Parameters(); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->with(new Parameters()) + ->willReturn($this->mockStream); + + $this->mockStream + ->expects($this->once()) + ->method('eof') + ->willReturn(false); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + $this->assertTrue($connection->hasDataToRead()); + } + + public function simpleDataTypesProvider(): array + { + return [ + 'simple_string' => [ + "+OK\r\n", + 'OK', + ], + 'error' => [ + "-Error message\r\n", + new ErrorResponse('Error message'), + ], + 'integer' => [ + ":1000\r\n", + 1000, + ], + 'null' => [ + "_\r\n", + null, + ], + 'double' => [ + ",1.23\r\n", + 1.23, + ], + 'boolean' => [ + "#f\r\n", + false, + ], + 'big_number' => [ + "(3492890328409238509324850943850943825024385\r\n", + 3492890328409238509324850943850943825024385, + ], + ]; + } + + public function aggregateDataTypesProvider(): array + { + return [ + 'bulk_string' => [ + 2, + [[-1], [8]], + ["$6\r\nfoobar\r\n", "foobar\r\n"], + 'foobar', + ], + 'array' => [ + 5, + [[-1], [-1], [5], [-1], [5]], + [ + "*2\r\n$3\r\nfoo\r\n$3\r\nbar\r\n", + "$3\r\nfoo\r\n", + "foo\r\n", + "$3\r\nbar\r\n", + "bar\r\n", + ], + ['foo', 'bar'], + ], + 'verbatim_string' => [ + 2, + [[-1], [17]], + ["=15\r\ntxt:Some string\r\n", "txt:Some string\r\n"], + 'Some string', + ], + 'blob_error' => [ + 2, + [[-1], [23]], + ["!21\r\nSYNTAX invalid syntax\r\n", "SYNTAX invalid syntax\r\n"], + new ErrorResponse('SYNTAX invalid syntax'), + ], + 'map' => [ + 5, + [[-1], [-1], [-1], [-1], [-1]], + [ + "%2\r\n+first\r\n:1\r\n+second\r\n:2\r\n", + "+first\r\n", + ":1\r\n", + "+second\r\n", + ":2\r\n", + ], + ['first' => 1, 'second' => 2], + ], + 'set' => [ + 6, + [[-1], [-1], [-1], [-1], [-1], [-1]], + [ + "~5\r\n+orange\r\n+apple\r\n#t\r\n:100\r\n:999\r\n", + "+orange\r\n", + "+apple\r\n", + "#t\r\n", + ":100\r\n", + ":999\r\n", + ], + ['orange', 'apple', true, 100, 999], + ], + 'push' => [ + 4, + [[-1], [-1], [-1], [-1]], + [ + ">3\r\n+message\r\n+somechannel\r\n+this is the message\r\n", + "+message\r\n", + "+somechannel\r\n", + "+this is the message\r\n", + ], + new PushResponse(['message', 'somechannel', 'this is the message']), + ], + ]; } // ******************************************************************** // @@ -92,16 +632,16 @@ class StreamConnectionTest extends PredisConnectionTestCase public function testPersistentParameterWithFalseLikeValues(): void { $connection1 = $this->createConnectionWithParams(['persistent' => 0]); - $this->assertNonPersistentConnection($connection1); + $this->assertNonPersistentConnection($connection1->getResource()->detach()); $connection2 = $this->createConnectionWithParams(['persistent' => false]); - $this->assertNonPersistentConnection($connection2); + $this->assertNonPersistentConnection($connection2->getResource()->detach()); $connection3 = $this->createConnectionWithParams(['persistent' => '0']); - $this->assertNonPersistentConnection($connection3); + $this->assertNonPersistentConnection($connection3->getResource()->detach()); $connection4 = $this->createConnectionWithParams(['persistent' => 'false']); - $this->assertNonPersistentConnection($connection4); + $this->assertNonPersistentConnection($connection4->getResource()->detach()); } /** @@ -111,16 +651,16 @@ class StreamConnectionTest extends PredisConnectionTestCase public function testPersistentParameterWithTrueLikeValues(): void { $connection1 = $this->createConnectionWithParams(['persistent' => 1]); - $this->assertPersistentConnection($connection1); + $this->assertPersistentConnection($connection1->getResource()->detach()); $connection2 = $this->createConnectionWithParams(['persistent' => true]); - $this->assertPersistentConnection($connection2); + $this->assertPersistentConnection($connection2->getResource()->detach()); $connection3 = $this->createConnectionWithParams(['persistent' => '1']); - $this->assertPersistentConnection($connection3); + $this->assertPersistentConnection($connection3->getResource()->detach()); $connection4 = $this->createConnectionWithParams(['persistent' => 'true']); - $this->assertPersistentConnection($connection4); + $this->assertPersistentConnection($connection4->getResource()->detach()); $connection1->disconnect(); } @@ -134,10 +674,10 @@ class StreamConnectionTest extends PredisConnectionTestCase $connection1 = $this->createConnectionWithParams(['persistent' => true]); $connection2 = $this->createConnectionWithParams(['persistent' => true]); - $this->assertPersistentConnection($connection1); - $this->assertPersistentConnection($connection2); + $this->assertPersistentConnection($connection1->getResource()->detach()); + $this->assertPersistentConnection($connection2->getResource()->detach()); - $this->assertSame($connection1->getResource(), $connection2->getResource()); + $this->assertEquals($connection1->getResource(), $connection2->getResource()); $connection1->disconnect(); } @@ -151,19 +691,90 @@ class StreamConnectionTest extends PredisConnectionTestCase $connection1 = $this->createConnectionWithParams(['persistent' => 'conn1']); $connection2 = $this->createConnectionWithParams(['persistent' => 'conn2']); - $this->assertPersistentConnection($connection1); - $this->assertPersistentConnection($connection2); + $this->assertPersistentConnection($connection1->getResource()->detach()); + $this->assertPersistentConnection($connection2->getResource()->detach()); $this->assertNotSame($connection1->getResource(), $connection2->getResource()); } + /** + * @group connected + * @group slow + */ + public function testThrowsExceptionOnConnectionTimeout(): void + { + // FACTORY TEST + $this->expectException(StreamInitException::class); + $this->expectExceptionMessageMatches('/.* \[tcp:\/\/169.254.10.10:6379\]/'); + + $connection = $this->createConnectionWithParams([ + 'host' => '169.254.10.10', + 'timeout' => 0.1, + ], false); + + $connection->connect(); + } + + /** + * @group connected + * @group slow + */ + public function testThrowsExceptionOnConnectionTimeoutIPv6(): void + { + // FACTORY TEST + $this->expectException(StreamInitException::class); + $this->expectExceptionMessageMatches('/.* \[tcp:\/\/\[0:0:0:0:0:ffff:a9fe:a0a\]:6379\]/'); + + $connection = $this->createConnectionWithParams([ + 'host' => '0:0:0:0:0:ffff:a9fe:a0a', + 'timeout' => 0.1, + ], false); + + $connection->connect(); + } + + /** + * @group connected + * @group slow + */ + public function testThrowsExceptionOnUnixDomainSocketNotFound(): void + { + // FACTORY TEST + $this->expectException(StreamInitException::class); + $this->expectExceptionMessageMatches('/.* \[unix:\/tmp\/nonexistent\/redis\.sock]/'); + + $connection = $this->createConnectionWithParams([ + 'scheme' => 'unix', + 'path' => '/tmp/nonexistent/redis.sock', + ], false); + + $connection->connect(); + } + + /** + * @medium + * @group connected + */ + public function testThrowsExceptionOnProtocolDesynchronizationErrors(): void + { + $this->expectException('Predis\Protocol\ProtocolException'); + + $connection = $this->createConnection(); + $stream = $connection->getResource(); + + $connection->writeRequest($this->getCommandFactory()->create('ping')); + $stream->read(1); + + $connection->read(); + } + /** * @group connected */ public function testTcpNodelayParameterSetsContextFlagWhenTrue() { $connection = $this->createConnectionWithParams(['tcp_nodelay' => true]); - $options = stream_context_get_options($connection->getResource()); + $options = stream_context_get_options($connection->getResource()->detach()); $this->assertIsArray($options); $this->assertArrayHasKey('socket', $options); @@ -177,7 +788,7 @@ class StreamConnectionTest extends PredisConnectionTestCase public function testTcpNodelayParameterDoesNotSetContextFlagWhenFalse() { $connection = $this->createConnectionWithParams(['tcp_nodelay' => false]); - $options = stream_context_get_options($connection->getResource()); + $options = stream_context_get_options($connection->getResource()->detach()); $this->assertIsArray($options); $this->assertArrayHasKey('socket', $options); @@ -191,7 +802,7 @@ class StreamConnectionTest extends PredisConnectionTestCase public function testTcpDelayContextFlagIsNotSetByDefault() { $connection = $this->createConnectionWithParams([]); - $options = stream_context_get_options($connection->getResource()); + $options = stream_context_get_options($connection->getResource()->detach()); $this->assertIsArray($options); $this->assertArrayHasKey('socket', $options);