From db30d5231d97fd9ecd014b7f905d2161e679376c Mon Sep 17 00:00:00 2001 From: LounisBou Date: Fri, 18 Sep 2026 15:03:37 +0200 Subject: [PATCH] fix(redis): ping the subscription connection to detect half-open sockets --- src/Hub/Transport/Redis/RedisTransport.php | 19 ++++++------ .../Hub/Transport/Redis/RedisClientStub.php | 12 ++++++++ .../Transport/Redis/RedisTransportTest.php | 30 +++++++++++++++++++ 3 files changed, 52 insertions(+), 9 deletions(-) diff --git a/src/Hub/Transport/Redis/RedisTransport.php b/src/Hub/Transport/Redis/RedisTransport.php index 904d84f..06655f7 100644 --- a/src/Hub/Transport/Redis/RedisTransport.php +++ b/src/Hub/Transport/Redis/RedisTransport.php @@ -55,17 +55,18 @@ public function __construct( $this->subscriber->on('unsubscribe', fn () => Hub::die(new RuntimeException('Redis connection lost'))); } - /** - * @codeCoverageIgnore - */ private function ping(): void { - /** @var PromiseInterface $ping */ - $ping = $this->redis->ping(); // @phpstan-ignore-line - $ping = maybeTimeout($ping, $this->options['readTimeout']); - $ping->then( - onRejected: Hub::die(...), - ); + // The subscriber connection only ever reads, so a half-open socket goes + // unnoticed there until something is written to it. + foreach ([$this->redis, $this->subscriber] as $client) { + /** @var PromiseInterface $ping */ + $ping = $client->ping(); // @phpstan-ignore-line + $ping = maybeTimeout($ping, $this->options['readTimeout']); + $ping->then( + onRejected: Hub::die(...), + ); + } } public function subscribe(callable $callback): void diff --git a/tests/Unit/Hub/Transport/Redis/RedisClientStub.php b/tests/Unit/Hub/Transport/Redis/RedisClientStub.php index 1485c00..3fc6824 100644 --- a/tests/Unit/Hub/Transport/Redis/RedisClientStub.php +++ b/tests/Unit/Hub/Transport/Redis/RedisClientStub.php @@ -9,6 +9,7 @@ use Evenement\EventEmitter; use Evenement\EventEmitterInterface; use Pest\Exceptions\ShouldNotHappen; +use React\Promise\Promise; use React\Promise\PromiseInterface; use function abs; @@ -21,12 +22,23 @@ final class RedisClientStub implements Client { public array $subscribedChannels = []; + public int $pings = 0; + + public bool $answersPings = true; + public function __construct( public readonly ArrayObject $storage = new ArrayObject(), private EventEmitterInterface $eventEmitter = new EventEmitter(), ) { } + public function ping(): PromiseInterface + { + ++$this->pings; + + return $this->answersPings ? resolve(true) : new Promise(static fn () => null); + } + public function subscribe(string $channel): void { $this->subscribedChannels[] = $channel; diff --git a/tests/Unit/Hub/Transport/Redis/RedisTransportTest.php b/tests/Unit/Hub/Transport/Redis/RedisTransportTest.php index 42335d4..4b911b1 100644 --- a/tests/Unit/Hub/Transport/Redis/RedisTransportTest.php +++ b/tests/Unit/Hub/Transport/Redis/RedisTransportTest.php @@ -111,3 +111,33 @@ expect($client->storage->getArrayCopy()['mercureUpdates'])->toHaveCount(3); }); + +it('pings the subscription connection as well as the command connection', function () { + $subscriber = new RedisClientStub(); + $redis = new RedisClientStub(); + + new RedisTransport($subscriber, $redis, options: ['pingInterval' => 0.01]); + + expect($redis->pings)->toBeGreaterThan(0) + ->and($subscriber->pings)->toBeGreaterThan(0); +}); + +it('kills the hub when the subscription connection stops answering', function () { + $subscriber = new RedisClientStub(); + $subscriber->answersPings = false; + $redis = new RedisClientStub(); + + new RedisTransport($subscriber, $redis, options: [ + 'pingInterval' => 0.01, + 'readTimeout' => 0.01, + ]); + + $stoppedByGuard = false; + Loop::addTimer(1.0, function () use (&$stoppedByGuard) { + $stoppedByGuard = true; + Loop::stop(); + }); + Loop::run(); + + expect($stoppedByGuard)->toBeFalse(); +});