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'); + } +}