Skip to content

Commit 3d0131c

Browse files
loks0nclaude
andcommitted
Require both Redis connections explicitly; drop the $work fallback
No implicit "$commands defaults to $receive" — callers pass both connections, which makes the two-connection model explicit and lets both be promoted. Pass the same connection twice when one suffices (sequential/inline). All call sites (workers, tests) updated accordingly. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 95064b1 commit 3d0131c

9 files changed

Lines changed: 21 additions & 20 deletions

File tree

‎src/Queue/Broker/Redis.php‎

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,6 @@ class Redis implements Publisher, Consumer
1414
private const int RECONNECT_BACKOFF_MS = 100;
1515
private const int RECONNECT_MAX_BACKOFF_MS = 5_000;
1616

17-
/** Carries acks and publishing; wrap in Locking when shared by coroutines. */
18-
private readonly Connection $commands;
19-
2017
private bool $closed = false;
2118
private int $reconnectAttempt = 0;
2219
private int $reconnectBackoffMs = self::RECONNECT_BACKOFF_MS;
@@ -30,11 +27,11 @@ class Redis implements Publisher, Consumer
3027
private $reconnectSuccessCallback = null;
3128

3229
public function __construct(
33-
// Drives the blocking receive loop and its claim writes (single caller).
30+
// Blocking receive loop + claim writes (single caller).
3431
private readonly Connection $receive,
35-
?Connection $commands = null,
32+
// Acks and publishing; wrap in Locking when shared by coroutines.
33+
private readonly Connection $commands,
3634
) {
37-
$this->commands = $commands ?? $receive;
3835
}
3936

4037
public function setReconnectCallback(?callable $callback): self

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@ class PoolTest extends Base
1515
protected function getPublisher(): Publisher
1616
{
1717
$pool = new UtopiaPool(new Stack(), 'redis', 1, function () {
18-
return new RedisBroker(new Redis('redis', 6379));
18+
return new RedisBroker(new Redis('redis', 6379), new Redis('redis', 6379));
1919
});
2020

2121
return new Pool($pool, $pool);

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ public function testReconnectCallbackReceivesAttemptContext(): void
1313
{
1414
$queue = new Queue('reconnect-callback');
1515
$connection = new FailingRedisConnection();
16-
$broker = new RedisBroker($connection);
16+
$broker = new RedisBroker($connection, $connection);
1717
$calls = [];
1818

1919
$broker->setReconnectCallback(function (Queue $queue, \Throwable $error, int $attempt, int $sleepMs) use (&$calls, $broker): void {
@@ -47,7 +47,7 @@ public function testReconnectSuccessCallbackReceivesAttemptCount(): void
4747
{
4848
$queue = new Queue('reconnect-success-callback');
4949
$connection = new RecoveringRedisConnection();
50-
$broker = new RedisBroker($connection);
50+
$broker = new RedisBroker($connection, $connection);
5151
$calls = [];
5252

5353
$broker->setReconnectCallback(fn () => null);

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,7 @@ public function testOneCoroutineNeverOverlaps(): void
3737
private function runWorker(int $messages, int $maxCoroutines): array
3838
{
3939
$connection = new InMemoryConnection();
40-
$broker = new Redis($connection);
40+
$broker = new Redis($connection, $connection);
4141
$queue = new Queue(self::QUEUE, self::NAMESPACE);
4242

4343
$active = 0;

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ private function getConnection(): RedisCluster
2020

2121
protected function getPublisher(): Publisher
2222
{
23-
return new Redis($this->getConnection());
23+
return new Redis($this->getConnection(), $this->getConnection());
2424
}
2525

2626
protected function getQueue(): Queue

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ private function getConnection(): Redis
1616

1717
protected function getPublisher(): Publisher
1818
{
19-
return new RedisBroker($this->getConnection());
19+
return new RedisBroker($this->getConnection(), $this->getConnection());
2020
}
2121

2222
protected function getQueue(): Queue

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,8 +11,7 @@ class WorkermanTest extends Base
1111
{
1212
protected function getPublisher(): Publisher
1313
{
14-
$connection = new Redis('redis', 6379);
15-
return new RedisPublisher($connection);
14+
return new RedisPublisher(new Redis('redis', 6379), new Redis('redis', 6379));
1615
}
1716

1817
protected function getQueue(): Queue

‎tests/Queue/servers/SwooleRedisCluster/worker.php‎

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -9,12 +9,14 @@
99
use Utopia\Queue\Server;
1010
use Utopia\Validator\Text;
1111

12+
$nodes = [
13+
'redis-cluster-0:6379',
14+
'redis-cluster-1:6379',
15+
'redis-cluster-2:6379',
16+
];
1217
$consumer = new Redis(
13-
new RedisCluster([
14-
'redis-cluster-0:6379',
15-
'redis-cluster-1:6379',
16-
'redis-cluster-2:6379',
17-
]),
18+
receive: new RedisCluster($nodes),
19+
commands: new RedisCluster($nodes),
1820
);
1921
$adapter = new Swoole($consumer, 12, 'swoole-redis-cluster');
2022
$server = new Server($adapter);

‎tests/Queue/servers/Workerman/worker.php‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,10 @@
99
use Utopia\Queue\Broker\Redis;
1010
use Utopia\Validator\Text;
1111

12-
$consumer = new Redis(new RedisConnection('redis'));
12+
$consumer = new Redis(
13+
receive: new RedisConnection('redis'),
14+
commands: new RedisConnection('redis'),
15+
);
1316
$adapter = new Workerman($consumer, 12, 'wokerman');
1417
$server = new Queue\Server($adapter);
1518

0 commit comments

Comments
 (0)