From 7c7cf901523d6ce738c8b83406da5f45d0f7e6c8 Mon Sep 17 00:00:00 2001 From: Vladyslav Vildanov <117659936+vladvildanov@users.noreply.github.com> Date: Thu, 8 Jan 2026 13:56:55 +0200 Subject: [PATCH] Improve connection handshake by pipelining commands (#1622) * Improve connection handshake by pipelining commands * Added exceptions handling for Redis < 6.0 * Updated CHANGELOG.md --- CHANGELOG.md | 1 + src/Connection/StreamConnection.php | 38 +++++++++- tests/PHPUnit/PredisConnectionTestCase.php | 1 + .../Connection/StreamConnectionTest.php | 74 ++++++++++++++++--- 4 files changed, 100 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b7d9ed5e..2eeafc97 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ ### Fixed - Fixed wrong `@param` annotation in `Parameters` (#1614) - Make `ZRANDMEMBER` prefixable (#1621) +- Improve connection handshake by pipelining commands (#1622) ### Added - Added retry support (#1616) diff --git a/src/Connection/StreamConnection.php b/src/Connection/StreamConnection.php index 67ed8ab4..adb8dfc2 100644 --- a/src/Connection/StreamConnection.php +++ b/src/Connection/StreamConnection.php @@ -90,14 +90,46 @@ class StreamConnection extends AbstractConnection public function connect() { if (parent::connect() && $this->initCommands) { - foreach ($this->initCommands as $command) { - $response = $this->executeCommand($command); + $responses = $this->sendPipeline($this->initCommands); - $this->handleOnConnectResponse($response, $command); + if ($responses[0][0] instanceof ErrorResponseInterface) { + // Error in HELLO command, Redis < 6.0. + // We need to handle it separately and re-send other commands. + $this->handleOnConnectResponse($responses[0][0], $responses[0][1]); + $responses = $this->sendPipeline(array_slice($this->initCommands, 1)); + } + + foreach ($responses as $response) { + $this->handleOnConnectResponse($response[0], $response[1]); } } } + /** + * Sends commands to the server as pipeline and returns responses. + * + * @param CommandInterface[] $commands + * @return array + * @throws CommunicationException + */ + protected function sendPipeline(array $commands): array + { + $serialisedCommands = ''; + + foreach ($commands as $command) { + $serialisedCommands .= $command->serializeCommand(); + } + + $this->write($serialisedCommands); + $responses = []; + + foreach ($commands as $command) { + $responses[] = [$this->readResponse($command), $command]; + } + + return $responses; + } + /** * {@inheritdoc} */ diff --git a/tests/PHPUnit/PredisConnectionTestCase.php b/tests/PHPUnit/PredisConnectionTestCase.php index 9cfd2134..b5405ab4 100644 --- a/tests/PHPUnit/PredisConnectionTestCase.php +++ b/tests/PHPUnit/PredisConnectionTestCase.php @@ -341,6 +341,7 @@ abstract class PredisConnectionTestCase extends PredisTestCase /** * @group connected + * @requiresRedisVersion >= 6.0.0 */ public function testSendsInitializationCommandsOnConnection(): void { diff --git a/tests/Predis/Connection/StreamConnectionTest.php b/tests/Predis/Connection/StreamConnectionTest.php index d9796f2c..6ca97ada 100644 --- a/tests/Predis/Connection/StreamConnectionTest.php +++ b/tests/Predis/Connection/StreamConnectionTest.php @@ -129,19 +129,16 @@ class StreamConnectionTest extends PredisConnectionTestCase ->withAnyParameters() ->willReturn($this->mockStream); + // All handshake commands should be pipelined in a single write + $pipelinedCommands = $command1->serializeCommand() + . $command2->serializeCommand() + . $command3->serializeCommand(); + $this->mockStream - ->expects($this->exactly(3)) + ->expects($this->once()) ->method('write') - ->withConsecutive( - [$command1->serializeCommand()], - [$command2->serializeCommand()], - [$command3->serializeCommand()] - ) - ->willReturnOnConsecutiveCalls( - strlen($command1->serializeCommand()), - strlen($command2->serializeCommand()), - strlen($command3->serializeCommand()) - ); + ->with($pipelinedCommands) + ->willReturn(strlen($pipelinedCommands)); $this->mockStream ->expects($this->exactly(3)) @@ -158,6 +155,61 @@ class StreamConnectionTest extends PredisConnectionTestCase $connection->connect(); } + /** + * @group disconnected + */ + public function testHandshakeCommandsArePipelinedInSingleNetworkRoundTrip(): void + { + $parameters = new Parameters(); + $command1 = new RawCommand('AUTH', ['username', 'password']); + $command2 = new RawCommand('SELECT', [5]); + $command3 = new RawCommand('CLIENT', ['SETNAME', 'predis']); + $command4 = new RawCommand('CLIENT', ['SETINFO', 'LIB-NAME', 'predis']); + + $this->mockStreamFactory + ->expects($this->once()) + ->method('createStream') + ->withAnyParameters() + ->willReturn($this->mockStream); + + // Verify that all handshake commands are serialized and sent together + // in a single write operation to reduce initial handshake latency + $pipelinedCommands = $command1->serializeCommand() + . $command2->serializeCommand() + . $command3->serializeCommand() + . $command4->serializeCommand(); + + $this->mockStream + ->expects($this->once()) + ->method('write') + ->with($pipelinedCommands) + ->willReturn(strlen($pipelinedCommands)); + + // Verify that responses are read separately for each command + $this->mockStream + ->expects($this->exactly(4)) + ->method('read') + ->with(-1) + ->willReturnOnConsecutiveCalls( + '+OK\r\n', + '+OK\r\n', + '+OK\r\n', + '+OK\r\n' + ); + + $connection = new StreamConnection($parameters, $this->mockStreamFactory); + + $connection->addConnectCommand($command1); + $connection->addConnectCommand($command2); + $connection->addConnectCommand($command3); + $connection->addConnectCommand($command4); + + $connection->connect(); + + // Verify connection is established + $this->assertTrue($connection->isConnected()); + } + /** * @group disconnected */