Skip to content

Commit 1bee005

Browse files
committed
feat(redis): expose reconnect success callback
1 parent 81b7779 commit 1bee005

2 files changed

Lines changed: 71 additions & 0 deletions

File tree

‎src/Queue/Broker/Redis.php‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,10 @@ class Redis implements Publisher, Consumer
1919
* @var (callable(Queue, \Throwable, int, int): void)|null
2020
*/
2121
private $reconnectCallback = null;
22+
/**
23+
* @var (callable(Queue, int): void)|null
24+
*/
25+
private $reconnectSuccessCallback = null;
2226

2327
public function __construct(private readonly Connection $connection)
2428
{
@@ -31,6 +35,13 @@ public function setReconnectCallback(?callable $callback): self
3135
return $this;
3236
}
3337

38+
public function setReconnectSuccessCallback(?callable $callback): self
39+
{
40+
$this->reconnectSuccessCallback = $callback;
41+
42+
return $this;
43+
}
44+
3445
public function consume(Queue $queue, callable $messageCallback, callable $successCallback, callable $errorCallback): void
3546
{
3647
$reconnectBackoffMs = self::RECONNECT_BACKOFF_MS;
@@ -42,6 +53,10 @@ public function consume(Queue $queue, callable $messageCallback, callable $succe
4253
*/
4354
try {
4455
$nextMessage = $this->connection->rightPopArray("{$queue->namespace}.queue.{$queue->name}", self::POP_TIMEOUT);
56+
if ($reconnectAttempt > 0) {
57+
$this->triggerReconnectSuccessCallback($queue, $reconnectAttempt);
58+
}
59+
4560
$reconnectBackoffMs = self::RECONNECT_BACKOFF_MS;
4661
$reconnectAttempt = 0;
4762
} catch (\RedisException|\RedisClusterException $e) {
@@ -147,6 +162,18 @@ private function triggerReconnectCallback(Queue $queue, \Throwable $error, int $
147162
}
148163
}
149164

165+
private function triggerReconnectSuccessCallback(Queue $queue, int $attempts): void
166+
{
167+
if (!\is_callable($this->reconnectSuccessCallback)) {
168+
return;
169+
}
170+
171+
try {
172+
($this->reconnectSuccessCallback)($queue, $attempts);
173+
} catch (\Throwable) {
174+
}
175+
}
176+
150177
public function enqueue(Queue $queue, array $payload, bool $priority = false): bool
151178
{
152179
$payload = [

‎tests/Queue/E2E/Adapter/RedisReconnectCallbackTest.php‎

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,36 @@ public function testReconnectCallbackReceivesAttemptContext(): void
4343
$this->assertGreaterThanOrEqual(0, $calls[0]['sleepMs']);
4444
$this->assertLessThanOrEqual(100, $calls[0]['sleepMs']);
4545
}
46+
47+
public function testReconnectSuccessCallbackReceivesAttemptCount(): void
48+
{
49+
$queue = new Queue('reconnect-success-callback');
50+
$connection = new RecoveringRedisConnection();
51+
$broker = new RedisBroker($connection);
52+
$calls = [];
53+
54+
$broker->setReconnectCallback(fn () => null);
55+
$broker->setReconnectSuccessCallback(function (Queue $queue, int $attempts) use (&$calls, $broker): void {
56+
$calls[] = [
57+
'queue' => $queue,
58+
'attempts' => $attempts,
59+
];
60+
61+
$broker->close();
62+
});
63+
64+
$broker->consume(
65+
$queue,
66+
fn () => null,
67+
fn () => null,
68+
fn () => null,
69+
);
70+
71+
$this->assertSame(2, $connection->popAttempts);
72+
$this->assertCount(1, $calls);
73+
$this->assertSame($queue, $calls[0]['queue']);
74+
$this->assertSame(1, $calls[0]['attempts']);
75+
}
4676
}
4777

4878
class FailingRedisConnection implements Connection
@@ -160,3 +190,17 @@ public function close(): void
160190
{
161191
}
162192
}
193+
194+
class RecoveringRedisConnection extends FailingRedisConnection
195+
{
196+
public function rightPopArray(string $queue, int $timeout): array|false
197+
{
198+
$this->popAttempts++;
199+
200+
if ($this->popAttempts === 1) {
201+
throw new \RedisException('Redis is unavailable.');
202+
}
203+
204+
return false;
205+
}
206+
}

0 commit comments

Comments
 (0)