Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions packages/queue/phpunit.xml
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
>
<testsuites>
<testsuite name="unit">
<file>./tests/Queue/E2E/Adapter/ConsumerResilienceTest.php</file>
<file>./tests/Queue/E2E/Adapter/LockingTest.php</file>
<file>./tests/Queue/E2E/Adapter/RedisReconnectCallbackTest.php</file>
<file>./tests/Queue/E2E/Adapter/ServerTelemetryTest.php</file>
Expand Down
43 changes: 41 additions & 2 deletions packages/queue/src/Queue/Adapter.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}

Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion packages/queue/src/Queue/Adapter/Swoole.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
101 changes: 101 additions & 0 deletions packages/queue/tests/Queue/E2E/Adapter/ConsumerResilienceTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
<?php

declare(strict_types=1);

namespace Tests\E2E\Adapter;

use PHPUnit\Framework\TestCase;
use Utopia\Queue\Adapter\Swoole;
use Utopia\Queue\Broker\Redis;
use Utopia\Queue\Consumer;
use Utopia\Queue\Message;
use Utopia\Queue\Queue;

/**
* A broker outage must not cost the worker. `receive()` was called unguarded
* from the consume loop, so anything it threw unwound out of consume() and ended
* the process. Nothing retried at that level: the connection pool's internal
* reconnect attempts happened to delay each failure by tens of seconds, which
* looked like backoff but was only hiding the missing guard.
*
* The failure is reported through $errorCallback with a null message, which is
* what its nullable signature is for, so a consumer logs it through the same
* Server::error() hook it already uses for handler failures.
*/
final class ConsumerResilienceTest extends TestCase
{
private const string QUEUE = 'resilience';

private const string NAMESPACE = 'tests';

public function testConsumeSurvivesBrokerFailuresAndResumes(): void
{
$connection = new InMemoryConnection();
$broker = new Redis($connection, $connection);
$queue = new Queue(self::QUEUE, self::NAMESPACE);

// Fails the first two receives, then delegates to the working broker.
$flaky = new class ($broker) implements Consumer {
public int $failures = 0;

public function __construct(private readonly Redis $inner) {}

public function receive(Queue $queue, int $timeout): ?Message
{
if ($this->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<string> $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');
}
}
Loading