From d6b032c241ef2038047cd3dd032a85a3d64bd536 Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Mon, 27 Jul 2026 20:01:16 +0100 Subject: [PATCH] fix(queue): keep the consumer alive when the broker fails MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit consume() called $this->consumer->receive() unguarded, so anything the broker threw unwound out of the loop and ended the worker. A single unreachable moment cost the process, and the message currently in flight went with it. Nothing at that level retried. What looked like resilience came from the connection pool underneath: its reconnect attempts delayed each failed receive by tens of seconds, which acted as an accidental backoff in front of the missing guard. Pools 2.0 removes those retries, so the gap becomes a crash loop rather than a slow stall — but the guard was missing either way, and a consumer should survive a broker outage regardless of what its pool happens to do. Route both consume loops through nextMessage(), which reports the failure and pauses for RECEIVE_BACKOFF before the next attempt. The loop then carries on, so the worker outlives the outage and resumes when the broker returns. Reporting goes to the existing $errorCallback with a null message rather than to error_log. Its signature is already nullable and Server already guards for null before setting the message on the context, so this is the case that shape was for: a consumer logs a broker failure through the same Server::error() hook it already registers for handler failures, instead of the library deciding on its behalf. The callback is invoked defensively — a reporting hook that throws must not cost the worker either, which is the whole point of the change. One helper on the base class rather than a try/catch in each loop, since Swoole overrides consume() and the Workerman path inherits it. The consume() callbacks are now documented, including that $errorCallback can receive a null message. The added test drives the Swoole loop with a consumer that fails twice and then delegates to a working broker, asserting both failures reach the callback with a null message. Before this change the test dies with "Uncaught RuntimeException: broker unreachable"; after it, both failures are absorbed and the queued message is processed. Co-Authored-By: Claude Opus 5 (1M context) --- packages/queue/phpunit.xml | 1 + packages/queue/src/Queue/Adapter.php | 43 +++++++- packages/queue/src/Queue/Adapter/Swoole.php | 2 +- .../E2E/Adapter/ConsumerResilienceTest.php | 101 ++++++++++++++++++ 4 files changed, 144 insertions(+), 3 deletions(-) create mode 100644 packages/queue/tests/Queue/E2E/Adapter/ConsumerResilienceTest.php diff --git a/packages/queue/phpunit.xml b/packages/queue/phpunit.xml index 64067c8e4..fc19630a3 100644 --- a/packages/queue/phpunit.xml +++ b/packages/queue/phpunit.xml @@ -6,6 +6,7 @@ > + ./tests/Queue/E2E/Adapter/ConsumerResilienceTest.php ./tests/Queue/E2E/Adapter/LockingTest.php ./tests/Queue/E2E/Adapter/RedisReconnectCallbackTest.php ./tests/Queue/E2E/Adapter/ServerTelemetryTest.php diff --git a/packages/queue/src/Queue/Adapter.php b/packages/queue/src/Queue/Adapter.php index 6989bced0..ce3ad8e7e 100644 --- a/packages/queue/src/Queue/Adapter.php +++ b/packages/queue/src/Queue/Adapter.php @@ -8,6 +8,12 @@ abstract class Adapter { protected const int RECEIVE_TIMEOUT = 2; + /** + * Pause before asking again after the broker failed to answer, so an + * unreachable broker is retried at a steady rate rather than in a tight loop. + */ + protected const int RECEIVE_BACKOFF = 1; + public Queue $queue; protected ?Container $context = null; protected bool $stopped = false; @@ -38,14 +44,20 @@ protected function isStopped(): bool return $this->stopped; } + /** + * @param callable(Message): void $messageCallback + * @param callable(Message): void $successCallback + * @param callable(?Message, \Throwable): void $errorCallback Receives null when + * the failure was in obtaining a message rather than handling one. + */ public function consume(callable $messageCallback, callable $successCallback, callable $errorCallback): void { $this->stopped = false; while (!$this->isStopped()) { - $message = $this->consumer->receive($this->queue, static::RECEIVE_TIMEOUT); + $message = $this->nextMessage($errorCallback); - if (!$message instanceof \Utopia\Queue\Message) { + if (!$message instanceof Message) { continue; } @@ -54,6 +66,33 @@ public function consume(callable $messageCallback, callable $successCallback, ca } } + /** + * Never throws: a broker that cannot be reached is reported to + * $errorCallback and retried after RECEIVE_BACKOFF. Losing the worker to a + * transient outage is worse than waiting for the broker to come back. + * + * $errorCallback takes a nullable message for exactly this case — the + * failure is in obtaining one, so there is none to report alongside it. + * + * @param callable(?Message, \Throwable): void $errorCallback + */ + protected function nextMessage(callable $errorCallback): ?Message + { + try { + return $this->consumer->receive($this->queue, static::RECEIVE_TIMEOUT); + } catch (\Throwable $error) { + // A reporting hook that throws must not cost the worker either. + try { + $errorCallback(null, $error); + } catch (\Throwable) { + } + + sleep(static::RECEIVE_BACKOFF); + + return null; + } + } + /** * Never throws: a failed handler is rejected and reported to $errorCallback; * a failing reject or callback is swallowed rather than left to escape (and diff --git a/packages/queue/src/Queue/Adapter/Swoole.php b/packages/queue/src/Queue/Adapter/Swoole.php index 9bf90bad4..f7eac7761 100644 --- a/packages/queue/src/Queue/Adapter/Swoole.php +++ b/packages/queue/src/Queue/Adapter/Swoole.php @@ -96,7 +96,7 @@ public function consume(callable $messageCallback, callable $successCallback, ca $waitGroup = new WaitGroup(); while (!$this->isStopped()) { - $message = $this->consumer->receive($this->queue, static::RECEIVE_TIMEOUT); + $message = $this->nextMessage($errorCallback); if (!$message instanceof \Utopia\Queue\Message) { continue; diff --git a/packages/queue/tests/Queue/E2E/Adapter/ConsumerResilienceTest.php b/packages/queue/tests/Queue/E2E/Adapter/ConsumerResilienceTest.php new file mode 100644 index 000000000..c61616ce4 --- /dev/null +++ b/packages/queue/tests/Queue/E2E/Adapter/ConsumerResilienceTest.php @@ -0,0 +1,101 @@ +failures < 2) { + ++$this->failures; + + throw new \RuntimeException('broker unreachable'); + } + + return $this->inner->receive($queue, $timeout); + } + + public function commit(Queue $queue, Message $message): void + { + $this->inner->commit($queue, $message); + } + + public function reject(Queue $queue, Message $message): void + { + $this->inner->reject($queue, $message); + } + + public function close(): void + { + $this->inner->close(); + } + }; + + $processed = 0; + /** @var list $reported */ + $reported = []; + $reportedMessages = []; + + \Swoole\Coroutine\run(function () use ($broker, $flaky, $queue, &$processed, &$reported, &$reportedMessages): void { + $broker->enqueue($queue, ['n' => 1]); + + $adapter = new class ($flaky, 1, self::QUEUE, self::NAMESPACE) extends Swoole { + // Keep the test quick; the production pause is RECEIVE_BACKOFF seconds. + protected const int RECEIVE_BACKOFF = 0; + }; + + $adapter->consume( + function () use ($adapter, &$processed): void { + ++$processed; + $adapter->stop(); + }, + fn(): null => null, + function (?Message $message, \Throwable $error) use (&$reported, &$reportedMessages): void { + $reported[] = $error->getMessage(); + $reportedMessages[] = $message; + }, + ); + }); + + $this->assertSame(2, $flaky->failures, 'both failures were absorbed rather than escaping'); + $this->assertSame(1, $processed, 'the loop resumed and drained the queue'); + $this->assertSame(['broker unreachable', 'broker unreachable'], $reported, 'each failure was reported'); + $this->assertSame([null, null], $reportedMessages, 'reported without a message, since none was obtained'); + } +}