mirror of
https://github.com/predis/predis.git
synced 2026-10-10 22:07:31 +00:00
Remove Webdis and phpiredis (#1291)
* remove webdis and phpiredis * update changelog * add pr number
This commit is contained in:
@@ -1,42 +0,0 @@
|
||||
<?php
|
||||
|
||||
/*
|
||||
* This file is part of the Predis package.
|
||||
*
|
||||
* (c) 2009-2020 Daniele Alessandri
|
||||
* (c) 2021-2023 Till Krüss
|
||||
*
|
||||
* For the full copyright and license information, please view the LICENSE
|
||||
* file that was distributed with this source code.
|
||||
*/
|
||||
|
||||
namespace Predis\Cluster\Hash;
|
||||
|
||||
use Predis\NotSupportedException;
|
||||
|
||||
/**
|
||||
* Hash generator implementing the CRC-CCITT-16 algorithm used by redis-cluster.
|
||||
*
|
||||
* @deprecated 2.1.2
|
||||
*/
|
||||
class PhpiredisCRC16 implements HashGeneratorInterface
|
||||
{
|
||||
public function __construct()
|
||||
{
|
||||
if (!function_exists('phpiredis_utils_crc16')) {
|
||||
// @codeCoverageIgnoreStart
|
||||
throw new NotSupportedException(
|
||||
'This hash generator requires a compatible version of ext-phpiredis'
|
||||
);
|
||||
// @codeCoverageIgnoreEnd
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function hash($value)
|
||||
{
|
||||
return phpiredis_utils_crc16($value);
|
||||
}
|
||||
}
|
||||
@@ -26,7 +26,7 @@ class CRC16 implements OptionInterface
|
||||
* Returns an hash generator instance from a descriptive name.
|
||||
*
|
||||
* @param OptionsInterface $options Client options.
|
||||
* @param string $description Identifier of a hash generator (`predis`, `phpiredis`)
|
||||
* @param string $description Identifier of a hash generator (`predis`)
|
||||
*
|
||||
* @return callable
|
||||
*/
|
||||
@@ -34,11 +34,9 @@ class CRC16 implements OptionInterface
|
||||
{
|
||||
if ($description === 'predis') {
|
||||
return new Hash\CRC16();
|
||||
} elseif ($description === 'phpiredis') {
|
||||
return new Hash\PhpiredisCRC16();
|
||||
} else {
|
||||
throw new InvalidArgumentException(
|
||||
'String value for the crc16 option must be either `predis` or `phpiredis`'
|
||||
'String value for the crc16 option must be either `predis`'
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -67,8 +65,6 @@ class CRC16 implements OptionInterface
|
||||
*/
|
||||
public function getDefault(OptionsInterface $options)
|
||||
{
|
||||
return function_exists('phpiredis_utils_crc16')
|
||||
? new Hash\PhpiredisCRC16()
|
||||
: new Hash\CRC16();
|
||||
return new Hash\CRC16();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,8 +17,6 @@ use Predis\Configuration\OptionInterface;
|
||||
use Predis\Configuration\OptionsInterface;
|
||||
use Predis\Connection\Factory;
|
||||
use Predis\Connection\FactoryInterface;
|
||||
use Predis\Connection\PhpiredisSocketConnection;
|
||||
use Predis\Connection\PhpiredisStreamConnection;
|
||||
use Predis\Connection\RelayConnection;
|
||||
|
||||
/**
|
||||
@@ -87,9 +85,6 @@ class Connections implements OptionInterface
|
||||
* string that identifies specific configurations of schemes and connection
|
||||
* classes. Supported configuration values are:
|
||||
*
|
||||
* - "phpiredis-stream" maps tcp, redis, unix to PhpiredisStreamConnection
|
||||
* - "phpiredis-socket" maps tcp, redis, unix to PhpiredisSocketConnection
|
||||
* - "phpiredis" is an alias of "phpiredis-stream"
|
||||
* - "relay" maps tcp, redis, unix, tls, rediss to RelayConnection
|
||||
*
|
||||
* @param OptionsInterface $options Client options
|
||||
@@ -105,19 +100,6 @@ class Connections implements OptionInterface
|
||||
$factory = $this->getDefault($options);
|
||||
|
||||
switch (strtolower($value)) {
|
||||
case 'phpiredis':
|
||||
case 'phpiredis-stream':
|
||||
$factory->define('tcp', PhpiredisStreamConnection::class);
|
||||
$factory->define('redis', PhpiredisStreamConnection::class);
|
||||
$factory->define('unix', PhpiredisStreamConnection::class);
|
||||
break;
|
||||
|
||||
case 'phpiredis-socket':
|
||||
$factory->define('tcp', PhpiredisSocketConnection::class);
|
||||
$factory->define('redis', PhpiredisSocketConnection::class);
|
||||
$factory->define('unix', PhpiredisSocketConnection::class);
|
||||
break;
|
||||
|
||||
case 'relay':
|
||||
$factory->define('tcp', RelayConnection::class);
|
||||
$factory->define('redis', RelayConnection::class);
|
||||
|
||||
@@ -30,7 +30,6 @@ class Factory implements FactoryInterface
|
||||
'tls' => 'Predis\Connection\StreamConnection',
|
||||
'redis' => 'Predis\Connection\StreamConnection',
|
||||
'rediss' => 'Predis\Connection\StreamConnection',
|
||||
'http' => 'Predis\Connection\WebdisConnection',
|
||||
];
|
||||
|
||||
/**
|
||||
|
||||
@@ -1,420 +0,0 @@
|
||||
<?php
|
||||
|
||||
/*
|
||||
* This file is part of the Predis package.
|
||||
*
|
||||
* (c) 2009-2020 Daniele Alessandri
|
||||
* (c) 2021-2023 Till Krüss
|
||||
*
|
||||
* For the full copyright and license information, please view the LICENSE
|
||||
* file that was distributed with this source code.
|
||||
*/
|
||||
|
||||
namespace Predis\Connection;
|
||||
|
||||
use Closure;
|
||||
use InvalidArgumentException;
|
||||
use Predis\Command\CommandInterface;
|
||||
use Predis\NotSupportedException;
|
||||
use Predis\Response\Error as ErrorResponse;
|
||||
use Predis\Response\ErrorInterface as ErrorResponseInterface;
|
||||
use Predis\Response\Status as StatusResponse;
|
||||
|
||||
/**
|
||||
* This class provides the implementation of a Predis connection that uses the
|
||||
* PHP socket extension for network communication and wraps the phpiredis C
|
||||
* extension (PHP bindings for hiredis) to parse the Redis protocol.
|
||||
*
|
||||
* This class is intended to provide an optional low-overhead alternative for
|
||||
* processing responses from Redis compared to the standard pure-PHP classes.
|
||||
* Differences in speed when dealing with short inline responses are practically
|
||||
* nonexistent, the actual speed boost is for big multibulk responses when this
|
||||
* protocol processor can parse and return responses very fast.
|
||||
*
|
||||
* For instructions on how to build and install the phpiredis extension, please
|
||||
* consult the repository of the project.
|
||||
*
|
||||
* The connection parameters supported by this class are:
|
||||
*
|
||||
* - scheme: it can be either 'redis', 'tcp' or 'unix'.
|
||||
* - host: hostname or IP address of the server.
|
||||
* - port: TCP port of the server.
|
||||
* - path: path of a UNIX domain socket when scheme is 'unix'.
|
||||
* - timeout: timeout to perform the connection (default is 5 seconds).
|
||||
* - read_write_timeout: timeout of read / write operations.
|
||||
*
|
||||
* @see http://github.com/nrk/phpiredis
|
||||
* @deprecated 2.1.2
|
||||
*/
|
||||
class PhpiredisSocketConnection extends AbstractConnection
|
||||
{
|
||||
private $reader;
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __construct(ParametersInterface $parameters)
|
||||
{
|
||||
$this->assertExtensions();
|
||||
|
||||
parent::__construct($parameters);
|
||||
|
||||
$this->reader = $this->createReader();
|
||||
}
|
||||
|
||||
/**
|
||||
* Disconnects from the server and destroys the underlying resource and the
|
||||
* protocol reader resource when PHP's garbage collector kicks in.
|
||||
*/
|
||||
public function __destruct()
|
||||
{
|
||||
parent::__destruct();
|
||||
|
||||
phpiredis_reader_destroy($this->reader);
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks if the socket and phpiredis extensions are loaded in PHP.
|
||||
*/
|
||||
protected function assertExtensions()
|
||||
{
|
||||
if (!extension_loaded('sockets')) {
|
||||
throw new NotSupportedException(
|
||||
'The "sockets" extension is required by this connection backend.'
|
||||
);
|
||||
}
|
||||
|
||||
if (!extension_loaded('phpiredis')) {
|
||||
throw new NotSupportedException(
|
||||
'The "phpiredis" extension is required by this connection backend.'
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
protected function assertParameters(ParametersInterface $parameters)
|
||||
{
|
||||
switch ($parameters->scheme) {
|
||||
case 'tcp':
|
||||
case 'redis':
|
||||
case 'unix':
|
||||
break;
|
||||
|
||||
default:
|
||||
throw new InvalidArgumentException("Invalid scheme: '$parameters->scheme'.");
|
||||
}
|
||||
|
||||
if (isset($parameters->persistent)) {
|
||||
throw new NotSupportedException(
|
||||
'Persistent connections are not supported by this connection backend.'
|
||||
);
|
||||
}
|
||||
|
||||
return $parameters;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new instance of the protocol reader resource.
|
||||
*
|
||||
* @return resource
|
||||
*/
|
||||
private function createReader()
|
||||
{
|
||||
$reader = phpiredis_reader_create();
|
||||
|
||||
phpiredis_reader_set_status_handler($reader, $this->getStatusHandler());
|
||||
phpiredis_reader_set_error_handler($reader, $this->getErrorHandler());
|
||||
|
||||
return $reader;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the underlying protocol reader resource.
|
||||
*
|
||||
* @return resource
|
||||
*/
|
||||
protected function getReader()
|
||||
{
|
||||
return $this->reader;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the handler used by the protocol reader for inline responses.
|
||||
*
|
||||
* @return Closure
|
||||
*/
|
||||
protected function getStatusHandler()
|
||||
{
|
||||
static $statusHandler;
|
||||
|
||||
if (!$statusHandler) {
|
||||
$statusHandler = function ($payload) {
|
||||
return StatusResponse::get($payload);
|
||||
};
|
||||
}
|
||||
|
||||
return $statusHandler;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the handler used by the protocol reader for error responses.
|
||||
*
|
||||
* @return Closure
|
||||
*/
|
||||
protected function getErrorHandler()
|
||||
{
|
||||
static $errorHandler;
|
||||
|
||||
if (!$errorHandler) {
|
||||
$errorHandler = function ($errorMessage) {
|
||||
return new ErrorResponse($errorMessage);
|
||||
};
|
||||
}
|
||||
|
||||
return $errorHandler;
|
||||
}
|
||||
|
||||
/**
|
||||
* Helper method used to throw exceptions on socket errors.
|
||||
*/
|
||||
private function emitSocketError()
|
||||
{
|
||||
$errno = socket_last_error();
|
||||
$errstr = socket_strerror($errno);
|
||||
|
||||
$this->disconnect();
|
||||
|
||||
$this->onConnectionError(trim($errstr), $errno);
|
||||
}
|
||||
|
||||
/**
|
||||
* Gets the address of an host from connection parameters.
|
||||
*
|
||||
* @param ParametersInterface $parameters Parameters used to initialize the connection.
|
||||
*
|
||||
* @return string
|
||||
*/
|
||||
protected static function getAddress(ParametersInterface $parameters)
|
||||
{
|
||||
if (filter_var($host = $parameters->host, FILTER_VALIDATE_IP)) {
|
||||
return $host;
|
||||
}
|
||||
|
||||
if ($host === $address = gethostbyname($host)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return $address;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
protected function createResource()
|
||||
{
|
||||
$parameters = $this->parameters;
|
||||
|
||||
if ($parameters->scheme === 'unix') {
|
||||
$address = $parameters->path;
|
||||
$domain = AF_UNIX;
|
||||
$protocol = 0;
|
||||
} else {
|
||||
if (false === $address = self::getAddress($parameters)) {
|
||||
$this->onConnectionError("Cannot resolve the address of '$parameters->host'.");
|
||||
}
|
||||
|
||||
$domain = filter_var($address, FILTER_VALIDATE_IP, FILTER_FLAG_IPV6) ? AF_INET6 : AF_INET;
|
||||
$protocol = SOL_TCP;
|
||||
}
|
||||
|
||||
if (false === $socket = @socket_create($domain, SOCK_STREAM, $protocol)) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
|
||||
$this->setSocketOptions($socket, $parameters);
|
||||
$this->connectWithTimeout($socket, $address, $parameters);
|
||||
|
||||
return $socket;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sets options on the socket resource from the connection parameters.
|
||||
*
|
||||
* @param resource $socket Socket resource.
|
||||
* @param ParametersInterface $parameters Parameters used to initialize the connection.
|
||||
*/
|
||||
private function setSocketOptions($socket, ParametersInterface $parameters)
|
||||
{
|
||||
if ($parameters->scheme !== 'unix') {
|
||||
if (!socket_set_option($socket, SOL_TCP, TCP_NODELAY, 1)) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
|
||||
if (!socket_set_option($socket, SOL_SOCKET, SO_REUSEADDR, 1)) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
}
|
||||
|
||||
if (isset($parameters->read_write_timeout)) {
|
||||
$rwtimeout = (float) $parameters->read_write_timeout;
|
||||
$timeoutSec = floor($rwtimeout);
|
||||
$timeoutUsec = ($rwtimeout - $timeoutSec) * 1000000;
|
||||
|
||||
$timeout = [
|
||||
'sec' => $timeoutSec,
|
||||
'usec' => $timeoutUsec,
|
||||
];
|
||||
|
||||
if (!socket_set_option($socket, SOL_SOCKET, SO_SNDTIMEO, $timeout)) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
|
||||
if (!socket_set_option($socket, SOL_SOCKET, SO_RCVTIMEO, $timeout)) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Opens the actual connection to the server with a timeout.
|
||||
*
|
||||
* @param resource $socket Socket resource.
|
||||
* @param string $address IP address (DNS-resolved from hostname)
|
||||
* @param ParametersInterface $parameters Parameters used to initialize the connection.
|
||||
*
|
||||
* @return void
|
||||
*/
|
||||
private function connectWithTimeout($socket, $address, ParametersInterface $parameters)
|
||||
{
|
||||
socket_set_nonblock($socket);
|
||||
|
||||
if (@socket_connect($socket, $address, (int) $parameters->port) === false) {
|
||||
$error = socket_last_error();
|
||||
|
||||
if ($error != SOCKET_EINPROGRESS && $error != SOCKET_EALREADY) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
}
|
||||
|
||||
socket_set_block($socket);
|
||||
|
||||
$null = null;
|
||||
$selectable = [$socket];
|
||||
|
||||
$timeout = (isset($parameters->timeout) ? (float) $parameters->timeout : 5.0);
|
||||
$timeoutSecs = floor($timeout);
|
||||
$timeoutUSecs = ($timeout - $timeoutSecs) * 1000000;
|
||||
|
||||
$selected = socket_select($selectable, $selectable, $null, $timeoutSecs, $timeoutUSecs);
|
||||
|
||||
if ($selected === 2) {
|
||||
$this->onConnectionError('Connection refused.', SOCKET_ECONNREFUSED);
|
||||
}
|
||||
|
||||
if ($selected === 0) {
|
||||
$this->onConnectionError('Connection timed out.', SOCKET_ETIMEDOUT);
|
||||
}
|
||||
|
||||
if ($selected === false) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function connect()
|
||||
{
|
||||
if (parent::connect() && $this->initCommands) {
|
||||
foreach ($this->initCommands as $command) {
|
||||
$response = $this->executeCommand($command);
|
||||
|
||||
if ($response instanceof ErrorResponseInterface) {
|
||||
$this->onConnectionError("`{$command->getId()}` failed: {$response->getMessage()}", 0);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function disconnect()
|
||||
{
|
||||
if ($this->isConnected()) {
|
||||
phpiredis_reader_reset($this->reader);
|
||||
socket_close($this->getResource());
|
||||
|
||||
parent::disconnect();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
protected function write($buffer)
|
||||
{
|
||||
$socket = $this->getResource();
|
||||
|
||||
while (($length = strlen($buffer)) > 0) {
|
||||
$written = socket_write($socket, $buffer, $length);
|
||||
|
||||
if ($length === $written) {
|
||||
return;
|
||||
}
|
||||
|
||||
if ($written === false) {
|
||||
$this->onConnectionError('Error while writing bytes to the server.');
|
||||
}
|
||||
|
||||
$buffer = substr($buffer, $written);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function read()
|
||||
{
|
||||
$socket = $this->getResource();
|
||||
$reader = $this->reader;
|
||||
|
||||
while (PHPIREDIS_READER_STATE_INCOMPLETE === $state = phpiredis_reader_get_state($reader)) {
|
||||
if (@socket_recv($socket, $buffer, 4096, 0) === false || $buffer === '' || $buffer === null) {
|
||||
$this->emitSocketError();
|
||||
}
|
||||
|
||||
phpiredis_reader_feed($reader, $buffer);
|
||||
}
|
||||
|
||||
if ($state === PHPIREDIS_READER_STATE_COMPLETE) {
|
||||
return phpiredis_reader_get_reply($reader);
|
||||
} else {
|
||||
$this->onProtocolError(phpiredis_reader_get_error($reader));
|
||||
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function writeRequest(CommandInterface $command)
|
||||
{
|
||||
$arguments = $command->getArguments();
|
||||
array_unshift($arguments, $command->getId());
|
||||
|
||||
$this->write(phpiredis_format_command($arguments));
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __wakeup()
|
||||
{
|
||||
$this->assertExtensions();
|
||||
$this->reader = $this->createReader();
|
||||
}
|
||||
}
|
||||
@@ -1,262 +0,0 @@
|
||||
<?php
|
||||
|
||||
/*
|
||||
* This file is part of the Predis package.
|
||||
*
|
||||
* (c) 2009-2020 Daniele Alessandri
|
||||
* (c) 2021-2023 Till Krüss
|
||||
*
|
||||
* For the full copyright and license information, please view the LICENSE
|
||||
* file that was distributed with this source code.
|
||||
*/
|
||||
|
||||
namespace Predis\Connection;
|
||||
|
||||
use Closure;
|
||||
use InvalidArgumentException;
|
||||
use Predis\Command\CommandInterface;
|
||||
use Predis\NotSupportedException;
|
||||
use Predis\Response\Error as ErrorResponse;
|
||||
use Predis\Response\Status as StatusResponse;
|
||||
|
||||
/**
|
||||
* This class provides the implementation of a Predis connection that uses PHP's
|
||||
* streams for network communication and wraps the phpiredis C extension (PHP
|
||||
* bindings for hiredis) to parse and serialize the Redis protocol.
|
||||
*
|
||||
* This class is intended to provide an optional low-overhead alternative for
|
||||
* processing responses from Redis compared to the standard pure-PHP classes.
|
||||
* Differences in speed when dealing with short inline responses are practically
|
||||
* nonexistent, the actual speed boost is for big multibulk responses when this
|
||||
* protocol processor can parse and return responses very fast.
|
||||
*
|
||||
* For instructions on how to build and install the phpiredis extension, please
|
||||
* consult the repository of the project.
|
||||
*
|
||||
* The connection parameters supported by this class are:
|
||||
*
|
||||
* - scheme: it can be either 'redis', 'tcp' or 'unix'.
|
||||
* - host: hostname or IP address of the server.
|
||||
* - port: TCP port of the server.
|
||||
* - path: path of a UNIX domain socket when scheme is 'unix'.
|
||||
* - timeout: timeout to perform the connection.
|
||||
* - read_write_timeout: timeout of read / write operations.
|
||||
* - async_connect: performs the connection asynchronously.
|
||||
* - tcp_nodelay: enables or disables Nagle's algorithm for coalescing.
|
||||
* - persistent: the connection is left intact after a GC collection.
|
||||
*
|
||||
* @see https://github.com/nrk/phpiredis
|
||||
* @deprecated 2.1.2
|
||||
*/
|
||||
class PhpiredisStreamConnection extends StreamConnection
|
||||
{
|
||||
private $reader;
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __construct(ParametersInterface $parameters)
|
||||
{
|
||||
$this->assertExtensions();
|
||||
|
||||
parent::__construct($parameters);
|
||||
|
||||
$this->reader = $this->createReader();
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __destruct()
|
||||
{
|
||||
parent::__destruct();
|
||||
|
||||
phpiredis_reader_destroy($this->reader);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function disconnect()
|
||||
{
|
||||
phpiredis_reader_reset($this->reader);
|
||||
|
||||
parent::disconnect();
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks if the phpiredis extension is loaded in PHP.
|
||||
*/
|
||||
private function assertExtensions()
|
||||
{
|
||||
if (!extension_loaded('phpiredis')) {
|
||||
throw new NotSupportedException(
|
||||
'The "phpiredis" extension is required by this connection backend.'
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
protected function assertParameters(ParametersInterface $parameters)
|
||||
{
|
||||
switch ($parameters->scheme) {
|
||||
case 'tcp':
|
||||
case 'redis':
|
||||
case 'unix':
|
||||
break;
|
||||
|
||||
case 'tls':
|
||||
case 'rediss':
|
||||
throw new InvalidArgumentException('SSL encryption is not supported by this connection backend.');
|
||||
default:
|
||||
throw new InvalidArgumentException("Invalid scheme: '$parameters->scheme'.");
|
||||
}
|
||||
|
||||
return $parameters;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
protected function createStreamSocket(ParametersInterface $parameters, $address, $flags)
|
||||
{
|
||||
$socket = null;
|
||||
$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) && function_exists('socket_import_stream')) {
|
||||
$rwtimeout = (float) $parameters->read_write_timeout;
|
||||
$rwtimeout = $rwtimeout > 0 ? $rwtimeout : -1;
|
||||
|
||||
$timeout = [
|
||||
'sec' => $timeoutSeconds = floor($rwtimeout),
|
||||
'usec' => ($rwtimeout - $timeoutSeconds) * 1000000,
|
||||
];
|
||||
|
||||
$socket = $socket ?: socket_import_stream($resource);
|
||||
@socket_set_option($socket, SOL_SOCKET, SO_SNDTIMEO, $timeout);
|
||||
@socket_set_option($socket, SOL_SOCKET, SO_RCVTIMEO, $timeout);
|
||||
}
|
||||
|
||||
if (isset($parameters->tcp_nodelay) && function_exists('socket_import_stream')) {
|
||||
$socket = $socket ?: socket_import_stream($resource);
|
||||
socket_set_option($socket, SOL_TCP, TCP_NODELAY, (int) $parameters->tcp_nodelay);
|
||||
}
|
||||
|
||||
return $resource;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new instance of the protocol reader resource.
|
||||
*
|
||||
* @return resource
|
||||
*/
|
||||
private function createReader()
|
||||
{
|
||||
$reader = phpiredis_reader_create();
|
||||
|
||||
phpiredis_reader_set_status_handler($reader, $this->getStatusHandler());
|
||||
phpiredis_reader_set_error_handler($reader, $this->getErrorHandler());
|
||||
|
||||
return $reader;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the underlying protocol reader resource.
|
||||
*
|
||||
* @return resource
|
||||
*/
|
||||
protected function getReader()
|
||||
{
|
||||
return $this->reader;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the handler used by the protocol reader for inline responses.
|
||||
*
|
||||
* @return Closure
|
||||
*/
|
||||
protected function getStatusHandler()
|
||||
{
|
||||
static $statusHandler;
|
||||
|
||||
if (!$statusHandler) {
|
||||
$statusHandler = function ($payload) {
|
||||
return StatusResponse::get($payload);
|
||||
};
|
||||
}
|
||||
|
||||
return $statusHandler;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the handler used by the protocol reader for error responses.
|
||||
*
|
||||
* @return Closure
|
||||
*/
|
||||
protected function getErrorHandler()
|
||||
{
|
||||
static $errorHandler;
|
||||
|
||||
if (!$errorHandler) {
|
||||
$errorHandler = function ($errorMessage) {
|
||||
return new ErrorResponse($errorMessage);
|
||||
};
|
||||
}
|
||||
|
||||
return $errorHandler;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function read()
|
||||
{
|
||||
$socket = $this->getResource();
|
||||
$reader = $this->reader;
|
||||
|
||||
while (PHPIREDIS_READER_STATE_INCOMPLETE === $state = phpiredis_reader_get_state($reader)) {
|
||||
$buffer = stream_socket_recvfrom($socket, 4096);
|
||||
|
||||
if ($buffer === false || $buffer === '') {
|
||||
$this->onConnectionError('Error while reading bytes from the server.');
|
||||
}
|
||||
|
||||
phpiredis_reader_feed($reader, $buffer);
|
||||
}
|
||||
|
||||
if ($state === PHPIREDIS_READER_STATE_COMPLETE) {
|
||||
return phpiredis_reader_get_reply($reader);
|
||||
} else {
|
||||
$this->onProtocolError(phpiredis_reader_get_error($reader));
|
||||
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function writeRequest(CommandInterface $command)
|
||||
{
|
||||
$arguments = $command->getArguments();
|
||||
array_unshift($arguments, $command->getId());
|
||||
|
||||
$this->write(phpiredis_format_command($arguments));
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __wakeup()
|
||||
{
|
||||
$this->assertExtensions();
|
||||
$this->reader = $this->createReader();
|
||||
}
|
||||
}
|
||||
@@ -1,366 +0,0 @@
|
||||
<?php
|
||||
|
||||
/*
|
||||
* This file is part of the Predis package.
|
||||
*
|
||||
* (c) 2009-2020 Daniele Alessandri
|
||||
* (c) 2021-2023 Till Krüss
|
||||
*
|
||||
* For the full copyright and license information, please view the LICENSE
|
||||
* file that was distributed with this source code.
|
||||
*/
|
||||
|
||||
namespace Predis\Connection;
|
||||
|
||||
use Closure;
|
||||
use InvalidArgumentException;
|
||||
use Predis\Command\CommandInterface;
|
||||
use Predis\NotSupportedException;
|
||||
use Predis\Protocol\ProtocolException;
|
||||
use Predis\Response\Error as ErrorResponse;
|
||||
use Predis\Response\Status as StatusResponse;
|
||||
|
||||
/**
|
||||
* This class implements a Predis connection that actually talks with Webdis
|
||||
* instead of connecting directly to Redis. It relies on the cURL extension to
|
||||
* communicate with the web server and the phpiredis extension to parse the
|
||||
* protocol for responses returned in the http response bodies.
|
||||
*
|
||||
* Some features are not yet available or they simply cannot be implemented:
|
||||
* - Pipelining commands.
|
||||
* - Publish / Subscribe.
|
||||
* - MULTI / EXEC transactions (not yet supported by Webdis).
|
||||
*
|
||||
* The connection parameters supported by this class are:
|
||||
*
|
||||
* - scheme: must be 'http'.
|
||||
* - host: hostname or IP address of the server.
|
||||
* - port: TCP port of the server.
|
||||
* - timeout: timeout to perform the connection (default is 5 seconds).
|
||||
* - user: username for authentication.
|
||||
* - pass: password for authentication.
|
||||
*
|
||||
* @see http://webd.is
|
||||
* @see http://github.com/nicolasff/webdis
|
||||
* @see http://github.com/seppo0010/phpiredis
|
||||
* @deprecated 2.1.2
|
||||
*/
|
||||
class WebdisConnection implements NodeConnectionInterface
|
||||
{
|
||||
private $parameters;
|
||||
private $resource;
|
||||
private $reader;
|
||||
|
||||
/**
|
||||
* @param ParametersInterface $parameters Initialization parameters for the connection.
|
||||
*
|
||||
* @throws InvalidArgumentException
|
||||
*/
|
||||
public function __construct(ParametersInterface $parameters)
|
||||
{
|
||||
$this->assertExtensions();
|
||||
|
||||
if ($parameters->scheme !== 'http') {
|
||||
throw new InvalidArgumentException("Invalid scheme: '{$parameters->scheme}'.");
|
||||
}
|
||||
|
||||
$this->parameters = $parameters;
|
||||
|
||||
$this->resource = $this->createCurl();
|
||||
$this->reader = $this->createReader();
|
||||
}
|
||||
|
||||
/**
|
||||
* Frees the underlying cURL and protocol reader resources when the garbage
|
||||
* collector kicks in.
|
||||
*/
|
||||
public function __destruct()
|
||||
{
|
||||
curl_close($this->resource);
|
||||
phpiredis_reader_destroy($this->reader);
|
||||
}
|
||||
|
||||
/**
|
||||
* Helper method used to throw on unsupported methods.
|
||||
*
|
||||
* @param string $method Name of the unsupported method.
|
||||
*
|
||||
* @throws NotSupportedException
|
||||
*/
|
||||
private function throwNotSupportedException($method)
|
||||
{
|
||||
$class = __CLASS__;
|
||||
throw new NotSupportedException("The method $class::$method() is not supported.");
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks if the cURL and phpiredis extensions are loaded in PHP.
|
||||
*/
|
||||
private function assertExtensions()
|
||||
{
|
||||
if (!extension_loaded('curl')) {
|
||||
throw new NotSupportedException(
|
||||
'The "curl" extension is required by this connection backend.'
|
||||
);
|
||||
}
|
||||
|
||||
if (!extension_loaded('phpiredis')) {
|
||||
throw new NotSupportedException(
|
||||
'The "phpiredis" extension is required by this connection backend.'
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Initializes cURL.
|
||||
*
|
||||
* @return resource
|
||||
*/
|
||||
private function createCurl()
|
||||
{
|
||||
$parameters = $this->getParameters();
|
||||
$timeout = (isset($parameters->timeout) ? (float) $parameters->timeout : 5.0) * 1000;
|
||||
|
||||
if (filter_var($host = $parameters->host, FILTER_VALIDATE_IP, FILTER_FLAG_IPV6)) {
|
||||
$host = "[$host]";
|
||||
}
|
||||
|
||||
$options = [
|
||||
CURLOPT_FAILONERROR => true,
|
||||
CURLOPT_CONNECTTIMEOUT_MS => $timeout,
|
||||
CURLOPT_URL => "$parameters->scheme://$host:$parameters->port",
|
||||
CURLOPT_HTTP_VERSION => CURL_HTTP_VERSION_1_1,
|
||||
CURLOPT_POST => true,
|
||||
CURLOPT_WRITEFUNCTION => [$this, 'feedReader'],
|
||||
];
|
||||
|
||||
if (isset($parameters->user, $parameters->pass)) {
|
||||
$options[CURLOPT_USERPWD] = "{$parameters->user}:{$parameters->pass}";
|
||||
}
|
||||
|
||||
curl_setopt_array($resource = curl_init(), $options);
|
||||
|
||||
return $resource;
|
||||
}
|
||||
|
||||
/**
|
||||
* Initializes the phpiredis protocol reader.
|
||||
*
|
||||
* @return resource
|
||||
*/
|
||||
private function createReader()
|
||||
{
|
||||
$reader = phpiredis_reader_create();
|
||||
|
||||
phpiredis_reader_set_status_handler($reader, $this->getStatusHandler());
|
||||
phpiredis_reader_set_error_handler($reader, $this->getErrorHandler());
|
||||
|
||||
return $reader;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the handler used by the protocol reader for inline responses.
|
||||
*
|
||||
* @return Closure
|
||||
*/
|
||||
protected function getStatusHandler()
|
||||
{
|
||||
static $statusHandler;
|
||||
|
||||
if (!$statusHandler) {
|
||||
$statusHandler = function ($payload) {
|
||||
return StatusResponse::get($payload);
|
||||
};
|
||||
}
|
||||
|
||||
return $statusHandler;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the handler used by the protocol reader for error responses.
|
||||
*
|
||||
* @return Closure
|
||||
*/
|
||||
protected function getErrorHandler()
|
||||
{
|
||||
static $errorHandler;
|
||||
|
||||
if (!$errorHandler) {
|
||||
$errorHandler = function ($errorMessage) {
|
||||
return new ErrorResponse($errorMessage);
|
||||
};
|
||||
}
|
||||
|
||||
return $errorHandler;
|
||||
}
|
||||
|
||||
/**
|
||||
* Feeds the phpredis reader resource with the data read from the network.
|
||||
*
|
||||
* @param resource $resource Reader resource.
|
||||
* @param string $buffer Buffer of data read from a connection.
|
||||
*
|
||||
* @return int
|
||||
*/
|
||||
protected function feedReader($resource, $buffer)
|
||||
{
|
||||
phpiredis_reader_feed($this->reader, $buffer);
|
||||
|
||||
return strlen($buffer);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function connect()
|
||||
{
|
||||
// NOOP
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function disconnect()
|
||||
{
|
||||
// NOOP
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function isConnected()
|
||||
{
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks if the specified command is supported by this connection class.
|
||||
*
|
||||
* @param CommandInterface $command Command instance.
|
||||
*
|
||||
* @return string
|
||||
* @throws NotSupportedException
|
||||
*/
|
||||
protected function getCommandId(CommandInterface $command)
|
||||
{
|
||||
switch ($commandID = $command->getId()) {
|
||||
case 'AUTH':
|
||||
case 'SELECT':
|
||||
case 'MULTI':
|
||||
case 'EXEC':
|
||||
case 'WATCH':
|
||||
case 'UNWATCH':
|
||||
case 'DISCARD':
|
||||
case 'MONITOR':
|
||||
throw new NotSupportedException("Command '$commandID' is not allowed by Webdis.");
|
||||
default:
|
||||
return $commandID;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function writeRequest(CommandInterface $command)
|
||||
{
|
||||
$this->throwNotSupportedException(__FUNCTION__);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function readResponse(CommandInterface $command)
|
||||
{
|
||||
$this->throwNotSupportedException(__FUNCTION__);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function executeCommand(CommandInterface $command)
|
||||
{
|
||||
$resource = $this->resource;
|
||||
$commandId = $this->getCommandId($command);
|
||||
|
||||
if ($arguments = $command->getArguments()) {
|
||||
$arguments = implode('/', array_map('urlencode', $arguments));
|
||||
$serializedCommand = "$commandId/$arguments.raw";
|
||||
} else {
|
||||
$serializedCommand = "$commandId.raw";
|
||||
}
|
||||
|
||||
curl_setopt($resource, CURLOPT_POSTFIELDS, $serializedCommand);
|
||||
|
||||
if (curl_exec($resource) === false) {
|
||||
$error = trim(curl_error($resource));
|
||||
$errno = curl_errno($resource);
|
||||
|
||||
throw new ConnectionException($this, "$error{$this->getParameters()}]", $errno);
|
||||
}
|
||||
|
||||
if (phpiredis_reader_get_state($this->reader) !== PHPIREDIS_READER_STATE_COMPLETE) {
|
||||
throw new ProtocolException($this, phpiredis_reader_get_error($this->reader));
|
||||
}
|
||||
|
||||
return phpiredis_reader_get_reply($this->reader);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function getResource()
|
||||
{
|
||||
return $this->resource;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function getParameters()
|
||||
{
|
||||
return $this->parameters;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function addConnectCommand(CommandInterface $command)
|
||||
{
|
||||
$this->throwNotSupportedException(__FUNCTION__);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function read()
|
||||
{
|
||||
$this->throwNotSupportedException(__FUNCTION__);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __toString()
|
||||
{
|
||||
return "{$this->parameters->host}:{$this->parameters->port}";
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __sleep()
|
||||
{
|
||||
return ['parameters'];
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc}
|
||||
*/
|
||||
public function __wakeup()
|
||||
{
|
||||
$this->assertExtensions();
|
||||
|
||||
$this->resource = $this->createCurl();
|
||||
$this->reader = $this->createReader();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user