Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 10 additions & 9 deletions src/Hub/Transport/Redis/RedisTransport.php
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions tests/Unit/Hub/Transport/Redis/RedisClientStub.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down
30 changes: 30 additions & 0 deletions tests/Unit/Hub/Transport/Redis/RedisTransportTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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();
});
Loading