Compare commits

...

33 Commits

Author SHA1 Message Date
Daniele Alessandri ca422b0300 Bump VERSION, update CHANGELOG and release. 2011-02-12 11:46:08 +01:00
Daniele Alessandri 79372cb99b Add inline (p)subscribe via options on Predis\PubSubContext initialization. 2011-02-12 11:43:37 +01:00
Daniele Alessandri 6a2e4d0396 Bump year in LICENSE file. 2011-01-26 16:25:05 +01:00
Daniele Alessandri 78814a473c Update CHANGELOG. 2011-01-26 16:24:50 +01:00
Daniele Alessandri 0d7fe31110 Minor indentation fix. 2011-01-26 16:24:34 +01:00
Daniele Alessandri 118af2809c Use only one check for replies that should not be passed to a reply parser. 2011-01-14 17:07:00 +01:00
Daniele Alessandri 8dcd10dbbc Fix missing registration for the Redis 2.2 profile. 2011-01-13 12:32:31 +01:00
Daniele Alessandri 3a6907241b Promote the profile for Redis 2.2 as stable (the default for the client is still 2.0). 2011-01-13 12:26:19 +01:00
Daniele Alessandri eae29bc3da Do not use is_numeric() when it is not really needed (it is relatively slow). 2011-01-13 12:18:24 +01:00
Daniele Alessandri 727feb3f27 Remove usage of constants in the protocol handlers. 2011-01-13 12:13:46 +01:00
Daniele Alessandri 3924235501 Apply minor changes. 2011-01-13 12:09:46 +01:00
Daniele Alessandri dc704c8cc4 Remove unused variables. 2011-01-13 12:07:54 +01:00
Daniele Alessandri 8d6f65d3dd Add the 'on_retry' callback as an option for Predis\MultiExecBlock. 2011-01-13 12:06:11 +01:00
Daniele Alessandri 40dfabb139 Do not perform useless read operations in the bulk reply handler. 2011-01-13 00:17:03 +01:00
Daniele Alessandri a605190354 Update README. 2011-01-06 11:49:15 +01:00
Daniele Alessandri b5ce81a030 Update CHANGELOG. 2011-01-06 11:40:39 +01:00
Daniele Alessandri a9e1d6b86b Reflect the status of the current branch. 2011-01-06 11:37:12 +01:00
Daniele Alessandri 923d998f35 Fix namespacing issue. 2011-01-06 00:59:12 +01:00
Daniele Alessandri eae8fb8971 Minor code style change. 2011-01-06 00:55:38 +01:00
Daniele Alessandri 4fc5ee65fd Look for a closing curly brace in a string only after the index of the opening one when using key tags. 2011-01-06 00:55:38 +01:00
Daniele Alessandri 007ddaecfa Optimize when a pair of curly brackets for key tagging is missing. 2011-01-06 00:55:38 +01:00
Daniele Alessandri 455e56927a Create a new class to handle the old response type of KEYS (Redis v1.2). 2011-01-06 00:55:38 +01:00
Daniele Alessandri 4b1302a93a Use a faster method to detect errors when reading a line from the server. 2011-01-06 00:55:00 +01:00
Daniele Alessandri d46b0e0785 Reuse code a bit. 2011-01-06 00:55:00 +01:00
Daniele Alessandri cf522ff2b4 Use a local cache for handlers while iterating chunks of a multi-bulk reply. 2011-01-06 00:54:10 +01:00
Daniele Alessandri db92c7a9b8 Reduce overhead by directly managing reply handlers for multibulk replies. 2011-01-06 00:54:10 +01:00
Daniele Alessandri 30926def60 Rewrite the bulk reply handler (more compact and faster code). 2011-01-06 00:54:10 +01:00
Daniele Alessandri 7465daa0eb More strict check of the length argument for a network read. 2011-01-06 00:52:56 +01:00
Daniele Alessandri 2c761d6c95 Remove a useless check for the payload length. 2011-01-06 00:52:56 +01:00
Daniele Alessandri 1e32f8aaa8 Use an explicit index to add values to the multi-bulk reply array. 2011-01-06 00:52:55 +01:00
Daniele Alessandri 2a25b0e3f3 Nothing fancy, and hardly an optimization. 2011-01-06 00:52:55 +01:00
Daniele Alessandri 5a6a48fa17 Avoid allocating an iterator to traverse the arguments list of a command. 2011-01-06 00:52:55 +01:00
Daniele Alessandri 64951b1799 Preallocate an empty array for empty command arguments. 2011-01-06 00:52:04 +01:00
6 changed files with 145 additions and 96 deletions
+11
View File
@@ -1,3 +1,14 @@
v0.6.4 (2011-02-12)
* Various performance improvements (15% ~ 25%) especially when dealing with
long multibulk replies or when using clustered connections.
* Added the "on_retry" option to Predis\MultiExecBlock that can be used to
specify an external callback (or any callable object) that gets invoked
whenever a transaction is aborted by the server.
* Added inline (p)subscribtion via options when initializing an instance of
Predis\PubSubContext.
v0.6.3 (2011-01-01)
* New commands available in the Redis v2.2 profile (dev):
- Strings: SETRANGE, GETRANGE, SETBIT, GETBIT
+1 -1
View File
@@ -1,4 +1,4 @@
Copyright (c) 2009-2010 Daniele Alessandri
Copyright (c) 2009-2011 Daniele Alessandri
Permission is hereby granted, free of charge, to any person
obtaining a copy of this software and associated documentation
+1 -1
View File
@@ -17,7 +17,7 @@ to be implemented soon in Predis.
## Main features ##
- Full support for Redis 2.0. Different versions of Redis are supported via server profiles.
- Full support for Redis 2.0 and 2.2. Different versions of Redis are supported via server profiles.
- Client-side sharding (support for consistent hashing and custom distribution strategies).
- Command pipelining on single and multiple connections (transparent).
- Abstraction for Redis transactions (>= 2.0) with support for CAS operations (>= 2.2).
+1 -1
View File
@@ -1 +1 @@
0.6.3
0.6.4
+129 -93
View File
@@ -287,8 +287,8 @@ class Client {
return $transBlock !== null ? $multi->execute($transBlock) : $multi;
}
public function pubSubContext() {
return new PubSubContext($this);
public function pubSubContext(Array $options = null) {
return new PubSubContext($this, $options);
}
}
@@ -417,7 +417,8 @@ class Protocol {
}
abstract class Command {
private $_arguments, $_hash;
private $_hash;
private $_arguments = array();
public abstract function getCommandId();
@@ -431,22 +432,24 @@ abstract class Command {
if (isset($this->_hash)) {
return $this->_hash;
}
else {
if (isset($this->_arguments[0])) {
// TODO: should we throw an exception if the command does
// not support sharding?
$key = $this->_arguments[0];
$start = strpos($key, '{');
$end = strpos($key, '}');
if ($start !== false && $end !== false) {
if (isset($this->_arguments[0])) {
// TODO: should we throw an exception if the command does
// not support sharding?
$key = $this->_arguments[0];
$start = strpos($key, '{');
if ($start !== false) {
$end = strpos($key, '}', $start);
if ($end !== false) {
$key = substr($key, ++$start, $end - $start);
}
$this->_hash = $distributor->generateKey($key);
return $this->_hash;
}
$this->_hash = $distributor->generateKey($key);
return $this->_hash;
}
return null;
}
@@ -469,11 +472,13 @@ abstract class Command {
}
public function getArguments() {
return isset($this->_arguments) ? $this->_arguments : array();
return $this->_arguments;
}
public function getArgument($index = 0) {
return isset($this->_arguments[$index]) ? $this->_arguments[$index] : null;
if (isset($this->_arguments[$index]) === true) {
return $this->_arguments[$index];
}
}
public function parseResponse($data) {
@@ -491,8 +496,7 @@ abstract class InlineCommand extends Command {
$arguments[0] = implode($arguments[0], ' ');
}
return $command . (count($arguments) > 0
? ' ' . implode($arguments, ' ') . Protocol::NEWLINE
: Protocol::NEWLINE
? ' ' . implode($arguments, ' ') . "\r\n" : "\r\n"
);
}
}
@@ -504,7 +508,7 @@ abstract class BulkCommand extends Command {
$data = implode($data, ' ');
}
return $command . ' ' . implode($arguments, ' ') . ' ' . strlen($data) .
Protocol::NEWLINE . $data . Protocol::NEWLINE;
"\r\n" . $data . "\r\n";
}
}
@@ -521,14 +525,14 @@ abstract class MultiBulkCommand extends Command {
$cmd_args = $arguments;
}
$newline = Protocol::NEWLINE;
$cmdlen = strlen($command);
$reqlen = $argsc + 1;
$buffer = "*{$reqlen}{$newline}\${$cmdlen}{$newline}{$command}{$newline}";
foreach ($cmd_args as $argument) {
$buffer = "*{$reqlen}\r\n\${$cmdlen}\r\n{$command}\r\n";
for ($i = 0; $i < $reqlen - 1; $i++) {
$argument = $cmd_args[$i];
$arglen = strlen($argument);
$buffer .= "\${$arglen}{$newline}{$argument}{$newline}";
$buffer .= "\${$arglen}\r\n{$argument}\r\n";
}
return $buffer;
@@ -543,10 +547,10 @@ interface IResponseHandler {
class ResponseStatusHandler implements IResponseHandler {
public function handle(Connection $connection, $status) {
if ($status === Protocol::OK) {
if ($status === 'OK') {
return true;
}
else if ($status === Protocol::QUEUED) {
if ($status === 'QUEUED') {
return new ResponseQueued();
}
return $status;
@@ -566,44 +570,31 @@ class ResponseErrorSilentHandler implements IResponseHandler {
}
class ResponseBulkHandler implements IResponseHandler {
public function handle(Connection $connection, $dataLength) {
if (!is_numeric($dataLength)) {
public function handle(Connection $connection, $lengthString) {
$length = (int) $lengthString;
if ($length != $lengthString) {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, "Cannot parse '$dataLength' as data length"
$connection, "Cannot parse '$length' as data length"
));
}
if ($dataLength > 0) {
$value = $connection->readBytes($dataLength);
self::discardNewLine($connection);
return $value;
if ($length >= 0) {
return $length > 0 ? substr($connection->readBytes($length + 2), 0, -2) : '';
}
else if ($dataLength == 0) {
self::discardNewLine($connection);
return '';
}
return null;
}
private static function discardNewLine(Connection $connection) {
if ($connection->readBytes(2) !== Protocol::NEWLINE) {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, 'Did not receive a new-line at the end of a bulk response'
));
if ($length == -1) {
return null;
}
}
}
class ResponseMultiBulkHandler implements IResponseHandler {
public function handle(Connection $connection, $rawLength) {
if (!is_numeric($rawLength)) {
public function handle(Connection $connection, $lengthString) {
$listLength = (int) $lengthString;
if ($listLength != $lengthString) {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, "Cannot parse '$rawLength' as data length"
$connection, "Cannot parse '$lengthString' as data length"
));
}
$listLength = (int) $rawLength;
if ($listLength === -1) {
return null;
}
@@ -611,9 +602,19 @@ class ResponseMultiBulkHandler implements IResponseHandler {
$list = array();
if ($listLength > 0) {
$handlers = array();
$reader = $connection->getResponseReader();
for ($i = 0; $i < $listLength; $i++) {
$list[] = $reader->read($connection);
$header = $connection->readLine();
$prefix = $header[0];
if (isset($handlers[$prefix])) {
$handler = $handlers[$prefix];
}
else {
$handler = $reader->getHandler($prefix);
$handlers[$prefix] = $handler;
}
$list[$i] = $handler->handle($connection, substr($header, 1));
}
}
@@ -622,13 +623,14 @@ class ResponseMultiBulkHandler implements IResponseHandler {
}
class ResponseMultiBulkStreamHandler implements IResponseHandler {
public function handle(Connection $connection, $rawLength) {
if (!is_numeric($rawLength)) {
public function handle(Connection $connection, $lengthString) {
$listLength = (int) $lengthString;
if ($listLength != $lengthString) {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, "Cannot parse '$rawLength' as data length"
$connection, "Cannot parse '$lengthString' as data length"
));
}
return new Shared\MultiBulkResponseIterator($connection, (int)$rawLength);
return new Shared\MultiBulkResponseIterator($connection, $lengthString);
}
}
@@ -637,14 +639,12 @@ class ResponseIntegerHandler implements IResponseHandler {
if (is_numeric($number)) {
return (int) $number;
}
else {
if ($number !== Protocol::NULL) {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, "Cannot parse '$number' as numeric response"
));
}
return null;
if ($number !== 'nil') {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, "Cannot parse '$number' as numeric response"
));
}
return null;
}
}
@@ -678,26 +678,26 @@ class ResponseReader {
public function read(Connection $connection) {
$header = $connection->readLine();
if ($header === '') {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, 'Unexpected empty header'
));
$this->throwMalformedResponse($connection, 'Unexpected empty header');
}
$prefix = $header[0];
$payload = strlen($header) > 1 ? substr($header, 1) : '';
if (!isset($this->_prefixHandlers[$prefix])) {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, "Unknown prefix '$prefix'"
));
$this->throwMalformedResponse($connection, "Unknown prefix '$prefix'");
}
$handler = $this->_prefixHandlers[$prefix];
return $handler->handle($connection, $payload);
return $handler->handle($connection, substr($header, 1));
}
private function throwMalformedResponse(Connection $connection, $message) {
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$connection, $message
));
}
}
class ResponseError {
public $skipParse = true;
private $_message;
public function __construct($message) {
@@ -705,10 +705,10 @@ class ResponseError {
}
public function __get($property) {
if ($property == 'error') {
if ($property === 'error') {
return true;
}
if ($property == 'message') {
if ($property === 'message') {
return $this->_message;
}
}
@@ -723,11 +723,21 @@ class ResponseError {
}
class ResponseQueued {
public $queued = true;
public $skipParse = true;
public function __toString() {
return Protocol::QUEUED;
}
public function __get($property) {
if ($property === 'queued') {
return true;
}
}
public function __isset($property) {
return $property === 'queued';
}
}
/* ------------------------------------------------------------------------- */
@@ -874,7 +884,7 @@ class MultiExecBlock {
}
$command = $client->createCommand($method, $arguments);
$response = $client->executeCommand($command);
if (!isset($response->queued)) {
if (!$response instanceof \Predis\ResponseQueued) {
$this->malformedServerResponse(
'The server did not respond with a QUEUED status reply'
);
@@ -987,6 +997,9 @@ class MultiExecBlock {
);
}
$this->reset();
if (isset($this->_options['on_retry']) && is_callable($this->_options['on_retry'])) {
call_user_func($this->_options['on_retry'], $this, $attemptsLeft);
}
continue;
}
break;
@@ -1016,7 +1029,7 @@ class MultiExecBlock {
// Since a MULTI/EXEC block cannot be initialized over a clustered
// connection, we can safely assume that Predis\Client::getConnection()
// will always return an instance of Predis\Connection.
Utils::onCommunicationException(new MalformedServerResponse(
Shared\Utils::onCommunicationException(new MalformedServerResponse(
$this->_redisClient->getConnection(), $message
));
}
@@ -1034,12 +1047,16 @@ class PubSubContext implements \Iterator {
const STATUS_SUBSCRIBED = 0x0010;
const STATUS_PSUBSCRIBED = 0x0100;
private $_redisClient, $_subscriptions, $_isStillValid, $_position;
private $_redisClient, $_position, $_options;
public function __construct(Client $redisClient) {
public function __construct(Client $redisClient, Array $options = null) {
$this->checkCapabilities($redisClient);
$this->_options = $options ?: array();
$this->_redisClient = $redisClient;
$this->_statusFlags = self::STATUS_VALID;
$this->genericSubscribeInit('subscribe');
$this->genericSubscribeInit('psubscribe');
}
public function __destruct() {
@@ -1061,6 +1078,19 @@ class PubSubContext implements \Iterator {
}
}
private function genericSubscribeInit($subscribeAction) {
if (isset($this->_options[$subscribeAction])) {
if (is_array($this->_options[$subscribeAction])) {
foreach ($this->_options[$subscribeAction] as $subscription) {
$this->$subscribeAction($subscription);
}
}
else {
$this->$subscribeAction($this->_options[$subscribeAction]);
}
}
}
private function isFlagSet($value) {
return ($this->_statusFlags & $value) === $value;
}
@@ -1342,8 +1372,7 @@ class Connection implements IConnection {
public function readResponse(Command $command) {
$response = $this->_reader->read($this);
$skipparse = isset($response->queued) || isset($response->error);
return $skipparse ? $response : $command->parseResponse($response);
return isset($response->skipParse) ? $response : $command->parseResponse($response);
}
public function executeCommand(Command $command) {
@@ -1379,7 +1408,7 @@ class Connection implements IConnection {
}
public function readBytes($length) {
if ($length == 0) {
if ($length <= 0) {
throw new \InvalidArgumentException('Length parameter must be greater than 0');
}
$socket = $this->getSocket();
@@ -1400,12 +1429,12 @@ class Connection implements IConnection {
$value = '';
do {
$chunk = fgets($socket);
if ($chunk === false || strlen($chunk) == 0) {
if ($chunk === false || $chunk === '') {
$this->onCommunicationException('Error while reading line from the server');
}
$value .= $chunk;
}
while (substr($value, -2) !== Protocol::NEWLINE);
while (substr($value, -2) !== "\r\n");
return substr($value, 0, -2);
}
@@ -1528,6 +1557,7 @@ abstract class RedisServerProfile {
return array(
'1.2' => '\Predis\RedisServer_v1_2',
'2.0' => '\Predis\RedisServer_v2_0',
'2.2' => '\Predis\RedisServer_v2_2',
'default' => '\Predis\RedisServer_v2_0',
'dev' => '\Predis\RedisServer_vNext',
);
@@ -1655,7 +1685,7 @@ class RedisServer_v1_2 extends RedisServerProfile {
'type' => '\Predis\Commands\Type',
/* commands operating on the key space */
'keys' => '\Predis\Commands\Keys',
'keys' => '\Predis\Commands\Keys_v1_2',
'randomkey' => '\Predis\Commands\RandomKey',
'randomKey' => '\Predis\Commands\RandomKey',
'rename' => '\Predis\Commands\Rename',
@@ -1789,6 +1819,9 @@ class RedisServer_v2_0 extends RedisServer_v1_2 {
'append' => '\Predis\Commands\Append',
'substr' => '\Predis\Commands\Substr',
/* commands operating on the key space */
'keys' => '\Predis\Commands\Keys',
/* commands operating on lists */
'blpop' => '\Predis\Commands\ListPopFirstBlocking',
'popFirstBlocking' => '\Predis\Commands\ListPopFirstBlocking',
@@ -1849,8 +1882,8 @@ class RedisServer_v2_0 extends RedisServer_v1_2 {
}
}
class RedisServer_vNext extends RedisServer_v2_0 {
public function getVersion() { return '2.1'; }
class RedisServer_v2_2 extends RedisServer_v2_0 {
public function getVersion() { return '2.2'; }
public function getSupportedCommands() {
return array_merge(parent::getSupportedCommands(), array(
/* transactions */
@@ -1879,6 +1912,10 @@ class RedisServer_vNext extends RedisServer_v2_0 {
}
}
class RedisServer_vNext extends RedisServer_v2_2 {
public function getVersion() { return 'DEV'; }
}
/* ------------------------------------------------------------------------- */
namespace Predis\Pipeline;
@@ -2422,12 +2459,11 @@ class Strlen extends \Predis\MultiBulkCommand {
class Keys extends \Predis\MultiBulkCommand {
public function canBeHashed() { return false; }
public function getCommandId() { return 'KEYS'; }
public function parseResponse($data) {
// TODO: is this behaviour correct?
if (is_array($data) || $data instanceof \Iterator) {
return $data;
}
return strlen($data) > 0 ? explode(' ', $data) : array();
}
class Keys_v1_2 extends Keys {
public function parseResponse($data) {
return explode(' ', $data);
}
}
+2
View File
@@ -216,6 +216,7 @@ class PredisClientFeaturesTestSuite extends PHPUnit_Framework_TestCase {
function testResponseQueued() {
$response = new \Predis\ResponseQueued();
$this->assertTrue($response->skipParse);
$this->assertTrue($response->queued);
$this->assertEquals(\Predis\Protocol::QUEUED, (string)$response);
}
@@ -227,6 +228,7 @@ class PredisClientFeaturesTestSuite extends PHPUnit_Framework_TestCase {
$errorMessage = 'ERROR MESSAGE';
$response = new \Predis\ResponseError($errorMessage);
$this->assertTrue($response->skipParse);
$this->assertTrue($response->error);
$this->assertEquals($errorMessage, $response->message);
$this->assertEquals($errorMessage, (string)$response);