Skip to content

Commit 2be4c1b

Browse files
loks0nclaude
andcommitted
Trim comments to non-obvious WHYs; rename Redis $work to $commands
- Cut the explanatory doc comments down to the few that carry non-obvious rationale (process() never throws, the two-connection split, the call_user_func_array argument ordering); drop the rest that just restated the code. - Rename the Redis broker's second connection from $work to $commands, which says what it carries (acks + publishing) rather than a vague "work". Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 487b64e commit 2be4c1b

6 files changed

Lines changed: 32 additions & 59 deletions

File tree

‎src/Queue/Adapter.php‎

Lines changed: 3 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -6,13 +6,10 @@
66

77
abstract class Adapter
88
{
9-
/** Seconds to block for a message before re-checking the stop flag. */
109
protected const int RECEIVE_TIMEOUT = 2;
1110

1211
public Queue $queue;
1312
protected ?Container $context = null;
14-
15-
/** Set to break out of the receive loop. */
1613
protected bool $stopped = false;
1714

1815
public function __construct(
@@ -37,9 +34,6 @@ abstract public function start(): self;
3734
*/
3835
abstract public function stop(): self;
3936

40-
/**
41-
* Receive and process messages one at a time until stopped.
42-
*/
4337
public function consume(callable $messageCallback, callable $successCallback, callable $errorCallback): void
4438
{
4539
$this->stopped = false;
@@ -56,10 +50,8 @@ public function consume(callable $messageCallback, callable $successCallback, ca
5650
}
5751

5852
/**
59-
* Run the handler for one message, then commit or reject it. Never throws:
60-
* any failure — including a failing commit/reject or error callback — is
61-
* routed to $errorCallback so it can't escape and be lost (e.g. swallowed
62-
* by a coroutine's default handler).
53+
* Never throws: a failing handler, commit, reject, or error callback is all
54+
* routed to $errorCallback so nothing escapes (and is lost) on a coroutine.
6355
*/
6456
protected function process(Message $message, callable $messageCallback, callable $successCallback, callable $errorCallback): void
6557
{
@@ -78,12 +70,11 @@ protected function process(Message $message, callable $messageCallback, callable
7870
try {
7971
$errorCallback($message, $error);
8072
} catch (\Throwable) {
81-
// Nothing left to do — the error callback itself failed.
73+
// the error callback itself failed; nothing left to do
8274
}
8375
}
8476
}
8577

86-
/** Install the per-message context container. */
8778
protected function setContext(Container $context): void
8879
{
8980
$this->context = $context;

‎src/Queue/Adapter/Swoole.php‎

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,6 @@ class Swoole extends Adapter
2323
/** @var callable[] */
2424
protected array $onWorkerStop = [];
2525

26-
/** Messages a worker may process concurrently. */
2726
protected int $maxCoroutines;
2827

2928
public function __construct(
@@ -85,9 +84,9 @@ protected function spawnWorker(int $workerId): void
8584
}
8685

8786
/**
88-
* Receive on a single loop and process each message on its own coroutine,
89-
* at most $maxCoroutines at a time. The channel is a semaphore: push()
90-
* blocks the loop once the pool is full until a handler frees a slot.
87+
* Receive on one loop, process each message on its own coroutine. The
88+
* channel caps concurrency at $maxCoroutines: push() blocks the loop while
89+
* the pool is full.
9190
*/
9291
public function consume(callable $messageCallback, callable $successCallback, callable $errorCallback): void
9392
{
@@ -109,8 +108,7 @@ public function consume(callable $messageCallback, callable $successCallback, ca
109108
try {
110109
$this->process($message, $messageCallback, $successCallback, $errorCallback);
111110
} catch (\Throwable $error) {
112-
// process() is total; last-resort net so a stray throw is
113-
// logged, not swallowed by Swoole's default handler.
111+
// process() is total; net for a stray throw so it isn't lost
114112
\error_log('Uncaught error while processing queue message: ' . $error->getMessage());
115113
} finally {
116114
$waitGroup->done();
@@ -119,13 +117,12 @@ public function consume(callable $messageCallback, callable $successCallback, ca
119117
});
120118
}
121119

122-
// Let in-flight handlers finish before returning.
123120
$waitGroup->wait();
124121
}
125122

126-
/** Keep the per-message container coroutine-local so handlers don't share it. */
127123
protected function setContext(Container $context): void
128124
{
125+
// coroutine-local so concurrent handlers don't share a context
129126
Coroutine::getContext()[self::CONTEXT_KEY] = $context;
130127
}
131128

‎src/Queue/Broker/Redis.php‎

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ class Redis implements Publisher, Consumer
1818
private readonly Connection $receive;
1919

2020
/** Carries acks and publishing; wrap in Locking when shared by coroutines. */
21-
private readonly Connection $work;
21+
private readonly Connection $commands;
2222

2323
private bool $closed = false;
2424
private int $reconnectAttempt = 0;
@@ -32,16 +32,12 @@ class Redis implements Publisher, Consumer
3232
*/
3333
private $reconnectSuccessCallback = null;
3434

35-
/**
36-
* @param Connection|null $work Defaults to $receive; pass a separate, locked
37-
* connection when processing concurrently.
38-
*/
3935
public function __construct(
4036
Connection $receive,
41-
?Connection $work = null,
37+
?Connection $commands = null,
4238
) {
4339
$this->receive = $receive;
44-
$this->work = $work ?? $receive;
40+
$this->commands = $commands ?? $receive;
4541
}
4642

4743
public function setReconnectCallback(?callable $callback): self
@@ -115,20 +111,20 @@ public function commit(Queue $queue, Message $message): void
115111
{
116112
$pid = $message->getPid();
117113

118-
$this->work->remove("{$queue->namespace}.jobs.{$queue->name}.{$pid}");
119-
$this->work->increment("{$queue->namespace}.stats.{$queue->name}.success");
120-
$this->work->listRemove("{$queue->namespace}.processing.{$queue->name}", $pid);
121-
$this->work->decrement("{$queue->namespace}.stats.{$queue->name}.processing");
114+
$this->commands->remove("{$queue->namespace}.jobs.{$queue->name}.{$pid}");
115+
$this->commands->increment("{$queue->namespace}.stats.{$queue->name}.success");
116+
$this->commands->listRemove("{$queue->namespace}.processing.{$queue->name}", $pid);
117+
$this->commands->decrement("{$queue->namespace}.stats.{$queue->name}.processing");
122118
}
123119

124120
public function reject(Queue $queue, Message $message): void
125121
{
126122
$pid = $message->getPid();
127123

128-
$this->work->leftPush("{$queue->namespace}.failed.{$queue->name}", $pid);
129-
$this->work->increment("{$queue->namespace}.stats.{$queue->name}.failed");
130-
$this->work->listRemove("{$queue->namespace}.processing.{$queue->name}", $pid);
131-
$this->work->decrement("{$queue->namespace}.stats.{$queue->name}.processing");
124+
$this->commands->leftPush("{$queue->namespace}.failed.{$queue->name}", $pid);
125+
$this->commands->increment("{$queue->namespace}.stats.{$queue->name}.failed");
126+
$this->commands->listRemove("{$queue->namespace}.processing.{$queue->name}", $pid);
127+
$this->commands->decrement("{$queue->namespace}.stats.{$queue->name}.processing");
132128
}
133129

134130
public function close(): void
@@ -169,9 +165,9 @@ public function enqueue(Queue $queue, array $payload, bool $priority = false): b
169165
'payload' => $payload
170166
];
171167
if ($priority) {
172-
return $this->work->rightPushArray("{$queue->namespace}.queue.{$queue->name}", $payload);
168+
return $this->commands->rightPushArray("{$queue->namespace}.queue.{$queue->name}", $payload);
173169
}
174-
return $this->work->leftPushArray("{$queue->namespace}.queue.{$queue->name}", $payload);
170+
return $this->commands->leftPushArray("{$queue->namespace}.queue.{$queue->name}", $payload);
175171
}
176172

177173
/**
@@ -184,7 +180,7 @@ public function retry(Queue $queue, ?int $limit = null): void
184180
$processed = 0;
185181

186182
while (true) {
187-
$pid = $this->work->rightPop("{$queue->namespace}.failed.{$queue->name}", self::POP_TIMEOUT);
183+
$pid = $this->commands->rightPop("{$queue->namespace}.failed.{$queue->name}", self::POP_TIMEOUT);
188184

189185
// No more jobs to retry
190186
if ($pid === false) {
@@ -215,7 +211,7 @@ public function retry(Queue $queue, ?int $limit = null): void
215211

216212
private function getJob(Queue $queue, string $pid): Message|false
217213
{
218-
$value = $this->work->get("{$queue->namespace}.jobs.{$queue->name}.{$pid}");
214+
$value = $this->commands->get("{$queue->namespace}.jobs.{$queue->name}.{$pid}");
219215

220216
// Missing/expired jobs return false or null depending on the driver.
221217
if (!\is_string($value)) {
@@ -233,6 +229,6 @@ public function getQueueSize(Queue $queue, bool $failedJobs = false): int
233229
if ($failedJobs) {
234230
$queueName = "{$queue->namespace}.failed.{$queue->name}";
235231
}
236-
return $this->work->listSize($queueName);
232+
return $this->commands->listSize($queueName);
237233
}
238234
}

‎src/Queue/Consumer.php‎

Lines changed: 4 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -4,24 +4,15 @@
44

55
interface Consumer
66
{
7-
/**
8-
* Block up to $timeout seconds for the next message and claim it. Returns
9-
* null on timeout so the caller can re-check its state.
10-
*/
7+
/** Block up to $timeout seconds for the next message and claim it, or null on timeout. */
118
public function receive(Queue $queue, int $timeout): ?Message;
129

13-
/**
14-
* Acknowledge a message as successfully processed.
15-
*/
10+
/** Acknowledge a processed message. */
1611
public function commit(Queue $queue, Message $message): void;
1712

18-
/**
19-
* Mark a message as failed.
20-
*/
13+
/** Mark a message as failed. */
2114
public function reject(Queue $queue, Message $message): void;
2215

23-
/**
24-
* Close the consumer and free any underlying resources.
25-
*/
16+
/** Close the consumer and free resources. */
2617
public function close(): void;
2718
}

‎src/Queue/Server.php‎

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -434,10 +434,8 @@ protected function getArguments(Container $context, Hook $hook, array $payload =
434434
);
435435
}
436436

437-
// Params and injections are collected in two passes, so $arguments ends
438-
// up keyed by declared order but not necessarily iterated in it. Sort by
439-
// key: call_user_func_array passes integer-keyed values positionally in
440-
// iteration order, so an unordered array would mis-assign arguments.
437+
// call_user_func_array passes integer keys in iteration order, not key
438+
// order, so sort the two-pass (params, then injections) array by key.
441439
\ksort($arguments);
442440

443441
return $arguments;

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,10 +10,10 @@
1010
use Utopia\Queue\Broker\Redis;
1111
use Utopia\Validator\Text;
1212

13-
// Separate receive (blocking pop) and work (locked, shared by coroutines) connections.
13+
// Dedicated blocking-receive connection; a separate locked connection for commands.
1414
$consumer = new Redis(
1515
receive: new RedisConnection('redis'),
16-
work: new Locking(new RedisConnection('redis')),
16+
commands: new Locking(new RedisConnection('redis')),
1717
);
1818
$adapter = new Swoole($consumer, 12, 'swoole', maxCoroutines: 5);
1919
$server = new Server($adapter);

0 commit comments

Comments
 (0)