PhpiredisConnection.php 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394
  1. <?php
  2. /*
  3. * This file is part of the Predis package.
  4. *
  5. * (c) Daniele Alessandri <suppakilla@gmail.com>
  6. *
  7. * For the full copyright and license information, please view the LICENSE
  8. * file that was distributed with this source code.
  9. */
  10. namespace Predis\Connection;
  11. use Predis\NotSupportedException;
  12. use Predis\ResponseError;
  13. use Predis\ResponseQueued;
  14. use Predis\Command\CommandInterface;
  15. /**
  16. * This class provides the implementation of a Predis connection that uses the
  17. * PHP socket extension for network communication and wraps the phpiredis C
  18. * extension (PHP bindings for hiredis) to parse the Redis protocol. Everything
  19. * is highly experimental (even the very same phpiredis since it is quite new),
  20. * so use it at your own risk.
  21. *
  22. * This class is mainly intended to provide an optional low-overhead alternative
  23. * for processing replies from Redis compared to the standard pure-PHP classes.
  24. * Differences in speed when dealing with short inline replies are practically
  25. * nonexistent, the actual speed boost is for long multibulk replies when this
  26. * protocol processor can parse and return replies very fast.
  27. *
  28. * For instructions on how to build and install the phpiredis extension, please
  29. * consult the repository of the project.
  30. *
  31. * The connection parameters supported by this class are:
  32. *
  33. * - scheme: it can be either 'tcp' or 'unix'.
  34. * - host: hostname or IP address of the server.
  35. * - port: TCP port of the server.
  36. * - path: path of a UNIX domain socket when scheme is 'unix'.
  37. * - timeout: timeout to perform the connection.
  38. * - read_write_timeout: timeout of read / write operations.
  39. *
  40. * @link http://github.com/nrk/phpiredis
  41. * @author Daniele Alessandri <suppakilla@gmail.com>
  42. */
  43. class PhpiredisConnection extends AbstractConnection
  44. {
  45. const ERR_MSG_EXTENSION = 'The %s extension must be loaded in order to be able to use this connection class';
  46. private $reader;
  47. /**
  48. * {@inheritdoc}
  49. */
  50. public function __construct(ConnectionParametersInterface $parameters)
  51. {
  52. $this->checkExtensions();
  53. $this->initializeReader();
  54. parent::__construct($parameters);
  55. }
  56. /**
  57. * Disconnects from the server and destroys the underlying resource and the
  58. * protocol reader resource when PHP's garbage collector kicks in.
  59. */
  60. public function __destruct()
  61. {
  62. phpiredis_reader_destroy($this->reader);
  63. parent::__destruct();
  64. }
  65. /**
  66. * Checks if the socket and phpiredis extensions are loaded in PHP.
  67. */
  68. private function checkExtensions()
  69. {
  70. if (!function_exists('socket_create')) {
  71. throw new NotSupportedException(sprintf(self::ERR_MSG_EXTENSION, 'socket'));
  72. }
  73. if (!function_exists('phpiredis_reader_create')) {
  74. throw new NotSupportedException(sprintf(self::ERR_MSG_EXTENSION, 'phpiredis'));
  75. }
  76. }
  77. /**
  78. * {@inheritdoc}
  79. */
  80. protected function checkParameters(ConnectionParametersInterface $parameters)
  81. {
  82. if (isset($parameters->iterable_multibulk)) {
  83. $this->onInvalidOption('iterable_multibulk', $parameters);
  84. }
  85. if (isset($parameters->persistent)) {
  86. $this->onInvalidOption('persistent', $parameters);
  87. }
  88. return parent::checkParameters($parameters);
  89. }
  90. /**
  91. * Initializes the protocol reader resource.
  92. */
  93. private function initializeReader()
  94. {
  95. $reader = phpiredis_reader_create();
  96. phpiredis_reader_set_status_handler($reader, $this->getStatusHandler());
  97. phpiredis_reader_set_error_handler($reader, $this->getErrorHandler());
  98. $this->reader = $reader;
  99. }
  100. /**
  101. * Gets the handler used by the protocol reader to handle status replies.
  102. *
  103. * @return \Closure
  104. */
  105. private function getStatusHandler()
  106. {
  107. return function ($payload) {
  108. switch ($payload) {
  109. case 'OK':
  110. return true;
  111. case 'QUEUED':
  112. return new ResponseQueued();
  113. default:
  114. return $payload;
  115. }
  116. };
  117. }
  118. /**
  119. * Gets the handler used by the protocol reader to handle Redis errors.
  120. *
  121. * @param Boolean $throw_errors Specify if Redis errors throw exceptions.
  122. * @return \Closure
  123. */
  124. private function getErrorHandler()
  125. {
  126. return function ($errorMessage) {
  127. return new ResponseError($errorMessage);
  128. };
  129. }
  130. /**
  131. * Helper method used to throw exceptions on socket errors.
  132. */
  133. private function emitSocketError()
  134. {
  135. $errno = socket_last_error();
  136. $errstr = socket_strerror($errno);
  137. $this->disconnect();
  138. $this->onConnectionError(trim($errstr), $errno);
  139. }
  140. /**
  141. * {@inheritdoc}
  142. */
  143. protected function createResource()
  144. {
  145. $parameters = $this->parameters;
  146. $isUnix = $this->parameters->scheme === 'unix';
  147. $domain = $isUnix ? AF_UNIX : AF_INET;
  148. $protocol = $isUnix ? 0 : SOL_TCP;
  149. $socket = @call_user_func('socket_create', $domain, SOCK_STREAM, $protocol);
  150. if (!is_resource($socket)) {
  151. $this->emitSocketError();
  152. }
  153. $this->setSocketOptions($socket, $parameters);
  154. return $socket;
  155. }
  156. /**
  157. * Sets options on the socket resource from the connection parameters.
  158. *
  159. * @param resource $socket Socket resource.
  160. * @param ConnectionParametersInterface $parameters Parameters used to initialize the connection.
  161. */
  162. private function setSocketOptions($socket, ConnectionParametersInterface $parameters)
  163. {
  164. if ($parameters->scheme !== 'tcp') {
  165. return;
  166. }
  167. if (!socket_set_option($socket, SOL_TCP, TCP_NODELAY, 1)) {
  168. $this->emitSocketError();
  169. }
  170. if (!socket_set_option($socket, SOL_SOCKET, SO_REUSEADDR, 1)) {
  171. $this->emitSocketError();
  172. }
  173. if (isset($parameters->read_write_timeout)) {
  174. $rwtimeout = $parameters->read_write_timeout;
  175. $timeoutSec = floor($rwtimeout);
  176. $timeoutUsec = ($rwtimeout - $timeoutSec) * 1000000;
  177. $timeout = array(
  178. 'sec' => $timeoutSec,
  179. 'usec' => $timeoutUsec,
  180. );
  181. if (!socket_set_option($socket, SOL_SOCKET, SO_SNDTIMEO, $timeout)) {
  182. $this->emitSocketError();
  183. }
  184. if (!socket_set_option($socket, SOL_SOCKET, SO_RCVTIMEO, $timeout)) {
  185. $this->emitSocketError();
  186. }
  187. }
  188. }
  189. /**
  190. * Gets the address from the connection parameters.
  191. *
  192. * @param ConnectionParametersInterface $parameters Parameters used to initialize the connection.
  193. * @return string
  194. */
  195. protected static function getAddress(ConnectionParametersInterface $parameters)
  196. {
  197. if ($parameters->scheme === 'unix') {
  198. return $parameters->path;
  199. }
  200. $host = $parameters->host;
  201. if (ip2long($host) === false) {
  202. if (false === $addresses = gethostbynamel($host)) {
  203. return false;
  204. }
  205. return $addresses[array_rand($addresses)];
  206. }
  207. return $host;
  208. }
  209. /**
  210. * Opens the actual connection to the server with a timeout.
  211. *
  212. * @param ConnectionParametersInterface $parameters Parameters used to initialize the connection.
  213. * @return string
  214. */
  215. private function connectWithTimeout(ConnectionParametersInterface $parameters)
  216. {
  217. if (false === $host = self::getAddress($parameters)) {
  218. $this->onConnectionError("Cannot resolve the address of '$parameters->host'.");
  219. }
  220. $socket = $this->getResource();
  221. socket_set_nonblock($socket);
  222. if (@socket_connect($socket, $host, $parameters->port) === false) {
  223. $error = socket_last_error();
  224. if ($error != SOCKET_EINPROGRESS && $error != SOCKET_EALREADY) {
  225. $this->emitSocketError();
  226. }
  227. }
  228. socket_set_block($socket);
  229. $null = null;
  230. $selectable = array($socket);
  231. $timeout = $parameters->timeout;
  232. $timeoutSecs = floor($timeout);
  233. $timeoutUSecs = ($timeout - $timeoutSecs) * 1000000;
  234. $selected = socket_select($selectable, $selectable, $null, $timeoutSecs, $timeoutUSecs);
  235. if ($selected === 2) {
  236. $this->onConnectionError('Connection refused', SOCKET_ECONNREFUSED);
  237. }
  238. if ($selected === 0) {
  239. $this->onConnectionError('Connection timed out', SOCKET_ETIMEDOUT);
  240. }
  241. if ($selected === false) {
  242. $this->emitSocketError();
  243. }
  244. }
  245. /**
  246. * {@inheritdoc}
  247. */
  248. public function connect()
  249. {
  250. parent::connect();
  251. $this->connectWithTimeout($this->parameters);
  252. if ($this->initCmds) {
  253. $this->sendInitializationCommands();
  254. }
  255. }
  256. /**
  257. * {@inheritdoc}
  258. */
  259. public function disconnect()
  260. {
  261. if ($this->isConnected()) {
  262. socket_close($this->getResource());
  263. parent::disconnect();
  264. }
  265. }
  266. /**
  267. * Sends the initialization commands to Redis when the connection is opened.
  268. */
  269. private function sendInitializationCommands()
  270. {
  271. foreach ($this->initCmds as $command) {
  272. $this->writeCommand($command);
  273. }
  274. foreach ($this->initCmds as $command) {
  275. $this->readResponse($command);
  276. }
  277. }
  278. /**
  279. * {@inheritdoc}
  280. */
  281. protected function write($buffer)
  282. {
  283. $socket = $this->getResource();
  284. while (($length = strlen($buffer)) > 0) {
  285. $written = socket_write($socket, $buffer, $length);
  286. if ($length === $written) {
  287. return;
  288. }
  289. if ($written === false) {
  290. $this->onConnectionError('Error while writing bytes to the server');
  291. }
  292. $buffer = substr($buffer, $written);
  293. }
  294. }
  295. /**
  296. * {@inheritdoc}
  297. */
  298. public function read()
  299. {
  300. $socket = $this->getResource();
  301. $reader = $this->reader;
  302. while (($state = phpiredis_reader_get_state($reader)) === PHPIREDIS_READER_STATE_INCOMPLETE) {
  303. if (@socket_recv($socket, $buffer, 4096, 0) === false || $buffer === '') {
  304. $this->emitSocketError();
  305. }
  306. phpiredis_reader_feed($reader, $buffer);
  307. }
  308. if ($state === PHPIREDIS_READER_STATE_COMPLETE) {
  309. return phpiredis_reader_get_reply($reader);
  310. } else {
  311. $this->onProtocolError(phpiredis_reader_get_error($reader));
  312. }
  313. }
  314. /**
  315. * {@inheritdoc}
  316. */
  317. public function writeCommand(CommandInterface $command)
  318. {
  319. $cmdargs = $command->getArguments();
  320. array_unshift($cmdargs, $command->getId());
  321. $this->write(phpiredis_format_command($cmdargs));
  322. }
  323. /**
  324. * {@inheritdoc}
  325. */
  326. public function __wakeup()
  327. {
  328. $this->checkExtensions();
  329. $this->initializeReader();
  330. }
  331. }