From f09e06ceb4242d1fbbb613ae1570cc5af5451248 Mon Sep 17 00:00:00 2001 From: Jake Barnby Date: Thu, 16 Jul 2026 00:41:09 +1200 Subject: [PATCH 1/4] (feat): add reliable Redis queue claims --- packages/queue/docker-compose.yml | 48 ++- packages/queue/phpunit.xml | 3 + packages/queue/src/Queue/Adapter.php | 4 +- packages/queue/src/Queue/Adapter/Swoole.php | 123 +++++- packages/queue/src/Queue/Broker/Redis.php | 365 +++++++++++++++++- .../queue/src/Queue/Broker/Redis/Script.php | 274 +++++++++++++ packages/queue/src/Queue/Claim.php | 13 + .../queue/src/Queue/Connection/Atomic.php | 18 + .../queue/src/Queue/Connection/Locking.php | 18 +- packages/queue/src/Queue/Connection/Redis.php | 12 +- .../queue/src/Queue/Consumer/Recoverable.php | 19 + packages/queue/src/Queue/Message.php | 7 +- packages/queue/src/Queue/Option/Reliable.php | 33 ++ packages/queue/src/Queue/Queue.php | 6 + .../Queue/E2E/Adapter/ReliableApiTest.php | 301 +++++++++++++++ .../Queue/E2E/Adapter/ReliableRedisTest.php | 357 +++++++++++++++++ .../Queue/E2E/Adapter/ReliableSwooleTest.php | 253 ++++++++++++ 17 files changed, 1840 insertions(+), 14 deletions(-) create mode 100644 packages/queue/src/Queue/Broker/Redis/Script.php create mode 100644 packages/queue/src/Queue/Claim.php create mode 100644 packages/queue/src/Queue/Connection/Atomic.php create mode 100644 packages/queue/src/Queue/Consumer/Recoverable.php create mode 100644 packages/queue/src/Queue/Option/Reliable.php create mode 100644 packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php create mode 100644 packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php create mode 100644 packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php diff --git a/packages/queue/docker-compose.yml b/packages/queue/docker-compose.yml index 6fc31d238..219257db2 100644 --- a/packages/queue/docker-compose.yml +++ b/packages/queue/docker-compose.yml @@ -9,13 +9,49 @@ services: timeout: 3s retries: 15 + dragonfly: + image: docker.dragonflydb.io/dragonflydb/dragonfly:v1.39.0 + ports: + - "16380:6379" + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 2s + timeout: 3s + retries: 15 + redis-cluster: - image: grokzen/redis-cluster:7.0.10 - environment: - IP: "0.0.0.0" - INITIAL_PORT: 17000 - MASTERS: 3 - SLAVES_PER_MASTER: 0 + image: redis:alpine + command: + - /bin/sh + - -c + - | + for port in 17000 17001 17002; do + mkdir -p "/tmp/redis-$$port" + redis-server \ + --port "$$port" \ + --bind 0.0.0.0 \ + --protected-mode no \ + --cluster-enabled yes \ + --cluster-config-file "nodes-$$port.conf" \ + --cluster-node-timeout 5000 \ + --cluster-announce-ip 127.0.0.1 \ + --cluster-announce-port "$$port" \ + --dir "/tmp/redis-$$port" \ + --appendonly no \ + --daemonize yes + done + for port in 17000 17001 17002; do + until redis-cli -p "$$port" ping >/dev/null 2>&1; do + sleep 0.1 + done + done + redis-cli --cluster create \ + 127.0.0.1:17000 \ + 127.0.0.1:17001 \ + 127.0.0.1:17002 \ + --cluster-replicas 0 \ + --cluster-yes + tail -f /dev/null ports: - "17000-17002:17000-17002" healthcheck: diff --git a/packages/queue/phpunit.xml b/packages/queue/phpunit.xml index 64067c8e4..1d2aca7cb 100644 --- a/packages/queue/phpunit.xml +++ b/packages/queue/phpunit.xml @@ -7,11 +7,14 @@ ./tests/Queue/E2E/Adapter/LockingTest.php + ./tests/Queue/E2E/Adapter/ReliableApiTest.php ./tests/Queue/E2E/Adapter/RedisReconnectCallbackTest.php ./tests/Queue/E2E/Adapter/ServerTelemetryTest.php ./tests/Queue/E2E/Adapter/SwooleConcurrencyTest.php + ./tests/Queue/E2E/Adapter/ReliableRedisTest.php + ./tests/Queue/E2E/Adapter/ReliableSwooleTest.php ./tests/Queue/E2E/Adapter/PoolTest.php ./tests/Queue/E2E/Adapter/SwooleTest.php ./tests/Queue/E2E/Adapter/SwooleRedisClusterTest.php diff --git a/packages/queue/src/Queue/Adapter.php b/packages/queue/src/Queue/Adapter.php index 6989bced0..9778392a5 100644 --- a/packages/queue/src/Queue/Adapter.php +++ b/packages/queue/src/Queue/Adapter.php @@ -3,6 +3,7 @@ namespace Utopia\Queue; use Utopia\DI\Container; +use Utopia\Queue\Option\Reliable; abstract class Adapter { @@ -18,8 +19,9 @@ public function __construct( string $queue, public string $namespace = 'utopia-queue', protected Container $resources = new Container(), + ?Reliable $reliable = null, ) { - $this->queue = new Queue($queue, $namespace); + $this->queue = new Queue($queue, $namespace, reliable: $reliable); } /** diff --git a/packages/queue/src/Queue/Adapter/Swoole.php b/packages/queue/src/Queue/Adapter/Swoole.php index 9bf90bad4..42ead27b6 100644 --- a/packages/queue/src/Queue/Adapter/Swoole.php +++ b/packages/queue/src/Queue/Adapter/Swoole.php @@ -9,6 +9,9 @@ use Utopia\DI\Container; use Utopia\Queue\Adapter; use Utopia\Queue\Consumer; +use Utopia\Queue\Consumer\Recoverable; +use Utopia\Queue\Message; +use Utopia\Queue\Option\Reliable; class Swoole extends Adapter { @@ -32,8 +35,9 @@ public function __construct( string $namespace = 'utopia-queue', int $maxCoroutines = 1, Container $resources = new Container(), + ?Reliable $reliable = null, ) { - parent::__construct($consumer, $workerNum, $queue, $namespace, $resources); + parent::__construct($consumer, $workerNum, $queue, $namespace, $resources, $reliable); $this->maxCoroutines = max(1, $maxCoroutines); } @@ -90,6 +94,21 @@ protected function spawnWorker(int $workerId): void */ #[\Override] public function consume(callable $messageCallback, callable $successCallback, callable $errorCallback): void + { + if ($this->queue->reliable instanceof Reliable) { + if (!$this->consumer instanceof Recoverable) { + throw new \LogicException('Reliable Swoole queues require a recoverable consumer.'); + } + + $this->consumeReliable($messageCallback, $successCallback, $errorCallback, $this->consumer); + + return; + } + + $this->consumeLegacy($messageCallback, $successCallback, $errorCallback); + } + + private function consumeLegacy(callable $messageCallback, callable $successCallback, callable $errorCallback): void { $this->stopped = false; $slots = new Channel($this->maxCoroutines); @@ -121,6 +140,108 @@ public function consume(callable $messageCallback, callable $successCallback, ca $waitGroup->wait(); } + private function consumeReliable( + callable $messageCallback, + callable $successCallback, + callable $errorCallback, + Recoverable $recoverable, + ): void { + $reliable = $this->queue->reliable + ?? throw new \LogicException('Reliable configuration is missing.'); + $this->stopped = false; + $slots = new Channel($this->maxCoroutines); + $handlers = new WaitGroup(); + $recoveryDone = new Channel(1); + $recovery = new WaitGroup(1); + + Coroutine::create(function () use ($recoverable, $reliable, $recoveryDone, $recovery): void { + try { + while ($recoveryDone->pop($reliable->scan) === false) { + try { + do { + $claims = $recoverable->expired($this->queue, $reliable->batch); + foreach ($claims as $claim) { + $recoverable->reclaim($this->queue, $claim); + } + if (\count($claims) === $reliable->batch) { + Coroutine::sleep(0.001); + } + } while (\count($claims) === $reliable->batch && !$this->isStopped()); + } catch (\Throwable $error) { + error_log('Queue recovery failed: ' . $error->getMessage()); + } + } + } finally { + $recovery->done(); + } + }); + + try { + while (!$this->isStopped()) { + $slots->push(true); + if ($this->isStopped()) { + $slots->pop(); + break; + } + + try { + $message = $this->consumer->receive($this->queue, static::RECEIVE_TIMEOUT); + } catch (\Throwable $error) { + $slots->pop(); + throw $error; + } + if (!$message instanceof Message) { + $slots->pop(); + continue; + } + + $handlers->add(); + Coroutine::create(function () use ( + $message, + $messageCallback, + $successCallback, + $errorCallback, + $recoverable, + $reliable, + $slots, + $handlers, + ): void { + $heartbeatDone = new Channel(1); + $heartbeat = new WaitGroup(1); + Coroutine::create(function () use ($message, $recoverable, $reliable, $heartbeatDone, $heartbeat): void { + try { + while ($heartbeatDone->pop($reliable->heartbeat) === false) { + if (!$recoverable->extend($this->queue, $message)) { + error_log("Queue lease was lost for message {$message->getPid()}."); + return; + } + } + } catch (\Throwable $error) { + error_log("Queue heartbeat failed for message {$message->getPid()}: {$error->getMessage()}"); + } finally { + $heartbeat->done(); + } + }); + + try { + $this->process($message, $messageCallback, $successCallback, $errorCallback); + } catch (\Throwable $error) { + error_log('Uncaught error while processing queue message: ' . $error->getMessage()); + } finally { + $heartbeatDone->push(true); + $heartbeat->wait(); + $handlers->done(); + $slots->pop(); + } + }); + } + } finally { + $handlers->wait(); + $recoveryDone->push(true); + $recovery->wait(); + } + } + #[\Override] public function context(): Container { diff --git a/packages/queue/src/Queue/Broker/Redis.php b/packages/queue/src/Queue/Broker/Redis.php index 65d38535b..0afba5191 100644 --- a/packages/queue/src/Queue/Broker/Redis.php +++ b/packages/queue/src/Queue/Broker/Redis.php @@ -2,13 +2,18 @@ namespace Utopia\Queue\Broker; +use Utopia\Queue\Broker\Redis\Script; +use Utopia\Queue\Claim; use Utopia\Queue\Connection; +use Utopia\Queue\Connection\Atomic; use Utopia\Queue\Consumer; +use Utopia\Queue\Consumer\Recoverable; use Utopia\Queue\Message; +use Utopia\Queue\Option\Reliable; use Utopia\Queue\Publisher; use Utopia\Queue\Queue; -class Redis implements Publisher, Consumer +class Redis implements Publisher, Consumer, Recoverable { private const int POP_TIMEOUT = 2; private const int RECONNECT_BACKOFF_MS = 100; @@ -48,6 +53,15 @@ public function setReconnectSuccessCallback(?callable $callback): self } public function receive(Queue $queue, int $timeout): ?Message + { + if ($queue->reliable instanceof Reliable) { + return $this->receiveReliable($queue, $timeout); + } + + return $this->receiveLegacy($queue, $timeout); + } + + private function receiveLegacy(Queue $queue, int $timeout): ?Message { if ($this->isClosed()) { return null; @@ -102,6 +116,12 @@ public function receive(Queue $queue, int $timeout): ?Message public function commit(Queue $queue, Message $message): void { + if ($queue->reliable instanceof Reliable) { + $this->commitReliable($queue, $message); + + return; + } + $pid = $message->getPid(); $this->commands->remove("{$queue->namespace}.jobs.{$queue->name}.{$pid}"); @@ -112,6 +132,12 @@ public function commit(Queue $queue, Message $message): void public function reject(Queue $queue, Message $message): void { + if ($queue->reliable instanceof Reliable) { + $this->rejectReliable($queue, $message); + + return; + } + $pid = $message->getPid(); $this->commands->leftPush("{$queue->namespace}.failed.{$queue->name}", $pid); @@ -157,6 +183,23 @@ private function triggerReconnectSuccessCallback(Queue $queue, int $attempts): v public function enqueue(Queue $queue, array $payload, bool $priority = false): bool { + if ($queue->reliable instanceof Reliable) { + $atomic = $this->atomic($this->commands); + $result = $atomic->evaluate( + Script::ENQUEUE, + [ + $this->pendingKey($queue), + $this->pid(), + $queue->name, + json_encode($payload, JSON_THROW_ON_ERROR), + $priority ? '1' : '0', + ], + 1, + ); + + return (int) $result === 1; + } + $payload = [ 'pid' => uniqid(more_entropy: true), 'queue' => $queue->name, @@ -175,6 +218,12 @@ public function enqueue(Queue $queue, array $payload, bool $priority = false): b */ public function retry(Queue $queue, ?int $limit = null): void { + if ($queue->reliable instanceof Reliable) { + $this->retryReliable($queue, $limit); + + return; + } + $start = time(); $processed = 0; @@ -224,10 +273,320 @@ private function getJob(Queue $queue, string $pid): Message|false public function getQueueSize(Queue $queue, bool $failedJobs = false): int { - $queueName = "{$queue->namespace}.queue.{$queue->name}"; + $queueName = $this->pendingKey($queue); if ($failedJobs) { - $queueName = "{$queue->namespace}.failed.{$queue->name}"; + if ($queue->reliable instanceof Reliable) { + $this->atomic($this->commands); + $queueName = $this->atomicKey($queue, 'failed'); + } else { + $queueName = "{$queue->namespace}.failed.{$queue->name}"; + } } return $this->commands->listSize($queueName); } + + public function extend(Queue $queue, Message $message): bool + { + $reliable = $this->reliable($queue); + $claimedAt = $this->claimedAt($message); + $result = $this->atomic($this->commands)->evaluate( + Script::EXTEND, + [ + $this->atomicKey($queue, 'jobs'), + $this->atomicKey($queue, 'processing'), + $message->getPid(), + $claimedAt, + (string) $reliable->visibility, + ], + 2, + ); + + return (int) $result === 1; + } + + public function expired(Queue $queue, int $limit): array + { + $this->reliable($queue); + if ($limit <= 0) { + return []; + } + + $result = $this->atomic($this->commands)->evaluate( + Script::EXPIRED, + [ + $this->atomicKey($queue, 'processing'), + $this->atomicKey($queue, 'jobs'), + (string) $limit, + ], + 2, + ); + if (!\is_array($result)) { + throw new \UnexpectedValueException('Reliable expired scan returned an invalid response.'); + } + + $claims = []; + $counter = \count($result); + for ($index = 0; $index + 1 < $counter; $index += 2) { + $pid = $result[$index]; + $claimedAt = $result[$index + 1]; + if (!\is_string($pid) || !\is_string($claimedAt)) { + throw new \UnexpectedValueException('Reliable expired scan returned an invalid claim.'); + } + $claims[] = new Claim($pid, $claimedAt === '' ? null : $claimedAt); + } + + return $claims; + } + + public function reclaim(Queue $queue, Claim $claim): ?Message + { + $this->reliable($queue); + $result = $this->atomic($this->commands)->evaluate( + Script::RECLAIM, + [ + $this->atomicKey($queue, 'processing'), + $this->atomicKey($queue, 'jobs'), + $this->pendingKey($queue), + $this->atomicKey($queue, 'quarantine'), + $this->statKey($queue, 'processing'), + $this->statKey($queue, 'reclaimed'), + $this->statKey($queue, 'quarantined'), + $claim->pid, + $claim->claimedAt ?? '', + $this->pid(), + ], + 7, + ); + + if (!\is_array($result) || !isset($result[0])) { + throw new \UnexpectedValueException('Reliable reclaim returned an invalid response.'); + } + if ((int) $result[0] !== 1) { + return null; + } + if (!isset($result[1]) || !\is_string($result[1])) { + throw new \UnexpectedValueException('Reliable reclaim returned an invalid message.'); + } + + return $this->message($result[1]); + } + + private function receiveReliable(Queue $queue, int $timeout): ?Message + { + if ($this->isClosed()) { + return null; + } + + $reliable = $this->reliable($queue); + $atomic = $this->atomic($this->receive); + $deadline = hrtime(true) + (max(0, $timeout) * 1_000_000_000); + $backoff = $reliable->pollMinimum; + + while (!$this->isClosed()) { + try { + $result = $atomic->evaluate( + Script::CLAIM, + [ + $this->pendingKey($queue), + $this->atomicKey($queue, 'jobs'), + $this->atomicKey($queue, 'processing'), + $this->atomicKey($queue, 'quarantine'), + $this->statKey($queue, 'total'), + $this->statKey($queue, 'processing'), + $this->statKey($queue, 'quarantined'), + (string) $reliable->visibility, + ], + 7, + ); + $this->reconnectSucceeded($queue); + } catch (\RedisException|\RedisClusterException $error) { + $remaining = max(0, $deadline - hrtime(true)); + $this->reconnectFailed($queue, $error, intdiv($remaining, 1_000_000)); + + if ($this->isClosed() || $timeout === 0 || hrtime(true) >= $deadline) { + return null; + } + + continue; + } + + if (!\is_array($result) || !isset($result[0])) { + throw new \UnexpectedValueException('Reliable claim returned an invalid response.'); + } + if ((int) $result[0] === 1) { + if (!isset($result[1], $result[2]) || !\is_string($result[1]) || !\is_string($result[2])) { + throw new \UnexpectedValueException('Reliable claim returned an invalid message.'); + } + + return $this->message($result[1], $result[2]); + } + + $remaining = $deadline - hrtime(true); + if ($timeout === 0 || $remaining <= 0) { + return null; + } + + $maximum = min($reliable->pollMaximum, $backoff); + $sleep = mt_rand($reliable->pollMinimum, $maximum) * 1_000_000; + $sleep = min($sleep, $remaining); + if ($sleep > 0) { + usleep(max(1, intdiv($sleep, 1_000))); + } + $backoff = min($reliable->pollMaximum, $backoff * 2); + } + + return null; + } + + private function commitReliable(Queue $queue, Message $message): void + { + $this->reliable($queue); + $this->atomic($this->commands)->evaluate( + Script::COMMIT, + [ + $this->atomicKey($queue, 'jobs'), + $this->atomicKey($queue, 'processing'), + $this->statKey($queue, 'success'), + $this->statKey($queue, 'processing'), + $message->getPid(), + $this->claimedAt($message), + ], + 4, + ); + } + + private function rejectReliable(Queue $queue, Message $message): void + { + $this->reliable($queue); + $this->atomic($this->commands)->evaluate( + Script::REJECT, + [ + $this->atomicKey($queue, 'jobs'), + $this->atomicKey($queue, 'processing'), + $this->atomicKey($queue, 'failed'), + $this->statKey($queue, 'failed'), + $this->statKey($queue, 'processing'), + $message->getPid(), + $this->claimedAt($message), + ], + 5, + ); + } + + private function retryReliable(Queue $queue, ?int $limit): void + { + $this->reliable($queue); + $atomic = $this->atomic($this->commands); + $processed = 0; + + while ($limit === null || $processed < $limit) { + $result = $atomic->evaluate( + Script::RETRY, + [ + $this->atomicKey($queue, 'failed'), + $this->atomicKey($queue, 'jobs'), + $this->pendingKey($queue), + $this->atomicKey($queue, 'quarantine'), + $this->statKey($queue, 'retried'), + $this->statKey($queue, 'quarantined'), + $this->pid(), + ], + 6, + ); + if (!\is_array($result) || !isset($result[0])) { + throw new \UnexpectedValueException('Reliable retry returned an invalid response.'); + } + + $status = (int) $result[0]; + if ($status === 0) { + return; + } + if ($status === 1) { + $processed++; + } + } + } + + private function atomic(Connection $connection): Atomic + { + if (!$connection instanceof Atomic || !$connection->supportsAtomic()) { + throw new \LogicException('Reliable queues require a single Redis connection with atomic scripting support.'); + } + + return $connection; + } + + private function reliable(Queue $queue): Reliable + { + if (!$queue->reliable instanceof Reliable) { + throw new \LogicException('Recovery operations require a reliable queue.'); + } + + return $queue->reliable; + } + + private function claimedAt(Message $message): string + { + return $message->getClaimedAt() + ?? throw new \LogicException('Reliable messages require claim metadata.'); + } + + private function message(string $encoded, ?string $claimedAt = null): Message + { + $value = json_decode($encoded, true, flags: JSON_THROW_ON_ERROR); + if (!\is_array($value)) { + throw new \UnexpectedValueException('Reliable message envelope must be an array.'); + } + $value['timestamp'] = (int) ($value['timestamp'] ?? 0); + + return new Message($value, $claimedAt); + } + + private function pendingKey(Queue $queue): string + { + return "{$queue->namespace}.queue.{$queue->name}"; + } + + private function atomicKey(Queue $queue, string $type): string + { + return "{$queue->namespace}.atomic.{$type}.{$queue->name}"; + } + + private function statKey(Queue $queue, string $stat): string + { + return "{$queue->namespace}.stats.{$queue->name}.{$stat}"; + } + + private function pid(): string + { + return uniqid(more_entropy: true); + } + + private function reconnectSucceeded(Queue $queue): void + { + if ($this->reconnectAttempt > 0) { + $this->triggerReconnectSuccessCallback($queue, $this->reconnectAttempt); + } + + $this->reconnectBackoffMs = self::RECONNECT_BACKOFF_MS; + $this->reconnectAttempt = 0; + } + + private function reconnectFailed(Queue $queue, \Throwable $error, int $maximumSleepMs): void + { + if ($this->isClosed()) { + return; + } + + $this->reconnectAttempt++; + try { + $this->receive->close(); + } catch (\Throwable) { + } + + $sleepMs = mt_rand(0, min($this->reconnectBackoffMs, max(0, $maximumSleepMs))); + $this->triggerReconnectCallback($queue, $error, $this->reconnectAttempt, $sleepMs); + usleep($sleepMs * 1000); + $this->reconnectBackoffMs = min(self::RECONNECT_MAX_BACKOFF_MS, $this->reconnectBackoffMs * 2); + } } diff --git a/packages/queue/src/Queue/Broker/Redis/Script.php b/packages/queue/src/Queue/Broker/Redis/Script.php new file mode 100644 index 000000000..81fe81fe6 --- /dev/null +++ b/packages/queue/src/Queue/Broker/Redis/Script.php @@ -0,0 +1,274 @@ + 0 then + redis.call('DECR', KEYS[4]) +end +return 1 +LUA; + + public const string REJECT = <<<'LUA' +local raw = redis.call('HGET', KEYS[1], ARGV[1]) +if not raw or not redis.call('ZSCORE', KEYS[2], ARGV[1]) then + return 0 +end + +local valid, record = pcall(cjson.decode, raw) +if not valid + or type(record) ~= 'table' + or record.state ~= 'processing' + or record.claimedAt ~= ARGV[2] then + return 0 +end + +record.state = 'failed' +redis.call('HSET', KEYS[1], ARGV[1], cjson.encode(record)) +redis.call('ZREM', KEYS[2], ARGV[1]) +redis.call('LPUSH', KEYS[3], ARGV[1]) +redis.call('INCR', KEYS[4]) +local processing = tonumber(redis.call('GET', KEYS[5]) or '0') +if processing > 0 then + redis.call('DECR', KEYS[5]) +end +return 1 +LUA; + + public const string RETRY = <<<'LUA' +local pid = redis.call('RPOP', KEYS[1]) +if not pid then + return {0} +end + +local raw = redis.call('HGET', KEYS[2], pid) +local valid, record = pcall(cjson.decode, raw or '') +if not raw + or not valid + or type(record) ~= 'table' + or record.state ~= 'failed' + or type(record.message) ~= 'table' + or type(record.message.queue) ~= 'string' + or type(record.message.payload) ~= 'table' then + redis.call('HDEL', KEYS[2], pid) + redis.call('LPUSH', KEYS[4], cjson.encode({reason = 'malformed_failed', pid = pid, raw = raw or false})) + redis.call('INCR', KEYS[6]) + return {2} +end + +if redis.call('HEXISTS', KEYS[2], ARGV[1]) == 1 then + redis.call('RPUSH', KEYS[1], pid) + return {3} +end + +local now = redis.call('TIME') +local message = { + pid = ARGV[1], + queue = record.message.queue, + timestamp = tonumber(now[1]), + payload = record.message.payload +} +local encoded = '{"pid":' .. cjson.encode(message.pid) + .. ',"queue":' .. cjson.encode(message.queue) + .. ',"timestamp":' .. tostring(message.timestamp) + .. ',"payload":' .. cjson.encode(message.payload) .. '}' +redis.call('LPUSH', KEYS[3], encoded) +redis.call('HDEL', KEYS[2], pid) +redis.call('INCR', KEYS[5]) +return {1, encoded} +LUA; + + public const string EXPIRED = <<<'LUA' +local now = redis.call('TIME') +local micros = (tonumber(now[1]) * 1000000) + tonumber(now[2]) +local pids = redis.call('ZRANGEBYSCORE', KEYS[1], '-inf', micros, 'LIMIT', 0, tonumber(ARGV[1])) +local result = {} + +for _, pid in ipairs(pids) do + local raw = redis.call('HGET', KEYS[2], pid) + local valid, record = pcall(cjson.decode, raw or '') + local claimedAt = '' + if raw + and valid + and type(record) == 'table' + and record.state == 'processing' + and type(record.claimedAt) == 'string' then + claimedAt = record.claimedAt + end + table.insert(result, pid) + table.insert(result, claimedAt) +end + +return result +LUA; + + public const string RECLAIM = <<<'LUA' +local score = redis.call('ZSCORE', KEYS[1], ARGV[1]) +if not score then + return {0} +end + +local now = redis.call('TIME') +local micros = (tonumber(now[1]) * 1000000) + tonumber(now[2]) +if tonumber(score) > micros then + return {0} +end + +local raw = redis.call('HGET', KEYS[2], ARGV[1]) +local valid, record = pcall(cjson.decode, raw or '') +if not raw + or not valid + or type(record) ~= 'table' + or record.state ~= 'processing' + or type(record.message) ~= 'table' + or type(record.message.queue) ~= 'string' + or type(record.message.timestamp) ~= 'number' + or type(record.message.payload) ~= 'table' then + redis.call('ZREM', KEYS[1], ARGV[1]) + redis.call('HDEL', KEYS[2], ARGV[1]) + redis.call('LPUSH', KEYS[4], cjson.encode({reason = 'missing_or_malformed_processing', pid = ARGV[1], raw = raw or false})) + local processing = tonumber(redis.call('GET', KEYS[5]) or '0') + if processing > 0 then + redis.call('DECR', KEYS[5]) + end + redis.call('INCR', KEYS[7]) + return {2} +end + +if ARGV[2] == '' or record.claimedAt ~= ARGV[2] then + return {0} +end + +if redis.call('HEXISTS', KEYS[2], ARGV[3]) == 1 then + return {3} +end + +local message = { + pid = ARGV[3], + queue = record.message.queue, + timestamp = record.message.timestamp, + payload = record.message.payload +} +local encoded = '{"pid":' .. cjson.encode(message.pid) + .. ',"queue":' .. cjson.encode(message.queue) + .. ',"timestamp":' .. tostring(message.timestamp) + .. ',"payload":' .. cjson.encode(message.payload) .. '}' +redis.call('RPUSH', KEYS[3], encoded) +redis.call('HDEL', KEYS[2], ARGV[1]) +redis.call('ZREM', KEYS[1], ARGV[1]) +local processing = tonumber(redis.call('GET', KEYS[5]) or '0') +if processing > 0 then + redis.call('DECR', KEYS[5]) +end +redis.call('INCR', KEYS[6]) +return {1, encoded} +LUA; + + public const string EXTEND = <<<'LUA' +local raw = redis.call('HGET', KEYS[1], ARGV[1]) +if not raw or not redis.call('ZSCORE', KEYS[2], ARGV[1]) then + return 0 +end + +local valid, record = pcall(cjson.decode, raw) +if not valid + or type(record) ~= 'table' + or record.state ~= 'processing' + or record.claimedAt ~= ARGV[2] then + return 0 +end + +local now = redis.call('TIME') +local micros = (tonumber(now[1]) * 1000000) + tonumber(now[2]) +local leaseUntil = micros + (tonumber(ARGV[3]) * 1000000) +record.leaseUntil = leaseUntil +redis.call('HSET', KEYS[1], ARGV[1], cjson.encode(record)) +redis.call('ZADD', KEYS[2], 'XX', leaseUntil, ARGV[1]) +return 1 +LUA; +} diff --git a/packages/queue/src/Queue/Claim.php b/packages/queue/src/Queue/Claim.php new file mode 100644 index 000000000..a204c0581 --- /dev/null +++ b/packages/queue/src/Queue/Claim.php @@ -0,0 +1,13 @@ + $arguments + */ + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed; +} diff --git a/packages/queue/src/Queue/Connection/Locking.php b/packages/queue/src/Queue/Connection/Locking.php index 97aa55e8e..a9b3de478 100644 --- a/packages/queue/src/Queue/Connection/Locking.php +++ b/packages/queue/src/Queue/Connection/Locking.php @@ -14,7 +14,7 @@ * Outside of a coroutine there is no preemption, so the lock degrades to a * plain in-process flag (see {@see Mutex}). */ -class Locking implements Connection +class Locking implements Connection, Atomic { /** * Wait forever when acquiring the lock; a command should never be dropped @@ -139,6 +139,22 @@ public function ping(): bool return $this->synchronize(fn(): bool => $this->connection->ping()); } + public function supportsAtomic(): bool + { + return $this->connection instanceof Atomic && $this->connection->supportsAtomic(); + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + if (!$this->connection instanceof Atomic || !$this->connection->supportsAtomic()) { + throw new \LogicException('The wrapped connection does not support atomic scripting.'); + } + + return $this->synchronize( + fn(): mixed => $this->connection->evaluate($script, $arguments, $keyCount), + ); + } + public function close(): void { $this->synchronize(fn() => $this->connection->close()); diff --git a/packages/queue/src/Queue/Connection/Redis.php b/packages/queue/src/Queue/Connection/Redis.php index ffd1c6597..9f8f688d9 100644 --- a/packages/queue/src/Queue/Connection/Redis.php +++ b/packages/queue/src/Queue/Connection/Redis.php @@ -4,7 +4,7 @@ use Utopia\Queue\Connection; -class Redis implements Connection +class Redis implements Connection, Atomic { protected const int CONNECT_MAX_ATTEMPTS = 5; protected const int CONNECT_BACKOFF_MS = 100; @@ -160,6 +160,16 @@ public function ping(): bool } } + public function supportsAtomic(): bool + { + return true; + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + return $this->getRedis()->eval($script, $arguments, $keyCount); + } + public function close(): void { try { diff --git a/packages/queue/src/Queue/Consumer/Recoverable.php b/packages/queue/src/Queue/Consumer/Recoverable.php new file mode 100644 index 000000000..f645b645e --- /dev/null +++ b/packages/queue/src/Queue/Consumer/Recoverable.php @@ -0,0 +1,19 @@ + */ + public function expired(Queue $queue, int $limit): array; + + public function reclaim(Queue $queue, Claim $claim): ?Message; +} diff --git a/packages/queue/src/Queue/Message.php b/packages/queue/src/Queue/Message.php index 7e6d885ff..1048c2264 100644 --- a/packages/queue/src/Queue/Message.php +++ b/packages/queue/src/Queue/Message.php @@ -9,7 +9,7 @@ class Message protected int $timestamp; protected array $payload; - public function __construct(array $array = []) + public function __construct(array $array = [], private readonly ?string $claimedAt = null) { if ($array === []) { return; @@ -69,6 +69,11 @@ public function getPayload(): array return $this->payload; } + public function getClaimedAt(): ?string + { + return $this->claimedAt; + } + public function asArray(): array { return [ diff --git a/packages/queue/src/Queue/Option/Reliable.php b/packages/queue/src/Queue/Option/Reliable.php new file mode 100644 index 000000000..116db0f50 --- /dev/null +++ b/packages/queue/src/Queue/Option/Reliable.php @@ -0,0 +1,33 @@ +visibility <= 0) { + throw new \InvalidArgumentException('Visibility must be greater than zero.'); + } + if ($this->heartbeat <= 0 || $this->heartbeat >= $this->visibility) { + throw new \InvalidArgumentException('Heartbeat must be greater than zero and less than visibility.'); + } + if ($this->scan <= 0) { + throw new \InvalidArgumentException('Recovery scan interval must be greater than zero.'); + } + if ($this->batch <= 0) { + throw new \InvalidArgumentException('Recovery batch must be greater than zero.'); + } + if ($this->pollMinimum <= 0 || $this->pollMaximum < $this->pollMinimum) { + throw new \InvalidArgumentException('Poll bounds must be positive and ordered.'); + } + } +} diff --git a/packages/queue/src/Queue/Queue.php b/packages/queue/src/Queue/Queue.php index 73c7a8d90..74bee28ad 100644 --- a/packages/queue/src/Queue/Queue.php +++ b/packages/queue/src/Queue/Queue.php @@ -4,15 +4,21 @@ namespace Utopia\Queue; +use Utopia\Queue\Option\Reliable; + readonly class Queue { public function __construct( public string $name, public string $namespace = 'utopia-queue', public int $jobTtl = 0, + public ?Reliable $reliable = null, ) { if ($this->name === '' || $this->name === '0') { throw new \InvalidArgumentException('Cannot create queue with empty name.'); } + if ($this->reliable instanceof Reliable && $this->jobTtl > 0) { + throw new \InvalidArgumentException('Reliable queues do not support job TTL.'); + } } } diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php new file mode 100644 index 000000000..338f9810e --- /dev/null +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php @@ -0,0 +1,301 @@ + 'pid', + 'queue' => 'queue', + 'timestamp' => 123, + 'payload' => ['value' => true], + ], '123:456'); + + $this->assertSame('123:456', $message->getClaimedAt()); + $this->assertSame([ + 'pid' => 'pid', + 'queue' => 'queue', + 'timestamp' => 123, + 'payload' => ['value' => true], + ], $message->asArray()); + } + + public function testReliableQueueRejectsJobTtlBeforeAnyConnectionCall(): void + { + $this->expectException(\InvalidArgumentException::class); + $this->expectExceptionMessage('Reliable queues do not support job TTL.'); + + new Queue('ttl', jobTtl: 1, reliable: new Reliable()); + } + + #[\PHPUnit\Framework\Attributes\DataProvider('unsupportedOperationProvider')] + public function testRedisClusterReliableOperationsFailBeforeConnectingOrMutating(string $operation): void + { + $connection = new RedisCluster(['127.0.0.1:1'], connectTimeout: 0.01, readTimeout: 0.01); + $broker = new RedisBroker($connection, $connection); + $queue = new Queue('cluster', reliable: new Reliable()); + $message = new Message([ + 'pid' => 'pid', + 'queue' => $queue->name, + 'timestamp' => 1, + 'payload' => [], + ], '1:1'); + + $this->expectException(\LogicException::class); + $this->expectExceptionMessage('Reliable queues require a single Redis connection with atomic scripting support.'); + + match ($operation) { + 'enqueue' => $broker->enqueue($queue, ['value' => true]), + 'receive' => $broker->receive($queue, 0), + 'commit' => $broker->commit($queue, $message), + 'reject' => $broker->reject($queue, $message), + 'retry' => $broker->retry($queue), + 'failed size' => $broker->getQueueSize($queue, failedJobs: true), + 'extend' => $broker->extend($queue, $message), + 'expired' => $broker->expired($queue, 1), + 'reclaim' => $broker->reclaim($queue, new Claim('pid', '1:1')), + default => throw new \InvalidArgumentException("Unknown operation: {$operation}"), + }; + } + + /** @return iterable */ + public static function unsupportedOperationProvider(): iterable + { + foreach (['enqueue', 'receive', 'commit', 'reject', 'retry', 'failed size', 'extend', 'expired', 'reclaim'] as $operation) { + yield $operation => [$operation]; + } + } + + public function testLockingRedisExposesAtomicCapability(): void + { + $this->assertInstanceOf(Atomic::class, new Locking(new Redis('127.0.0.1', 1))); + } + + public function testLockingDoesNotAdvertiseAtomicSupportForRedisCluster(): void + { + $locking = new Locking(new RedisCluster(['127.0.0.1:1'])); + + $this->assertInstanceOf(Atomic::class, $locking); + $this->assertFalse($locking->supportsAtomic()); + } + + public function testLockingSerializesAtomicEvaluation(): void + { + $recorder = new ReliableRecorder(); + $connection = new RecordingAtomicConnection($recorder); + $locking = new Locking($connection, new RecordingAtomicLock($recorder)); + + $result = $locking->evaluate('return ARGV[1]', ['key', 'value'], 1); + + $this->assertSame('value', $result); + $this->assertSame(['acquire', 'evaluate', 'release'], $recorder->events); + $this->assertSame([['return ARGV[1]', ['key', 'value'], 1]], $connection->calls); + } + + public function testDefaultLockPreventsConcurrentAtomicEvaluationFromInterleaving(): void + { + $connection = new ConcurrentAtomicConnection(); + $locking = new Locking($connection); + + \Swoole\Coroutine\run(function () use ($locking): void { + $waitGroup = new \Swoole\Coroutine\WaitGroup(); + for ($index = 0; $index < 2; $index++) { + $waitGroup->add(); + \Swoole\Coroutine::create(function () use ($locking, $waitGroup): void { + try { + $locking->evaluate('return 1'); + } finally { + $waitGroup->done(); + } + }); + } + $waitGroup->wait(); + }); + + $this->assertSame(1, $connection->peak); + $this->assertSame(['start', 'end', 'start', 'end'], $connection->events); + } + + public function testReliableEmptyReceiveUsesBoundedNonBusyPolling(): void + { + $connection = new PollingAtomicConnection(); + $broker = new RedisBroker($connection, $connection); + $queue = new Queue('polling', reliable: new Reliable( + visibility: 2, + heartbeat: 1, + scan: 1, + batch: 10, + pollMinimum: 10, + pollMaximum: 20, + )); + + $started = hrtime(true); + $this->assertNotInstanceOf(\Utopia\Queue\Message::class, $broker->receive($queue, 1)); + $elapsed = (hrtime(true) - $started) / 1_000_000_000; + + $this->assertGreaterThanOrEqual(0.9, $elapsed); + $this->assertLessThan(1.3, $elapsed); + $this->assertGreaterThanOrEqual(10, $connection->evaluations); + $this->assertLessThanOrEqual(110, $connection->evaluations); + } + + public function testReliableReceiveReconnectsAndResetsAfterSuccessfulEmptyEvaluation(): void + { + $connection = new PollingAtomicConnection(failFirst: true); + $broker = new RedisBroker($connection, $connection); + $queue = new Queue('reconnect', reliable: new Reliable( + visibility: 2, + heartbeat: 1, + scan: 1, + batch: 10, + pollMinimum: 10, + pollMaximum: 20, + )); + $failures = []; + $successes = []; + $broker->setReconnectCallback(function (Queue $queue, \Throwable $error, int $attempt, int $sleep) use (&$failures): void { + $failures[] = [$queue, $error, $attempt, $sleep]; + }); + $broker->setReconnectSuccessCallback(function (Queue $queue, int $attempts) use (&$successes): void { + $successes[] = [$queue, $attempts]; + }); + + $this->assertNotInstanceOf(\Utopia\Queue\Message::class, $broker->receive($queue, 1)); + + $this->assertSame(1, $connection->closes); + $this->assertGreaterThan(1, $connection->evaluations); + $this->assertCount(1, $failures); + $this->assertSame($queue, $failures[0][0]); + $this->assertInstanceOf(\RedisException::class, $failures[0][1]); + $this->assertSame(1, $failures[0][2]); + $this->assertGreaterThanOrEqual(0, $failures[0][3]); + $this->assertLessThanOrEqual(100, $failures[0][3]); + $this->assertSame([[$queue, 1]], $successes); + } +} + +final class RecordingAtomicConnection extends InMemoryConnection implements Atomic +{ + /** @var list, int}> */ + public array $calls = []; + + public function __construct(private readonly ReliableRecorder $recorder) {} + + public function supportsAtomic(): bool + { + return true; + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + $this->recorder->events[] = 'evaluate'; + $this->calls[] = [$script, $arguments, $keyCount]; + + return $arguments[$keyCount] ?? null; + } +} + +final readonly class RecordingAtomicLock implements Lock +{ + public function __construct(private ReliableRecorder $recorder) {} + + public function acquire(float $timeout = 0.0): bool + { + return true; + } + + public function tryAcquire(): bool + { + return true; + } + + public function release(): void {} + + public function withLock(callable $callback, float $timeout = 0.0): mixed + { + $this->recorder->events[] = 'acquire'; + + try { + return $callback(); + } finally { + $this->recorder->events[] = 'release'; + } + } +} + +final class ReliableRecorder +{ + /** @var list */ + public array $events = []; +} + +final class ConcurrentAtomicConnection extends InMemoryConnection implements Atomic +{ + public int $active = 0; + public int $peak = 0; + + /** @var list */ + public array $events = []; + + public function supportsAtomic(): bool + { + return true; + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + $this->events[] = 'start'; + $this->active++; + $this->peak = max($this->peak, $this->active); + \Swoole\Coroutine::sleep(0.01); + $this->active--; + $this->events[] = 'end'; + + return 1; + } +} + +final class PollingAtomicConnection extends InMemoryConnection implements Atomic +{ + public int $closes = 0; + public int $evaluations = 0; + + public function __construct(private readonly bool $failFirst = false) {} + + public function supportsAtomic(): bool + { + return true; + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + $this->evaluations++; + if ($this->failFirst && $this->evaluations === 1) { + throw new \RedisException('Redis is unavailable.'); + } + + return [0]; + } + + #[\Override] + public function close(): void + { + $this->closes++; + } +} diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php new file mode 100644 index 000000000..167703c8a --- /dev/null +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php @@ -0,0 +1,357 @@ +redis = new \Redis(); + $this->redis->connect('127.0.0.1', 16379, 1.0); + $this->queue = new Queue( + 'legacy-ttl-' . bin2hex(random_bytes(6)), + self::NAMESPACE, + jobTtl: 30, + ); + $broker = $this->broker(16379); + $broker->enqueue($this->queue, ['legacy' => true]); + + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $ttl = $this->redis->ttl(self::NAMESPACE . '.jobs.' . $this->queue->name . '.' . $message->getPid()); + $this->assertGreaterThan(0, $ttl); + $this->assertLessThanOrEqual(30, $ttl); + + $broker->commit($this->queue, $message); + } + + #[DataProvider('serverProvider')] + public function testRejectRetryClaimCommitLifecycleHasExactStateAndStats(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + + $this->assertTrue($broker->enqueue($this->queue, ['value' => 'payload'])); + $first = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $first); + $this->assertNotNull($first->getClaimedAt()); + + $firstPid = $first->getPid(); + $broker->reject($this->queue, $first); + $broker->reject($this->queue, $first); + + $this->assertSame(1, $broker->getQueueSize($this->queue, failedJobs: true)); + $this->assertSame(1, $this->stat('failed')); + $this->assertSame(0, $this->stat('processing')); + + $before = $this->serverSeconds(); + $broker->retry($this->queue, 1); + $after = $this->serverSeconds(); + + $this->assertSame(0, $broker->getQueueSize($this->queue, failedJobs: true)); + $retried = $this->pending(); + $this->assertSame(['pid', 'queue', 'timestamp', 'payload'], array_keys($retried)); + $this->assertNotSame($firstPid, $retried['pid']); + $this->assertGreaterThanOrEqual($before, $retried['timestamp']); + $this->assertLessThanOrEqual($after, $retried['timestamp']); + $this->assertSame(['value' => 'payload'], $retried['payload']); + $this->assertFalse($this->redis->hExists($this->key('jobs'), $firstPid)); + $this->assertSame(1, $this->stat('retried')); + + $second = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $second); + $this->assertNotSame($firstPid, $second->getPid()); + $broker->commit($this->queue, $second); + $broker->commit($this->queue, $second); + $broker->reject($this->queue, $second); + + $this->assertSame(2, $this->stat('total')); + $this->assertSame(0, $this->stat('processing')); + $this->assertSame(1, $this->stat('failed')); + $this->assertSame(1, $this->stat('retried')); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(0, $this->redis->hLen($this->key('jobs'))); + $this->assertSame(0, $this->redis->zCard($this->key('processing'))); + $this->assertSame(0, $broker->getQueueSize($this->queue, failedJobs: true)); + } + + #[DataProvider('serverProvider')] + public function testReliableEnqueuePreservesPriorityAndFifoOrdering(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + $broker->enqueue($this->queue, ['order' => 'normal-1']); + $broker->enqueue($this->queue, ['order' => 'normal-2']); + $broker->enqueue($this->queue, ['order' => 'priority'], priority: true); + + foreach (['priority', 'normal-1', 'normal-2'] as $expected) { + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $this->assertSame($expected, $message->getPayload()['order']); + $broker->commit($this->queue, $message); + } + + $this->assertSame(0, $broker->getQueueSize($this->queue)); + $this->assertSame(3, $this->stat('success')); + } + + #[DataProvider('serverProvider')] + public function testReclaimPreservesEnqueueDataUsesNewPidAndPlainEnvelope(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + $broker->enqueue($this->queue, ['number' => 42]); + + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $oldPid = $message->getPid(); + $timestamp = $message->getTimestamp(); + $this->redis->zAdd($this->key('processing'), 0, $oldPid); + + $claims = $broker->expired($this->queue, 10); + $this->assertCount(1, $claims); + $replacement = $broker->reclaim($this->queue, $claims[0]); + + $this->assertInstanceOf(Message::class, $replacement); + $this->assertNull($replacement->getClaimedAt()); + $this->assertNotSame($oldPid, $replacement->getPid()); + $this->assertSame($timestamp, $replacement->getTimestamp()); + $this->assertSame(['number' => 42], $replacement->getPayload()); + $this->assertSame(['pid', 'queue', 'timestamp', 'payload'], array_keys($this->pending())); + $this->assertFalse($this->redis->hExists($this->key('jobs'), $oldPid)); + $this->assertSame(0, $this->redis->zCard($this->key('processing'))); + $this->assertSame(1, $this->stat('reclaimed')); + + $broker->commit($this->queue, $message); + $broker->reject($this->queue, $message); + $this->assertSame(0, $this->stat('success')); + $this->assertSame(0, $this->stat('failed')); + $this->assertSame(0, $this->stat('processing')); + + $this->assertNotInstanceOf(\Utopia\Queue\Message::class, $broker->reclaim($this->queue, $claims[0])); + $this->assertSame(1, $broker->getQueueSize($this->queue)); + $this->assertSame(1, $this->stat('reclaimed')); + + $claimedAgain = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $claimedAgain); + $this->assertNotNull($claimedAgain->getClaimedAt()); + $broker->commit($this->queue, $claimedAgain); + } + + #[DataProvider('serverProvider')] + public function testHeartbeatRenewsOnlyLeaseAndOriginalTokenStillCommits(int $port): void + { + $this->prepare($port, visibility: 2); + $broker = $this->broker($port); + $broker->enqueue($this->queue, ['slow' => true]); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + + $token = $message->getClaimedAt(); + $firstScore = (float) $this->redis->zScore($this->key('processing'), $message->getPid()); + usleep(20_000); + $this->assertTrue($broker->extend($this->queue, $message)); + $secondScore = (float) $this->redis->zScore($this->key('processing'), $message->getPid()); + usleep(20_000); + $this->assertTrue($broker->extend($this->queue, $message)); + + $record = json_decode((string) $this->redis->hGet($this->key('jobs'), $message->getPid()), true); + $this->assertSame($token, $message->getClaimedAt()); + $this->assertSame($token, $record['claimedAt']); + $this->assertGreaterThan($firstScore, $secondScore); + $this->assertGreaterThan($secondScore, (float) $this->redis->zScore($this->key('processing'), $message->getPid())); + + $broker->commit($this->queue, $message); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(0, $this->stat('processing')); + } + + #[DataProvider('serverProvider')] + public function testRetryAndReclaimRacesHaveOneWinner(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + $broker->enqueue($this->queue, ['winner' => 'reclaim']); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $this->redis->zAdd($this->key('processing'), 0, $message->getPid()); + $claim = $broker->expired($this->queue, 1)[0]; + + $this->assertInstanceOf(Message::class, $broker->reclaim($this->queue, $claim)); + $broker->reject($this->queue, $message); + $broker->retry($this->queue); + $this->assertSame(1, $broker->getQueueSize($this->queue)); + $this->assertSame(0, $broker->getQueueSize($this->queue, failedJobs: true)); + + $this->cleanup(); + $broker->enqueue($this->queue, ['winner' => 'retry']); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $broker->reject($this->queue, $message); + $broker->retry($this->queue, 1); + $broker->retry($this->queue, 1); + + $this->assertSame([], $broker->expired($this->queue, 10)); + $this->assertSame(1, $broker->getQueueSize($this->queue)); + $this->assertSame(0, $broker->getQueueSize($this->queue, failedJobs: true)); + $this->assertSame(1, $this->stat('retried')); + } + + #[DataProvider('serverProvider')] + public function testMalformedPendingAndMissingProcessingStateAreQuarantined(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + $this->redis->lPush($this->key('queue'), '{invalid-json'); + + $started = hrtime(true); + $this->assertNotInstanceOf(\Utopia\Queue\Message::class, $broker->receive($this->queue, 1)); + $elapsed = (hrtime(true) - $started) / 1_000_000_000; + + $this->assertGreaterThanOrEqual(0.9, $elapsed); + $this->assertLessThan(1.5, $elapsed); + $this->assertSame(1, $this->redis->lLen($this->key('quarantine'))); + $this->assertSame(1, $this->stat('quarantined')); + + $missingPid = 'missing-pid'; + $this->redis->zAdd($this->key('processing'), 0, $missingPid); + $this->redis->set($this->statKey('processing'), '1'); + $claims = $broker->expired($this->queue, 10); + $this->assertCount(1, $claims); + $this->assertNull($claims[0]->claimedAt); + $this->assertNotInstanceOf(\Utopia\Queue\Message::class, $broker->reclaim($this->queue, $claims[0])); + + $this->assertSame(2, $this->redis->lLen($this->key('quarantine'))); + $this->assertSame(2, $this->stat('quarantined')); + $this->assertSame(0, $this->redis->zCard($this->key('processing'))); + $this->assertSame(0, $this->stat('processing')); + } + + #[DataProvider('serverProvider')] + public function testMissingFailedStateIsQuarantinedWithoutBlockingValidRetry(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + $broker->enqueue($this->queue, ['valid' => true]); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $broker->reject($this->queue, $message); + $this->redis->rPush($this->key('failed'), 'missing-pid'); + + $broker->retry($this->queue); + + $this->assertSame(0, $broker->getQueueSize($this->queue, failedJobs: true)); + $this->assertSame(1, $broker->getQueueSize($this->queue)); + $this->assertSame(1, $this->redis->lLen($this->key('quarantine'))); + $this->assertSame(1, $this->stat('quarantined')); + $this->assertSame(1, $this->stat('retried')); + } + + /** @return iterable */ + public static function serverProvider(): iterable + { + yield 'Redis' => [16379]; + yield 'Dragonfly' => [16380]; + } + + protected function tearDown(): void + { + if (isset($this->redis)) { + $this->cleanup(); + $this->redis->close(); + } + } + + private function prepare(int $port, int $visibility = 2): void + { + $this->redis = new \Redis(); + $this->redis->connect('127.0.0.1', $port, 1.0); + $name = 'atomic-' . bin2hex(random_bytes(6)); + $this->queue = new Queue( + $name, + self::NAMESPACE, + reliable: new Reliable(visibility: $visibility, heartbeat: 1, scan: 1, batch: 100), + ); + $this->cleanup(); + } + + private function broker(int $port): RedisBroker + { + return new RedisBroker( + new RedisConnection('127.0.0.1', $port, connectTimeout: 1.0, readTimeout: 2.0), + new Locking(new RedisConnection('127.0.0.1', $port, connectTimeout: 1.0, readTimeout: 2.0)), + ); + } + + private function cleanup(): void + { + $this->redis->del([ + $this->key('queue'), + $this->key('jobs'), + $this->key('processing'), + $this->key('failed'), + $this->key('quarantine'), + $this->statKey('total'), + $this->statKey('processing'), + $this->statKey('success'), + $this->statKey('failed'), + $this->statKey('retried'), + $this->statKey('reclaimed'), + $this->statKey('quarantined'), + ]); + } + + private function key(string $type): string + { + if ($type === 'queue') { + return self::NAMESPACE . '.queue.' . $this->queue->name; + } + + return self::NAMESPACE . '.atomic.' . $type . '.' . $this->queue->name; + } + + private function statKey(string $stat): string + { + return self::NAMESPACE . '.stats.' . $this->queue->name . '.' . $stat; + } + + private function stat(string $stat): int + { + return (int) ($this->redis->get($this->statKey($stat)) ?: 0); + } + + /** @return array{pid: string, queue: string, timestamp: int, payload: array} */ + private function pending(): array + { + $value = $this->redis->lIndex($this->key('queue'), 0); + $this->assertIsString($value); + $decoded = json_decode($value, true); + $this->assertIsArray($decoded); + + return $decoded; + } + + private function serverSeconds(): int + { + $time = $this->redis->time(); + $this->assertIsArray($time); + + return (int) ltrim((string) $time[0], ':'); + } +} diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php new file mode 100644 index 000000000..c95e5adf0 --- /dev/null +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php @@ -0,0 +1,253 @@ +prepare(); + $publisher = $this->broker(); + $messages = $maxCoroutines * 3; + for ($index = 0; $index < $messages; $index++) { + $publisher->enqueue($this->queue, ['index' => $index]); + } + + $peakClaims = 0; + $processed = 0; + + Coroutine\run(function () use ($maxCoroutines, $messages, &$peakClaims, &$processed): void { + $broker = $this->broker(); + $adapter = new Swoole( + $broker, + 1, + $this->queue->name, + self::NAMESPACE, + maxCoroutines: $maxCoroutines, + reliable: $this->queue->reliable, + ); + + $adapter->consume( + function () use ($adapter, $messages, &$peakClaims, &$processed): void { + $redis = new \Redis(); + $redis->connect('127.0.0.1', self::PORT, 1.0); + $peakClaims = max($peakClaims, $redis->zCard($this->key('processing'))); + $redis->close(); + Coroutine::sleep(0.03); + + if (++$processed === $messages) { + $adapter->stop(); + } + }, + static fn(): null => null, + static fn(): null => null, + ); + }); + + $this->assertSame($messages, $processed); + $this->assertLessThanOrEqual($maxCoroutines, $peakClaims); + $this->assertSame(0, $this->redis->zCard($this->key('processing'))); + } + + public function testHeartbeatKeepsSlowHandlerOwnedUntilCommit(): void + { + $this->prepare(visibility: 2, heartbeat: 1, scan: 1); + $this->broker()->enqueue($this->queue, ['slow' => true]); + $processed = 0; + $token = null; + + Coroutine\run(function () use (&$processed, &$token): void { + $adapter = new Swoole( + $this->broker(), + 1, + $this->queue->name, + self::NAMESPACE, + maxCoroutines: 1, + reliable: $this->queue->reliable, + ); + + $adapter->consume( + function (Message $message) use ($adapter, &$processed, &$token): void { + $token = $message->getClaimedAt(); + Coroutine::sleep(3.2); + $processed++; + $adapter->stop(); + }, + static fn(): null => null, + static fn(): null => null, + ); + }); + + $this->assertNotNull($token); + $this->assertSame(1, $processed); + $this->assertSame(0, $this->redis->lLen($this->key('queue'))); + $this->assertSame(0, $this->stat('reclaimed')); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(0, $this->stat('processing')); + } + + public function testConcurrentRecoveryLoopsProduceOneReplacement(): void + { + $this->prepare(); + $broker = $this->broker(); + $broker->enqueue($this->queue, ['race' => true]); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $this->redis->zAdd($this->key('processing'), 0, $message->getPid()); + $claim = $broker->expired($this->queue, 1)[0]; + $results = []; + + Coroutine\run(function () use ($claim, &$results): void { + $waitGroup = new WaitGroup(); + for ($index = 0; $index < 2; $index++) { + $waitGroup->add(); + Coroutine::create(function () use ($claim, $waitGroup, &$results): void { + try { + $results[] = $this->broker()->reclaim($this->queue, $claim); + } finally { + $waitGroup->done(); + } + }); + } + $waitGroup->wait(); + }); + + $messages = array_values(array_filter($results, static fn(mixed $result): bool => $result instanceof Message)); + $this->assertCount(1, $messages); + $this->assertSame(1, $this->redis->lLen($this->key('queue'))); + $this->assertSame(1, $this->stat('reclaimed')); + $this->assertSame(0, $this->stat('processing')); + } + + public function testRecoveryDrainsConsecutiveBoundedBatchesInOneScan(): void + { + $this->prepare(scan: 1, batch: 2); + $broker = $this->broker(); + for ($index = 0; $index < 5; $index++) { + $broker->enqueue($this->queue, ['index' => $index]); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $this->redis->zAdd($this->key('processing'), 0, $message->getPid()); + } + $processed = 0; + + Coroutine\run(function () use (&$processed): void { + $adapter = new Swoole( + $this->broker(), + 1, + $this->queue->name, + self::NAMESPACE, + maxCoroutines: 1, + reliable: $this->queue->reliable, + ); + $adapter->consume( + function () use ($adapter, &$processed): void { + if (++$processed === 5) { + $adapter->stop(); + } + }, + static fn(): null => null, + static fn(): null => null, + ); + }); + + $this->assertSame(5, $processed); + $this->assertSame(5, $this->stat('reclaimed')); + $this->assertSame(5, $this->stat('success')); + $this->assertSame(0, $this->stat('processing')); + $this->assertSame(0, $this->redis->lLen($this->key('queue'))); + } + + /** @return iterable */ + public static function capacityProvider(): iterable + { + yield 'one coroutine' => [1]; + yield 'three coroutines' => [3]; + } + + protected function tearDown(): void + { + if (isset($this->redis)) { + $this->cleanup(); + $this->redis->close(); + } + } + + private function prepare(int $visibility = 2, int $heartbeat = 1, int $scan = 1, int $batch = 100): void + { + $this->redis = new \Redis(); + $this->redis->connect('127.0.0.1', self::PORT, 1.0); + $this->queue = new Queue( + 'atomic-' . bin2hex(random_bytes(6)), + self::NAMESPACE, + reliable: new Reliable(visibility: $visibility, heartbeat: $heartbeat, scan: $scan, batch: $batch), + ); + $this->cleanup(); + } + + private function broker(): RedisBroker + { + return new RedisBroker( + new RedisConnection('127.0.0.1', self::PORT, connectTimeout: 1.0, readTimeout: 2.0), + new Locking(new RedisConnection('127.0.0.1', self::PORT, connectTimeout: 1.0, readTimeout: 2.0)), + ); + } + + private function cleanup(): void + { + $this->redis->del([ + $this->key('queue'), + $this->key('jobs'), + $this->key('processing'), + $this->key('failed'), + $this->key('quarantine'), + $this->statKey('total'), + $this->statKey('processing'), + $this->statKey('success'), + $this->statKey('failed'), + $this->statKey('retried'), + $this->statKey('reclaimed'), + $this->statKey('quarantined'), + ]); + } + + private function key(string $type): string + { + if ($type === 'queue') { + return self::NAMESPACE . '.queue.' . $this->queue->name; + } + + return self::NAMESPACE . '.atomic.' . $type . '.' . $this->queue->name; + } + + private function statKey(string $stat): string + { + return self::NAMESPACE . '.stats.' . $this->queue->name . '.' . $stat; + } + + private function stat(string $stat): int + { + return (int) ($this->redis->get($this->statKey($stat)) ?: 0); + } +} From 7f816c5d129a65684a294beb017f95475c82cd0d Mon Sep 17 00:00:00 2001 From: Jake Barnby Date: Thu, 16 Jul 2026 01:10:00 +1200 Subject: [PATCH 2/4] (fix): harden reliable Redis queue recovery --- packages/queue/src/Queue/Adapter/Swoole.php | 46 ++++- packages/queue/src/Queue/Broker/Redis.php | 5 +- .../queue/src/Queue/Broker/Redis/Script.php | 97 ++++++---- .../Queue/E2E/Adapter/ReliableApiTest.php | 53 +++++- .../Queue/E2E/Adapter/ReliableRedisTest.php | 176 +++++++++++++++++- .../Queue/E2E/Adapter/ReliableSwooleTest.php | 122 ++++++++++++ .../E2E/Adapter/reliable-create-failure.php | 55 ++++++ 7 files changed, 511 insertions(+), 43 deletions(-) create mode 100644 packages/queue/tests/Queue/E2E/Adapter/reliable-create-failure.php diff --git a/packages/queue/src/Queue/Adapter/Swoole.php b/packages/queue/src/Queue/Adapter/Swoole.php index 42ead27b6..754ef80e2 100644 --- a/packages/queue/src/Queue/Adapter/Swoole.php +++ b/packages/queue/src/Queue/Adapter/Swoole.php @@ -124,7 +124,7 @@ private function consumeLegacy(callable $messageCallback, callable $successCallb $slots->push(true); $waitGroup->add(); - Coroutine::create(function () use ($message, $messageCallback, $successCallback, $errorCallback, $slots, $waitGroup): void { + $coroutine = Coroutine::create(function () use ($message, $messageCallback, $successCallback, $errorCallback, $slots, $waitGroup): void { try { $this->process($message, $messageCallback, $successCallback, $errorCallback); } catch (\Throwable $error) { @@ -135,6 +135,11 @@ private function consumeLegacy(callable $messageCallback, callable $successCallb $slots->pop(); } }); + if ($coroutine === false) { + $waitGroup->done(); + $slots->pop(); + $this->process($message, $messageCallback, $successCallback, $errorCallback); + } } $waitGroup->wait(); @@ -154,7 +159,7 @@ private function consumeReliable( $recoveryDone = new Channel(1); $recovery = new WaitGroup(1); - Coroutine::create(function () use ($recoverable, $reliable, $recoveryDone, $recovery): void { + $recoveryCoroutine = Coroutine::create(function () use ($recoverable, $reliable, $recoveryDone, $recovery): void { try { while ($recoveryDone->pop($reliable->scan) === false) { try { @@ -175,6 +180,13 @@ private function consumeReliable( $recovery->done(); } }); + if ($recoveryCoroutine === false) { + $recovery->done(); + $recoveryDone->close(); + $this->stopConsumption(); + + throw new \RuntimeException('Failed to create queue recovery coroutine.'); + } try { while (!$this->isStopped()) { @@ -196,7 +208,7 @@ private function consumeReliable( } $handlers->add(); - Coroutine::create(function () use ( + $handlerCoroutine = Coroutine::create(function () use ( $message, $messageCallback, $successCallback, @@ -208,7 +220,7 @@ private function consumeReliable( ): void { $heartbeatDone = new Channel(1); $heartbeat = new WaitGroup(1); - Coroutine::create(function () use ($message, $recoverable, $reliable, $heartbeatDone, $heartbeat): void { + $heartbeatCoroutine = Coroutine::create(function () use ($message, $recoverable, $reliable, $heartbeatDone, $heartbeat): void { try { while ($heartbeatDone->pop($reliable->heartbeat) === false) { if (!$recoverable->extend($this->queue, $message)) { @@ -222,6 +234,15 @@ private function consumeReliable( $heartbeat->done(); } }); + if ($heartbeatCoroutine === false) { + $heartbeat->done(); + $heartbeatDone->close(); + $this->stopConsumption(); + $handlers->done(); + $slots->pop(); + + return; + } try { $this->process($message, $messageCallback, $successCallback, $errorCallback); @@ -234,6 +255,13 @@ private function consumeReliable( $slots->pop(); } }); + if ($handlerCoroutine === false) { + $handlers->done(); + $slots->pop(); + $this->stopConsumption(); + + throw new \RuntimeException('Failed to create queue handler coroutine.'); + } } } finally { $handlers->wait(); @@ -272,6 +300,16 @@ public function stop(): self return $this; } + private function stopConsumption(): void + { + $this->stopped = true; + + try { + $this->consumer->close(); + } catch (\Throwable) { + } + } + public function workerStart(callable $callback): self { $this->onWorkerStart[] = $callback; diff --git a/packages/queue/src/Queue/Broker/Redis.php b/packages/queue/src/Queue/Broker/Redis.php index 0afba5191..9d453eebd 100644 --- a/packages/queue/src/Queue/Broker/Redis.php +++ b/packages/queue/src/Queue/Broker/Redis.php @@ -273,10 +273,13 @@ private function getJob(Queue $queue, string $pid): Message|false public function getQueueSize(Queue $queue, bool $failedJobs = false): int { + if ($queue->reliable instanceof Reliable) { + $this->atomic($this->commands); + } + $queueName = $this->pendingKey($queue); if ($failedJobs) { if ($queue->reliable instanceof Reliable) { - $this->atomic($this->commands); $queueName = $this->atomicKey($queue, 'failed'); } else { $queueName = "{$queue->namespace}.failed.{$queue->name}"; diff --git a/packages/queue/src/Queue/Broker/Redis/Script.php b/packages/queue/src/Queue/Broker/Redis/Script.php index 81fe81fe6..d5bb8c06b 100644 --- a/packages/queue/src/Queue/Broker/Redis/Script.php +++ b/packages/queue/src/Queue/Broker/Redis/Script.php @@ -8,17 +8,14 @@ final class Script { public const string ENQUEUE = <<<'LUA' local now = redis.call('TIME') -local payload = cjson.decode(ARGV[3]) -local message = { - pid = ARGV[1], - queue = ARGV[2], - timestamp = tonumber(now[1]), - payload = payload -} -local encoded = '{"pid":' .. cjson.encode(message.pid) - .. ',"queue":' .. cjson.encode(message.queue) - .. ',"timestamp":' .. tostring(message.timestamp) - .. ',"payload":' .. cjson.encode(message.payload) .. '}' +local valid, payload = pcall(cjson.decode, ARGV[3]) +if not valid or type(payload) ~= 'table' then + return 0 +end +local encoded = '{"pid":' .. cjson.encode(ARGV[1]) + .. ',"queue":' .. cjson.encode(ARGV[2]) + .. ',"timestamp":' .. tostring(tonumber(now[1])) + .. ',"payload":' .. ARGV[3] .. '}' if ARGV[4] == '1' then redis.call('RPUSH', KEYS[1], encoded) else @@ -34,13 +31,24 @@ final class Script end local valid, message = pcall(cjson.decode, raw) +local marker = ',"payload":' +local payloadStart = string.find(raw, marker, 1, true) +local encodedPayload = nil +if payloadStart and string.sub(raw, -1) == '}' then + encodedPayload = string.sub(raw, payloadStart + string.len(marker), -2) +end +local payloadValid, payload = pcall(cjson.decode, encodedPayload or '') if not valid or type(message) ~= 'table' or type(message.pid) ~= 'string' or message.pid == '' or type(message.queue) ~= 'string' or type(message.timestamp) ~= 'number' - or type(message.payload) ~= 'table' then + or type(message.payload) ~= 'table' + or not encodedPayload + or encodedPayload == '' + or not payloadValid + or type(payload) ~= 'table' then redis.call('LPUSH', KEYS[4], cjson.encode({reason = 'malformed_pending', raw = raw})) redis.call('INCR', KEYS[7]) return {2} @@ -57,7 +65,9 @@ final class Script local micros = (tonumber(now[1]) * 1000000) + tonumber(now[2]) local leaseUntil = micros + (tonumber(ARGV[1]) * 1000000) local record = { - message = message, + queue = message.queue, + timestamp = message.timestamp, + payload = encodedPayload, state = 'processing', claimedAt = claimedAt, leaseUntil = leaseUntil @@ -128,13 +138,22 @@ final class Script local raw = redis.call('HGET', KEYS[2], pid) local valid, record = pcall(cjson.decode, raw or '') +local payloadValid = false +local payload = nil +if valid and type(record) == 'table' then + payloadValid, payload = pcall(cjson.decode, record.payload or '') +end if not raw or not valid or type(record) ~= 'table' or record.state ~= 'failed' - or type(record.message) ~= 'table' - or type(record.message.queue) ~= 'string' - or type(record.message.payload) ~= 'table' then + or type(record.queue) ~= 'string' + or record.queue == '' + or type(record.timestamp) ~= 'number' + or type(record.payload) ~= 'string' + or record.payload == '' + or not payloadValid + or type(payload) ~= 'table' then redis.call('HDEL', KEYS[2], pid) redis.call('LPUSH', KEYS[4], cjson.encode({reason = 'malformed_failed', pid = pid, raw = raw or false})) redis.call('INCR', KEYS[6]) @@ -147,16 +166,15 @@ final class Script end local now = redis.call('TIME') -local message = { +local replacement = { pid = ARGV[1], - queue = record.message.queue, + queue = record.queue, timestamp = tonumber(now[1]), - payload = record.message.payload } -local encoded = '{"pid":' .. cjson.encode(message.pid) - .. ',"queue":' .. cjson.encode(message.queue) - .. ',"timestamp":' .. tostring(message.timestamp) - .. ',"payload":' .. cjson.encode(message.payload) .. '}' +local encoded = '{"pid":' .. cjson.encode(replacement.pid) + .. ',"queue":' .. cjson.encode(replacement.queue) + .. ',"timestamp":' .. tostring(replacement.timestamp) + .. ',"payload":' .. record.payload .. '}' redis.call('LPUSH', KEYS[3], encoded) redis.call('HDEL', KEYS[2], pid) redis.call('INCR', KEYS[5]) @@ -201,14 +219,24 @@ final class Script local raw = redis.call('HGET', KEYS[2], ARGV[1]) local valid, record = pcall(cjson.decode, raw or '') +local payloadValid = false +local payload = nil +if valid and type(record) == 'table' then + payloadValid, payload = pcall(cjson.decode, record.payload or '') +end if not raw or not valid or type(record) ~= 'table' or record.state ~= 'processing' - or type(record.message) ~= 'table' - or type(record.message.queue) ~= 'string' - or type(record.message.timestamp) ~= 'number' - or type(record.message.payload) ~= 'table' then + or type(record.claimedAt) ~= 'string' + or record.claimedAt == '' + or type(record.queue) ~= 'string' + or record.queue == '' + or type(record.timestamp) ~= 'number' + or type(record.payload) ~= 'string' + or record.payload == '' + or not payloadValid + or type(payload) ~= 'table' then redis.call('ZREM', KEYS[1], ARGV[1]) redis.call('HDEL', KEYS[2], ARGV[1]) redis.call('LPUSH', KEYS[4], cjson.encode({reason = 'missing_or_malformed_processing', pid = ARGV[1], raw = raw or false})) @@ -228,16 +256,15 @@ final class Script return {3} end -local message = { +local replacement = { pid = ARGV[3], - queue = record.message.queue, - timestamp = record.message.timestamp, - payload = record.message.payload + queue = record.queue, + timestamp = record.timestamp, } -local encoded = '{"pid":' .. cjson.encode(message.pid) - .. ',"queue":' .. cjson.encode(message.queue) - .. ',"timestamp":' .. tostring(message.timestamp) - .. ',"payload":' .. cjson.encode(message.payload) .. '}' +local encoded = '{"pid":' .. cjson.encode(replacement.pid) + .. ',"queue":' .. cjson.encode(replacement.queue) + .. ',"timestamp":' .. tostring(replacement.timestamp) + .. ',"payload":' .. record.payload .. '}' redis.call('RPUSH', KEYS[3], encoded) redis.call('HDEL', KEYS[2], ARGV[1]) redis.call('ZREM', KEYS[1], ARGV[1]) diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php index 338f9810e..660d79176 100644 --- a/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php @@ -66,6 +66,7 @@ public function testRedisClusterReliableOperationsFailBeforeConnectingOrMutating 'commit' => $broker->commit($queue, $message), 'reject' => $broker->reject($queue, $message), 'retry' => $broker->retry($queue), + 'pending size' => $broker->getQueueSize($queue), 'failed size' => $broker->getQueueSize($queue, failedJobs: true), 'extend' => $broker->extend($queue, $message), 'expired' => $broker->expired($queue, 1), @@ -77,11 +78,38 @@ public function testRedisClusterReliableOperationsFailBeforeConnectingOrMutating /** @return iterable */ public static function unsupportedOperationProvider(): iterable { - foreach (['enqueue', 'receive', 'commit', 'reject', 'retry', 'failed size', 'extend', 'expired', 'reclaim'] as $operation) { + foreach (['enqueue', 'receive', 'commit', 'reject', 'retry', 'pending size', 'failed size', 'extend', 'expired', 'reclaim'] as $operation) { yield $operation => [$operation]; } } + #[\PHPUnit\Framework\Attributes\DataProvider('queueSizeProvider')] + public function testReliableQueueSizeFailsBeforeReadingUnsupportedConnection(bool $failedJobs): void + { + $connection = new UnsupportedAtomicConnection(); + $broker = new RedisBroker($connection, $connection); + $queue = new Queue('unsupported-size', reliable: new Reliable()); + + try { + $broker->getQueueSize($queue, $failedJobs); + $this->fail('Reliable queue size should reject unsupported atomic connections.'); + } catch (\LogicException $error) { + $this->assertSame( + 'Reliable queues require a single Redis connection with atomic scripting support.', + $error->getMessage(), + ); + } + + $this->assertSame(0, $connection->reads); + } + + /** @return iterable */ + public static function queueSizeProvider(): iterable + { + yield 'pending' => [false]; + yield 'failed' => [true]; + } + public function testLockingRedisExposesAtomicCapability(): void { $this->assertInstanceOf(Atomic::class, new Locking(new Redis('127.0.0.1', 1))); @@ -190,6 +218,29 @@ public function testReliableReceiveReconnectsAndResetsAfterSuccessfulEmptyEvalua } } +final class UnsupportedAtomicConnection extends InMemoryConnection implements Atomic +{ + public int $reads = 0; + + public function supportsAtomic(): bool + { + return false; + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + throw new \LogicException('Atomic evaluation should not be attempted.'); + } + + #[\Override] + public function listSize(string $key): int + { + $this->reads++; + + return parent::listSize($key); + } +} + final class RecordingAtomicConnection extends InMemoryConnection implements Atomic { /** @var list, int}> */ diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php index 167703c8a..f340678cb 100644 --- a/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php @@ -263,6 +263,169 @@ public function testMissingFailedStateIsQuarantinedWithoutBlockingValidRetry(int $this->assertSame(1, $this->stat('retried')); } + #[DataProvider('serverProvider')] + public function testMalformedClaimsFillingBatchAreQuarantinedWithoutStarvingValidRecovery(int $port): void + { + $batch = 3; + $this->prepare($port, batch: $batch); + $broker = $this->broker($port); + $broker->enqueue($this->queue, ['valid' => true]); + $valid = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $valid); + $this->redis->zAdd($this->key('processing'), 1, $valid->getPid()); + + $records = [ + 'missing-claimed-at' => null, + 'non-string-claimed-at' => 123, + 'empty-claimed-at' => '', + ]; + foreach ($records as $pid => $claimedAt) { + $payload = json_encode(['poison' => $pid], JSON_THROW_ON_ERROR); + $record = [ + 'queue' => $this->queue->name, + 'timestamp' => 1, + 'payload' => $payload, + 'state' => 'processing', + 'leaseUntil' => 0, + ]; + if ($pid !== 'missing-claimed-at') { + $record['claimedAt'] = $claimedAt; + } + $this->redis->hSet($this->key('jobs'), $pid, json_encode($record, JSON_THROW_ON_ERROR)); + $this->redis->zAdd($this->key('processing'), 0, $pid); + } + $this->redis->set($this->statKey('processing'), (string) (\count($records) + 1)); + + $replacements = []; + $seen = []; + for ($scan = 0; $scan < 2; $scan++) { + $claims = $broker->expired($this->queue, $batch); + array_push($seen, ...$claims); + foreach ($claims as $claim) { + $replacement = $broker->reclaim($this->queue, $claim); + if ($replacement instanceof Message) { + $replacements[] = $replacement; + } + } + if (\count($claims) < $batch) { + break; + } + } + + $this->assertCount(1, $replacements); + $this->assertSame(['valid' => true], $replacements[0]->getPayload()); + $this->assertSame(1, $broker->getQueueSize($this->queue)); + $this->assertSame(0, $this->redis->zCard($this->key('processing'))); + $this->assertSame(0, $this->redis->hLen($this->key('jobs'))); + $this->assertSame(3, $this->redis->lLen($this->key('quarantine'))); + $this->assertSame(0, $this->stat('processing')); + $this->assertSame(1, $this->stat('reclaimed')); + $this->assertSame(3, $this->stat('quarantined')); + + foreach ($seen as $claim) { + $this->assertNotInstanceOf(Message::class, $broker->reclaim($this->queue, $claim)); + } + $this->assertSame(1, $this->stat('reclaimed')); + $this->assertSame(3, $this->stat('quarantined')); + } + + #[DataProvider('serverProvider')] + public function testRetryPreservesFailedFifoOrderingAcrossMultipleJobs(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + $order = ['first', 'second', 'third']; + + foreach ($order as $value) { + $broker->enqueue($this->queue, ['order' => $value]); + } + foreach ($order as $value) { + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $this->assertSame($value, $message->getPayload()['order']); + $broker->reject($this->queue, $message); + } + + $this->assertSame(3, $broker->getQueueSize($this->queue, failedJobs: true)); + $broker->retry($this->queue); + $this->assertSame(0, $broker->getQueueSize($this->queue, failedJobs: true)); + + foreach ($order as $value) { + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $this->assertSame($value, $message->getPayload()['order']); + $broker->commit($this->queue, $message); + } + + $this->assertSame(3, $this->stat('retried')); + $this->assertSame(3, $this->stat('success')); + $this->assertSame(0, $this->stat('processing')); + } + + #[DataProvider('serverProvider')] + public function testPayloadEnvelopeRemainsOpaqueAcrossClaimRetryAndReclaim(int $port): void + { + $this->prepare($port); + $broker = $this->broker($port); + $payload = [ + 'large' => 9_007_199_254_740_993, + 'ordered' => [ + 'zebra' => 1, + 'alpha' => 2, + 'middle' => 3, + ], + ]; + $encodedPayload = json_encode($payload, JSON_THROW_ON_ERROR); + + $broker->enqueue($this->queue, $payload); + $envelope = $this->redis->lIndex($this->key('queue'), 0); + $this->assertIsString($envelope); + $this->assertSame($encodedPayload, $this->encodedPayload($envelope)); + $decoded = json_decode($envelope, true, flags: JSON_THROW_ON_ERROR); + $this->assertSame($payload, $decoded['payload']); + + $claimed = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $claimed); + $this->assertSame($payload, $claimed->getPayload()); + $record = json_decode( + (string) $this->redis->hGet($this->key('jobs'), $claimed->getPid()), + true, + flags: JSON_THROW_ON_ERROR, + ); + $this->assertArrayNotHasKey('message', $record); + $this->assertSame($this->queue->name, $record['queue']); + $this->assertSame($claimed->getTimestamp(), $record['timestamp']); + $this->assertSame($encodedPayload, $record['payload']); + + $broker->reject($this->queue, $claimed); + $before = $this->serverSeconds(); + $broker->retry($this->queue); + $after = $this->serverSeconds(); + $retried = $this->redis->lIndex($this->key('queue'), 0); + $this->assertIsString($retried); + $this->assertSame($encodedPayload, $this->encodedPayload($retried)); + $retriedEnvelope = json_decode($retried, true, flags: JSON_THROW_ON_ERROR); + $this->assertNotSame($claimed->getPid(), $retriedEnvelope['pid']); + $this->assertGreaterThanOrEqual($before, $retriedEnvelope['timestamp']); + $this->assertLessThanOrEqual($after, $retriedEnvelope['timestamp']); + $this->assertSame($payload, $retriedEnvelope['payload']); + $claimedAgain = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $claimedAgain); + $this->assertSame($retriedEnvelope['pid'], $claimedAgain->getPid()); + $this->assertSame($payload, $claimedAgain->getPayload()); + + $this->redis->zAdd($this->key('processing'), 0, $claimedAgain->getPid()); + $claim = $broker->expired($this->queue, 1)[0]; + $reclaimed = $broker->reclaim($this->queue, $claim); + $this->assertInstanceOf(Message::class, $reclaimed); + $this->assertNotSame($claimedAgain->getPid(), $reclaimed->getPid()); + $this->assertSame($claimedAgain->getTimestamp(), $reclaimed->getTimestamp()); + $this->assertSame($payload, $reclaimed->getPayload()); + $reclaimedEnvelope = $this->redis->lIndex($this->key('queue'), 0); + $this->assertIsString($reclaimedEnvelope); + $this->assertSame($encodedPayload, $this->encodedPayload($reclaimedEnvelope)); + } + /** @return iterable */ public static function serverProvider(): iterable { @@ -278,7 +441,7 @@ protected function tearDown(): void } } - private function prepare(int $port, int $visibility = 2): void + private function prepare(int $port, int $visibility = 2, int $batch = 100): void { $this->redis = new \Redis(); $this->redis->connect('127.0.0.1', $port, 1.0); @@ -286,7 +449,7 @@ private function prepare(int $port, int $visibility = 2): void $this->queue = new Queue( $name, self::NAMESPACE, - reliable: new Reliable(visibility: $visibility, heartbeat: 1, scan: 1, batch: 100), + reliable: new Reliable(visibility: $visibility, heartbeat: 1, scan: 1, batch: $batch), ); $this->cleanup(); } @@ -336,6 +499,15 @@ private function stat(string $stat): int return (int) ($this->redis->get($this->statKey($stat)) ?: 0); } + private function encodedPayload(string $envelope): string + { + $marker = '"payload":'; + $offset = strpos($envelope, $marker); + $this->assertNotFalse($offset); + + return substr($envelope, $offset + \strlen($marker), -1); + } + /** @return array{pid: string, queue: string, timestamp: int, payload: array} */ private function pending(): array { diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php index c95e5adf0..de4703515 100644 --- a/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php @@ -179,6 +179,70 @@ function () use ($adapter, &$processed): void { $this->assertSame(0, $this->redis->lLen($this->key('queue'))); } + #[DataProvider('createFailureProvider')] + public function testCoroutineCreationFailureUnwindsWithoutProcessingClaimedJob(string $stage): void + { + $this->prepare(); + if ($stage !== 'recovery') { + $this->broker()->enqueue($this->queue, ['stage' => $stage]); + } + + $result = $this->runCreateFailure($stage); + + $this->assertSame(0, $result['code'], $result['error']); + $this->assertStringContainsString('handled', $result['output']); + $this->assertStringNotContainsString('processed', $result['output']); + if ($stage === 'recovery') { + $this->assertSame(0, $this->stat('processing')); + $this->assertSame(0, $this->redis->lLen($this->key('failed'))); + + return; + } + + $this->assertSame(1, $this->stat('processing')); + $this->assertSame(1, $this->redis->zCard($this->key('processing'))); + $this->assertSame(1, $this->redis->hLen($this->key('jobs'))); + $this->assertSame(0, $this->redis->lLen($this->key('failed'))); + $this->assertSame(0, $this->stat('failed')); + $this->assertSame(0, $this->stat('success')); + $this->assertSame(0, $this->redis->lLen($this->key('queue'))); + + $pid = $this->redis->zRange($this->key('processing'), 0, 0)[0] ?? null; + $this->assertIsString($pid); + $this->redis->zAdd($this->key('processing'), 0, $pid); + $broker = $this->broker(); + $claims = $broker->expired($this->queue, 1); + $this->assertCount(1, $claims); + $this->assertInstanceOf(Message::class, $broker->reclaim($this->queue, $claims[0])); + $recovered = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $recovered); + $broker->commit($this->queue, $recovered); + + $this->assertSame(0, $this->stat('processing')); + $this->assertSame(1, $this->stat('reclaimed')); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(0, $this->redis->zCard($this->key('processing'))); + $this->assertSame(0, $this->redis->hLen($this->key('jobs'))); + } + + public function testLegacyHandlerCreationFailureFallsBackSynchronously(): void + { + $this->prepare(); + $legacy = new Queue($this->queue->name, self::NAMESPACE); + $this->broker()->enqueue($legacy, ['stage' => 'legacy']); + + $result = $this->runCreateFailure('legacy'); + + $this->assertSame(0, $result['code'], $result['error']); + $this->assertStringContainsString('handled', $result['output']); + $this->assertStringContainsString('processed', $result['output']); + $this->assertSame(0, $this->stat('processing')); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(0, $this->stat('failed')); + $this->assertSame(0, $this->redis->lLen(self::NAMESPACE . '.processing.' . $this->queue->name)); + $this->assertSame(0, $this->redis->lLen(self::NAMESPACE . '.failed.' . $this->queue->name)); + } + /** @return iterable */ public static function capacityProvider(): iterable { @@ -186,6 +250,14 @@ public static function capacityProvider(): iterable yield 'three coroutines' => [3]; } + /** @return iterable */ + public static function createFailureProvider(): iterable + { + yield 'recovery' => ['recovery']; + yield 'handler' => ['handler']; + yield 'heartbeat' => ['heartbeat']; + } + protected function tearDown(): void { if (isset($this->redis)) { @@ -214,6 +286,56 @@ private function broker(): RedisBroker ); } + /** @return array{code: int, output: string, error: string} */ + private function runCreateFailure(string $stage): array + { + $process = proc_open( + [PHP_BINARY, __DIR__ . '/reliable-create-failure.php', $stage, $this->queue->name], + [ + 0 => ['pipe', 'r'], + 1 => ['pipe', 'w'], + 2 => ['pipe', 'w'], + ], + $pipes, + \dirname(__DIR__, 4), + ); + $this->assertIsResource($process); + fclose($pipes[0]); + stream_set_blocking($pipes[1], false); + stream_set_blocking($pipes[2], false); + + $code = null; + $deadline = microtime(true) + 2.0; + while (microtime(true) < $deadline) { + $status = proc_get_status($process); + if (!$status['running']) { + $code = $status['exitcode']; + break; + } + usleep(10_000); + } + if ($code === null) { + proc_terminate($process); + usleep(50_000); + $status = proc_get_status($process); + if ($status['running']) { + proc_terminate($process, 9); + } + $code = 124; + } + + $output = stream_get_contents($pipes[1]); + $error = stream_get_contents($pipes[2]); + fclose($pipes[1]); + fclose($pipes[2]); + $closed = proc_close($process); + if ($code === -1) { + $code = $closed; + } + + return ['code' => $code, 'output' => $output, 'error' => $error]; + } + private function cleanup(): void { $this->redis->del([ diff --git a/packages/queue/tests/Queue/E2E/Adapter/reliable-create-failure.php b/packages/queue/tests/Queue/E2E/Adapter/reliable-create-failure.php new file mode 100644 index 000000000..3f214c279 --- /dev/null +++ b/packages/queue/tests/Queue/E2E/Adapter/reliable-create-failure.php @@ -0,0 +1,55 @@ + 1, + 'handler' => 2, + 'heartbeat' => 3, + 'legacy' => 1, + default => throw new InvalidArgumentException("Unknown creation failure stage: {$stage}"), +}; +$reliable = new Reliable(visibility: 2, heartbeat: 1, scan: 1, batch: 100); +$broker = new RedisBroker( + new RedisConnection('127.0.0.1', 16379, connectTimeout: 1.0, readTimeout: 2.0), + new Locking(new RedisConnection('127.0.0.1', 16379, connectTimeout: 1.0, readTimeout: 2.0)), +); +$adapter = new Swoole( + $broker, + 1, + $queue, + 'reliable-swoole-tests', + maxCoroutines: 1, + reliable: $stage === 'legacy' ? null : $reliable, +); + +Coroutine::set([ + 'hook_flags' => SWOOLE_HOOK_ALL, + 'max_coroutine' => $maximum, +]); +Coroutine\run(function () use ($adapter): void { + try { + $adapter->consume( + static function () use ($adapter): void { + echo "processed\n"; + $adapter->stop(); + }, + static fn(): null => null, + static fn(): null => null, + ); + } catch (RuntimeException) { + } + + echo "handled\n"; +}); From c500e0a374c92b9bc045107d2eaeb4f57ac227ea Mon Sep 17 00:00:00 2001 From: Jake Barnby Date: Thu, 16 Jul 2026 01:20:23 +1200 Subject: [PATCH 3/4] (fix): preserve claims across queue shutdown --- packages/queue/src/Queue/Adapter/Swoole.php | 4 + packages/queue/src/Queue/Broker/Redis.php | 4 + .../Queue/E2E/Adapter/ReliableApiTest.php | 42 +++++ .../Queue/E2E/Adapter/ReliableSwooleTest.php | 145 ++++++++++++++++++ 4 files changed, 195 insertions(+) diff --git a/packages/queue/src/Queue/Adapter/Swoole.php b/packages/queue/src/Queue/Adapter/Swoole.php index 754ef80e2..f64641a68 100644 --- a/packages/queue/src/Queue/Adapter/Swoole.php +++ b/packages/queue/src/Queue/Adapter/Swoole.php @@ -202,6 +202,10 @@ private function consumeReliable( $slots->pop(); throw $error; } + if ($this->isStopped()) { + $slots->pop(); + break; + } if (!$message instanceof Message) { $slots->pop(); continue; diff --git a/packages/queue/src/Queue/Broker/Redis.php b/packages/queue/src/Queue/Broker/Redis.php index 9d453eebd..2c484d13e 100644 --- a/packages/queue/src/Queue/Broker/Redis.php +++ b/packages/queue/src/Queue/Broker/Redis.php @@ -413,6 +413,10 @@ private function receiveReliable(Queue $queue, int $timeout): ?Message continue; } + if ($this->isClosed()) { + return null; + } + if (!\is_array($result) || !isset($result[0])) { throw new \UnexpectedValueException('Reliable claim returned an invalid response.'); } diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php index 660d79176..d81c5fac0 100644 --- a/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableApiTest.php @@ -216,6 +216,48 @@ public function testReliableReceiveReconnectsAndResetsAfterSuccessfulEmptyEvalua $this->assertLessThanOrEqual(100, $failures[0][3]); $this->assertSame([[$queue, 1]], $successes); } + + public function testReliableReceiveDropsClaimCompletedWhileBrokerCloses(): void + { + $connection = new ClosingClaimConnection(); + $broker = new RedisBroker($connection, $connection); + $queue = new Queue('closed-claim', reliable: new Reliable()); + $connection->closeWith($broker->close(...)); + + $message = $broker->receive($queue, 0); + + $this->assertTrue($connection->claimed); + $this->assertNotInstanceOf(\Utopia\Queue\Message::class, $message); + } +} + +final class ClosingClaimConnection extends InMemoryConnection implements Atomic +{ + public bool $claimed = false; + + private ?\Closure $close = null; + + public function closeWith(\Closure $close): void + { + $this->close = $close; + } + + public function supportsAtomic(): bool + { + return true; + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + $this->claimed = true; + ($this->close ?? throw new \LogicException('Close callback is missing.'))(); + + return [ + 1, + '{"pid":"claimed","queue":"closed-claim","timestamp":1,"payload":{"value":true}}', + '1:1', + ]; + } } final class UnsupportedAtomicConnection extends InMemoryConnection implements Atomic diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php index de4703515..e6beb486d 100644 --- a/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php @@ -7,11 +7,15 @@ use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\TestCase; use Swoole\Coroutine; +use Swoole\Coroutine\Channel; use Swoole\Coroutine\WaitGroup; use Utopia\Queue\Adapter\Swoole; use Utopia\Queue\Broker\Redis as RedisBroker; +use Utopia\Queue\Claim; use Utopia\Queue\Connection\Locking; use Utopia\Queue\Connection\Redis as RedisConnection; +use Utopia\Queue\Consumer; +use Utopia\Queue\Consumer\Recoverable; use Utopia\Queue\Message; use Utopia\Queue\Option\Reliable; use Utopia\Queue\Queue; @@ -179,6 +183,80 @@ function () use ($adapter, &$processed): void { $this->assertSame(0, $this->redis->lLen($this->key('queue'))); } + public function testStopAfterConcurrentClaimLeavesMessageForRecovery(): void + { + $this->prepare(); + $publisher = $this->broker(); + $publisher->enqueue($this->queue, ['id' => 'A']); + $publisher->enqueue($this->queue, ['id' => 'B']); + $callbacks = []; + $acks = []; + $errors = []; + + Coroutine\run(function () use (&$callbacks, &$acks, &$errors): void { + $consumer = new CoordinatedRecoverableConsumer($this->broker()); + $adapter = new Swoole( + $consumer, + 1, + $this->queue->name, + self::NAMESPACE, + maxCoroutines: 2, + reliable: $this->queue->reliable, + ); + + $adapter->consume( + function (Message $message) use ($adapter, $consumer, &$callbacks): void { + $id = $message->getPayload()['id']; + $callbacks[] = $id; + if ($id !== 'A') { + return; + } + + try { + $consumer->waitForSecondClaim(); + } finally { + $adapter->stop(); + $consumer->releaseSecondClaim(); + } + }, + static function (Message $message) use (&$acks): void { + $acks[] = $message->getPayload()['id']; + }, + static function (Message $message, \Throwable $error) use (&$errors): void { + $errors[] = [$message->getPayload()['id'], $error->getMessage()]; + }, + ); + }); + + $this->assertSame(['A'], $callbacks); + $this->assertSame(['A'], $acks); + $this->assertSame([], $errors); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(1, $this->stat('processing')); + $this->assertSame(1, $this->redis->hLen($this->key('jobs'))); + $this->assertSame(1, $this->redis->zCard($this->key('processing'))); + $this->assertSame(0, $this->redis->lLen($this->key('queue'))); + + $pid = $this->redis->zRange($this->key('processing'), 0, 0)[0] ?? null; + $this->assertIsString($pid); + $this->redis->zAdd($this->key('processing'), 0, $pid); + $broker = $this->broker(); + $claims = $broker->expired($this->queue, 1); + $this->assertCount(1, $claims); + $this->assertInstanceOf(Message::class, $broker->reclaim($this->queue, $claims[0])); + $recovered = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $recovered); + $this->assertSame(['id' => 'B'], $recovered->getPayload()); + $broker->commit($this->queue, $recovered); + + $this->assertSame(2, $this->stat('success')); + $this->assertSame(1, $this->stat('reclaimed')); + $this->assertSame(0, $this->stat('processing')); + $this->assertSame(0, $this->redis->hLen($this->key('jobs'))); + $this->assertSame(0, $this->redis->zCard($this->key('processing'))); + $this->assertSame(0, $this->redis->lLen($this->key('queue'))); + } + #[DataProvider('createFailureProvider')] public function testCoroutineCreationFailureUnwindsWithoutProcessingClaimedJob(string $stage): void { @@ -373,3 +451,70 @@ private function stat(string $stat): int return (int) ($this->redis->get($this->statKey($stat)) ?: 0); } } + +final class CoordinatedRecoverableConsumer implements Consumer, Recoverable +{ + private readonly Channel $claimed; + private readonly Channel $release; + private int $receives = 0; + + public function __construct(private readonly RedisBroker $consumer) + { + $this->claimed = new Channel(1); + $this->release = new Channel(1); + } + + public function receive(Queue $queue, int $timeout): ?Message + { + $message = $this->consumer->receive($queue, $timeout); + $this->receives++; + if ($this->receives === 2 && $message instanceof Message) { + $this->claimed->push($message); + $this->release->pop(); + } + + return $message; + } + + public function commit(Queue $queue, Message $message): void + { + $this->consumer->commit($queue, $message); + } + + public function reject(Queue $queue, Message $message): void + { + $this->consumer->reject($queue, $message); + } + + public function close(): void + { + $this->consumer->close(); + } + + public function extend(Queue $queue, Message $message): bool + { + return $this->consumer->extend($queue, $message); + } + + public function expired(Queue $queue, int $limit): array + { + return $this->consumer->expired($queue, $limit); + } + + public function reclaim(Queue $queue, Claim $claim): ?Message + { + return $this->consumer->reclaim($queue, $claim); + } + + public function waitForSecondClaim(): void + { + if (!$this->claimed->pop(2.0) instanceof Message) { + throw new \RuntimeException('Second claim was not completed.'); + } + } + + public function releaseSecondClaim(): void + { + $this->release->push(true); + } +} From 70eaa97c49247cedca96a18e1fd4c3d9dd7a4d92 Mon Sep 17 00:00:00 2001 From: Jake Barnby Date: Thu, 16 Jul 2026 02:11:22 +1200 Subject: [PATCH 4/4] (fix): recover reliable Redis claims safely --- packages/queue/src/Queue/Broker/Redis.php | 250 ++++++++++++------ .../queue/src/Queue/Broker/Redis/Script.php | 2 +- .../queue/src/Queue/Exception/ClaimLost.php | 15 ++ .../Queue/E2E/Adapter/ReliableApiTest.php | 155 +++++++++++ .../Queue/E2E/Adapter/ReliableRedisTest.php | 106 +++++++- .../Queue/E2E/Adapter/ReliableSwooleTest.php | 196 +++++++++++++- packages/queue/tests/e2e.sh | 6 +- 7 files changed, 635 insertions(+), 95 deletions(-) create mode 100644 packages/queue/src/Queue/Exception/ClaimLost.php diff --git a/packages/queue/src/Queue/Broker/Redis.php b/packages/queue/src/Queue/Broker/Redis.php index 2c484d13e..b3419aa3d 100644 --- a/packages/queue/src/Queue/Broker/Redis.php +++ b/packages/queue/src/Queue/Broker/Redis.php @@ -8,6 +8,7 @@ use Utopia\Queue\Connection\Atomic; use Utopia\Queue\Consumer; use Utopia\Queue\Consumer\Recoverable; +use Utopia\Queue\Exception\ClaimLost; use Utopia\Queue\Message; use Utopia\Queue\Option\Reliable; use Utopia\Queue\Publisher; @@ -16,12 +17,19 @@ class Redis implements Publisher, Consumer, Recoverable { private const int POP_TIMEOUT = 2; + private const int COMMAND_MAX_ATTEMPTS = 3; + private const int COMMAND_BACKOFF_MS = 100; + private const int COMMAND_MAX_BACKOFF_MS = 500; private const int RECONNECT_BACKOFF_MS = 100; private const int RECONNECT_MAX_BACKOFF_MS = 5_000; private bool $closed = false; private int $reconnectAttempt = 0; private int $reconnectBackoffMs = self::RECONNECT_BACKOFF_MS; + + /** @var \WeakMap */ + private \WeakMap $settled; + /** * @var (callable(Queue, \Throwable, int, int): void)|null */ @@ -36,7 +44,9 @@ public function __construct( private readonly Connection $receive, // Acks and publishing; wrap in Locking when shared by coroutines. private readonly Connection $commands, - ) {} + ) { + $this->settled = new \WeakMap(); + } public function setReconnectCallback(?callable $callback): self { @@ -184,17 +194,19 @@ private function triggerReconnectSuccessCallback(Queue $queue, int $attempts): v public function enqueue(Queue $queue, array $payload, bool $priority = false): bool { if ($queue->reliable instanceof Reliable) { - $atomic = $this->atomic($this->commands); - $result = $atomic->evaluate( - Script::ENQUEUE, - [ - $this->pendingKey($queue), - $this->pid(), - $queue->name, - json_encode($payload, JSON_THROW_ON_ERROR), - $priority ? '1' : '0', - ], - 1, + $result = $this->command( + fn(): mixed => $this->atomic($this->commands)->evaluate( + Script::ENQUEUE, + [ + $this->pendingKey($queue), + $this->pid(), + $queue->name, + json_encode($payload, JSON_THROW_ON_ERROR), + $priority ? '1' : '0', + ], + 1, + ), + retry: false, ); return (int) $result === 1; @@ -285,6 +297,13 @@ public function getQueueSize(Queue $queue, bool $failedJobs = false): int $queueName = "{$queue->namespace}.failed.{$queue->name}"; } } + if ($queue->reliable instanceof Reliable) { + return $this->command( + fn(): int => $this->commands->listSize($queueName), + retry: true, + ); + } + return $this->commands->listSize($queueName); } @@ -292,16 +311,19 @@ public function extend(Queue $queue, Message $message): bool { $reliable = $this->reliable($queue); $claimedAt = $this->claimedAt($message); - $result = $this->atomic($this->commands)->evaluate( - Script::EXTEND, - [ - $this->atomicKey($queue, 'jobs'), - $this->atomicKey($queue, 'processing'), - $message->getPid(), - $claimedAt, - (string) $reliable->visibility, - ], - 2, + $result = $this->command( + fn(): mixed => $this->atomic($this->commands)->evaluate( + Script::EXTEND, + [ + $this->atomicKey($queue, 'jobs'), + $this->atomicKey($queue, 'processing'), + $message->getPid(), + $claimedAt, + (string) $reliable->visibility, + ], + 2, + ), + retry: true, ); return (int) $result === 1; @@ -314,14 +336,17 @@ public function expired(Queue $queue, int $limit): array return []; } - $result = $this->atomic($this->commands)->evaluate( - Script::EXPIRED, - [ - $this->atomicKey($queue, 'processing'), - $this->atomicKey($queue, 'jobs'), - (string) $limit, - ], - 2, + $result = $this->command( + fn(): mixed => $this->atomic($this->commands)->evaluate( + Script::EXPIRED, + [ + $this->atomicKey($queue, 'processing'), + $this->atomicKey($queue, 'jobs'), + (string) $limit, + ], + 2, + ), + retry: true, ); if (!\is_array($result)) { throw new \UnexpectedValueException('Reliable expired scan returned an invalid response.'); @@ -344,21 +369,24 @@ public function expired(Queue $queue, int $limit): array public function reclaim(Queue $queue, Claim $claim): ?Message { $this->reliable($queue); - $result = $this->atomic($this->commands)->evaluate( - Script::RECLAIM, - [ - $this->atomicKey($queue, 'processing'), - $this->atomicKey($queue, 'jobs'), - $this->pendingKey($queue), - $this->atomicKey($queue, 'quarantine'), - $this->statKey($queue, 'processing'), - $this->statKey($queue, 'reclaimed'), - $this->statKey($queue, 'quarantined'), - $claim->pid, - $claim->claimedAt ?? '', - $this->pid(), - ], - 7, + $result = $this->command( + fn(): mixed => $this->atomic($this->commands)->evaluate( + Script::RECLAIM, + [ + $this->atomicKey($queue, 'processing'), + $this->atomicKey($queue, 'jobs'), + $this->pendingKey($queue), + $this->atomicKey($queue, 'quarantine'), + $this->statKey($queue, 'processing'), + $this->statKey($queue, 'reclaimed'), + $this->statKey($queue, 'quarantined'), + $claim->pid, + $claim->claimedAt ?? '', + $this->pid(), + ], + 7, + ), + retry: false, ); if (!\is_array($result) || !isset($result[0])) { @@ -448,36 +476,62 @@ private function receiveReliable(Queue $queue, int $timeout): ?Message private function commitReliable(Queue $queue, Message $message): void { $this->reliable($queue); - $this->atomic($this->commands)->evaluate( - Script::COMMIT, - [ - $this->atomicKey($queue, 'jobs'), - $this->atomicKey($queue, 'processing'), - $this->statKey($queue, 'success'), - $this->statKey($queue, 'processing'), - $message->getPid(), - $this->claimedAt($message), - ], - 4, + if (isset($this->settled[$message])) { + return; + } + + $claimedAt = $this->claimedAt($message); + $result = $this->command( + fn(): mixed => $this->atomic($this->commands)->evaluate( + Script::COMMIT, + [ + $this->atomicKey($queue, 'jobs'), + $this->atomicKey($queue, 'processing'), + $this->statKey($queue, 'success'), + $this->statKey($queue, 'processing'), + $message->getPid(), + $claimedAt, + ], + 4, + ), + retry: false, ); + if ((int) $result !== 1) { + throw new ClaimLost($message->getPid(), $claimedAt); + } + + $this->settled[$message] = true; } private function rejectReliable(Queue $queue, Message $message): void { $this->reliable($queue); - $this->atomic($this->commands)->evaluate( - Script::REJECT, - [ - $this->atomicKey($queue, 'jobs'), - $this->atomicKey($queue, 'processing'), - $this->atomicKey($queue, 'failed'), - $this->statKey($queue, 'failed'), - $this->statKey($queue, 'processing'), - $message->getPid(), - $this->claimedAt($message), - ], - 5, + if (isset($this->settled[$message])) { + return; + } + + $claimedAt = $this->claimedAt($message); + $result = $this->command( + fn(): mixed => $this->atomic($this->commands)->evaluate( + Script::REJECT, + [ + $this->atomicKey($queue, 'jobs'), + $this->atomicKey($queue, 'processing'), + $this->atomicKey($queue, 'failed'), + $this->statKey($queue, 'failed'), + $this->statKey($queue, 'processing'), + $message->getPid(), + $claimedAt, + ], + 5, + ), + retry: false, ); + if ((int) $result !== 1) { + throw new ClaimLost($message->getPid(), $claimedAt); + } + + $this->settled[$message] = true; } private function retryReliable(Queue $queue, ?int $limit): void @@ -487,18 +541,21 @@ private function retryReliable(Queue $queue, ?int $limit): void $processed = 0; while ($limit === null || $processed < $limit) { - $result = $atomic->evaluate( - Script::RETRY, - [ - $this->atomicKey($queue, 'failed'), - $this->atomicKey($queue, 'jobs'), - $this->pendingKey($queue), - $this->atomicKey($queue, 'quarantine'), - $this->statKey($queue, 'retried'), - $this->statKey($queue, 'quarantined'), - $this->pid(), - ], - 6, + $result = $this->command( + fn(): mixed => $atomic->evaluate( + Script::RETRY, + [ + $this->atomicKey($queue, 'failed'), + $this->atomicKey($queue, 'jobs'), + $this->pendingKey($queue), + $this->atomicKey($queue, 'quarantine'), + $this->statKey($queue, 'retried'), + $this->statKey($queue, 'quarantined'), + $this->pid(), + ], + 6, + ), + retry: false, ); if (!\is_array($result) || !isset($result[0])) { throw new \UnexpectedValueException('Reliable retry returned an invalid response.'); @@ -569,6 +626,39 @@ private function pid(): string return uniqid(more_entropy: true); } + /** + * Reset a failed shared command connection before either propagating an + * ambiguous mutation or replaying an idempotent/read-only operation. + * + * @template T + * @param callable(): T $command + * @return T + */ + private function command(callable $command, bool $retry): mixed + { + $attempt = 0; + $backoffMs = self::COMMAND_BACKOFF_MS; + + while (true) { + try { + return $command(); + } catch (\RedisException|\RedisClusterException $error) { + $attempt++; + try { + $this->commands->close(); + } catch (\Throwable) { + } + + if (!$retry || $attempt >= self::COMMAND_MAX_ATTEMPTS) { + throw $error; + } + + usleep(mt_rand(0, $backoffMs) * 1000); + $backoffMs = min(self::COMMAND_MAX_BACKOFF_MS, $backoffMs * 2); + } + } + } + private function reconnectSucceeded(Queue $queue): void { if ($this->reconnectAttempt > 0) { diff --git a/packages/queue/src/Queue/Broker/Redis/Script.php b/packages/queue/src/Queue/Broker/Redis/Script.php index d5bb8c06b..1a704323f 100644 --- a/packages/queue/src/Queue/Broker/Redis/Script.php +++ b/packages/queue/src/Queue/Broker/Redis/Script.php @@ -184,7 +184,7 @@ final class Script public const string EXPIRED = <<<'LUA' local now = redis.call('TIME') local micros = (tonumber(now[1]) * 1000000) + tonumber(now[2]) -local pids = redis.call('ZRANGEBYSCORE', KEYS[1], '-inf', micros, 'LIMIT', 0, tonumber(ARGV[1])) +local pids = redis.call('ZREVRANGEBYSCORE', KEYS[1], micros, '-inf', 'LIMIT', 0, tonumber(ARGV[1])) local result = {} for _, pid in ipairs(pids) do diff --git a/packages/queue/src/Queue/Exception/ClaimLost.php b/packages/queue/src/Queue/Exception/ClaimLost.php new file mode 100644 index 000000000..1edb0067d --- /dev/null +++ b/packages/queue/src/Queue/Exception/ClaimLost.php @@ -0,0 +1,15 @@ +assertTrue($connection->claimed); $this->assertNotInstanceOf(\Utopia\Queue\Message::class, $message); } + + public function testReliableCommitResetsAfterTransportFailureAndLaterAckIsIdempotent(): void + { + $connection = new RecoveringCommandConnection(1, new \RedisException('Redis is unavailable.')); + $locking = new Locking($connection); + $broker = new RedisBroker($locking, $locking); + $queue = new Queue('command-commit', reliable: new Reliable()); + $message = $this->claimedMessage($queue); + + try { + $broker->commit($queue, $message); + $this->fail('An acknowledgement with an unknown transport outcome must not report success.'); + } catch (\RedisException) { + } + + $this->assertSame(1, $connection->evaluations); + $this->assertSame(1, $connection->closes); + + $broker->commit($queue, $message); + $broker->commit($queue, $message); + + $this->assertSame(2, $connection->evaluations); + $this->assertSame(1, $connection->closes); + } + + public function testReliableRejectResetsAfterTransportFailureWithoutRetry(): void + { + $connection = new RecoveringCommandConnection(1, new \RedisException('Redis is unavailable.')); + $locking = new Locking($connection); + $broker = new RedisBroker($locking, $locking); + $queue = new Queue('command-reject', reliable: new Reliable()); + $message = $this->claimedMessage($queue); + + try { + $broker->reject($queue, $message); + $this->fail('A rejection with an unknown transport outcome must not be retried.'); + } catch (\RedisException) { + } + + $this->assertSame(1, $connection->evaluations); + $this->assertSame(1, $connection->closes); + + $broker->reject($queue, $message); + $broker->reject($queue, $message); + $broker->commit($queue, $message); + + $this->assertSame(2, $connection->evaluations); + $this->assertSame(1, $connection->closes); + } + + #[\PHPUnit\Framework\Attributes\DataProvider('retryableCommandProvider')] + public function testHeartbeatAndRecoveryCommandsResetAndRetryTransportFailures( + string $operation, + mixed $result, + \Throwable $failure, + ): void { + $connection = new RecoveringCommandConnection($result, $failure); + $locking = new Locking($connection); + $broker = new RedisBroker($locking, $locking); + $queue = new Queue("command-{$operation}", reliable: new Reliable()); + $message = $this->claimedMessage($queue); + + $actual = match ($operation) { + 'extend' => $broker->extend($queue, $message), + 'expired' => $broker->expired($queue, 1), + default => throw new \LogicException("Unknown operation: {$operation}"), + }; + + $expected = match ($operation) { + 'extend' => true, + 'expired' => [], + }; + $this->assertSame($expected, $actual); + $this->assertSame(2, $connection->evaluations); + $this->assertSame(1, $connection->closes); + } + + /** @return iterable */ + public static function retryableCommandProvider(): iterable + { + yield 'heartbeat Redis failure' => ['extend', 1, new \RedisException('Redis is unavailable.')]; + yield 'expired scan cluster failure' => ['expired', [], new \RedisClusterException('Redis Cluster is unavailable.')]; + } + + #[\PHPUnit\Framework\Attributes\DataProvider('unsafeCommandProvider')] + public function testNonIdempotentReliableCommandsResetWithoutRetry(string $operation): void + { + $connection = new RecoveringCommandConnection(1, new \RedisException('Redis is unavailable.')); + $locking = new Locking($connection); + $broker = new RedisBroker($locking, $locking); + $queue = new Queue("command-{$operation}", reliable: new Reliable()); + + try { + match ($operation) { + 'enqueue' => $broker->enqueue($queue, ['value' => true]), + 'retry' => $broker->retry($queue), + 'reclaim' => $broker->reclaim($queue, new Claim('pid', '1:1')), + default => throw new \LogicException("Unknown operation: {$operation}"), + }; + $this->fail('A command with an unknown transport outcome must not be retried.'); + } catch (\RedisException) { + } + + $this->assertSame(1, $connection->evaluations); + $this->assertSame(1, $connection->closes); + } + + /** @return iterable */ + public static function unsafeCommandProvider(): iterable + { + yield 'enqueue' => ['enqueue']; + yield 'retry' => ['retry']; + yield 'reclaim' => ['reclaim']; + } + + private function claimedMessage(Queue $queue): Message + { + return new Message([ + 'pid' => 'pid', + 'queue' => $queue->name, + 'timestamp' => 1, + 'payload' => [], + ], '1:1'); + } } final class ClosingClaimConnection extends InMemoryConnection implements Atomic @@ -392,3 +516,34 @@ public function close(): void $this->closes++; } } + +final class RecoveringCommandConnection extends InMemoryConnection implements Atomic +{ + public int $closes = 0; + public int $evaluations = 0; + + public function __construct( + private readonly mixed $result, + private readonly \Throwable $failure, + ) {} + + public function supportsAtomic(): bool + { + return true; + } + + public function evaluate(string $script, array $arguments = [], int $keyCount = 0): mixed + { + if (++$this->evaluations === 1) { + throw $this->failure; + } + + return $this->result; + } + + #[\Override] + public function close(): void + { + $this->closes++; + } +} diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php index f340678cb..dad256dee 100644 --- a/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableRedisTest.php @@ -6,9 +6,11 @@ use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\TestCase; +use Utopia\Queue\Adapter; use Utopia\Queue\Broker\Redis as RedisBroker; use Utopia\Queue\Connection\Locking; use Utopia\Queue\Connection\Redis as RedisConnection; +use Utopia\Queue\Exception\ClaimLost; use Utopia\Queue\Message; use Utopia\Queue\Option\Reliable; use Utopia\Queue\Queue; @@ -138,8 +140,20 @@ public function testReclaimPreservesEnqueueDataUsesNewPidAndPlainEnvelope(int $p $this->assertSame(0, $this->redis->zCard($this->key('processing'))); $this->assertSame(1, $this->stat('reclaimed')); - $broker->commit($this->queue, $message); - $broker->reject($this->queue, $message); + try { + $broker->commit($this->queue, $message); + $this->fail('Committing a reclaimed claim must report that ownership was lost.'); + } catch (ClaimLost $error) { + $this->assertSame($message->getPid(), $error->pid); + $this->assertSame($message->getClaimedAt(), $error->claimedAt); + } + try { + $broker->reject($this->queue, $message); + $this->fail('Rejecting a reclaimed claim must report that ownership was lost.'); + } catch (ClaimLost $error) { + $this->assertSame($message->getPid(), $error->pid); + $this->assertSame($message->getClaimedAt(), $error->claimedAt); + } $this->assertSame(0, $this->stat('success')); $this->assertSame(0, $this->stat('failed')); $this->assertSame(0, $this->stat('processing')); @@ -154,6 +168,56 @@ public function testReclaimPreservesEnqueueDataUsesNewPidAndPlainEnvelope(int $p $broker->commit($this->queue, $claimedAgain); } + #[DataProvider('serverProvider')] + public function testForcedReclaimReportsLostClaimWithoutCallingSuccess(int $port): void + { + $this->prepare($port, visibility: 3); + $broker = $this->broker($port); + $broker->enqueue($this->queue, ['slow' => true]); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $adapter = new ProcessingAdapter( + $broker, + 1, + $this->queue->name, + self::NAMESPACE, + reliable: $this->queue->reliable, + ); + $replacement = null; + $successes = []; + $errors = []; + + $adapter->processMessage( + $message, + function (Message $message) use ($broker, &$replacement): void { + usleep(50_000); + $this->redis->zAdd($this->key('processing'), 0, $message->getPid()); + $claim = $broker->expired($this->queue, 1)[0]; + $replacement = $broker->reclaim($this->queue, $claim); + }, + static function (Message $message) use (&$successes): void { + $successes[] = $message; + }, + static function (Message $message, \Throwable $error) use (&$errors): void { + $errors[] = [$message, $error]; + }, + ); + + $this->assertInstanceOf(Message::class, $replacement); + $this->assertSame([], $successes); + $this->assertCount(1, $errors); + $this->assertInstanceOf(ClaimLost::class, $errors[0][1]); + $this->assertSame(0, $this->stat('success')); + $this->assertSame(1, $this->stat('reclaimed')); + $this->assertSame(1, $this->redis->lLen($this->key('queue'))); + + $recovered = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $recovered); + $broker->commit($this->queue, $recovered); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(0, $this->stat('processing')); + } + #[DataProvider('serverProvider')] public function testHeartbeatRenewsOnlyLeaseAndOriginalTokenStillCommits(int $port): void { @@ -194,7 +258,11 @@ public function testRetryAndReclaimRacesHaveOneWinner(int $port): void $claim = $broker->expired($this->queue, 1)[0]; $this->assertInstanceOf(Message::class, $broker->reclaim($this->queue, $claim)); - $broker->reject($this->queue, $message); + try { + $broker->reject($this->queue, $message); + $this->fail('The reclaim winner must invalidate the original rejection token.'); + } catch (ClaimLost) { + } $broker->retry($this->queue); $this->assertSame(1, $broker->getQueueSize($this->queue)); $this->assertSame(0, $broker->getQueueSize($this->queue, failedJobs: true)); @@ -527,3 +595,35 @@ private function serverSeconds(): int return (int) ltrim((string) $time[0], ':'); } } + +final class ProcessingAdapter extends Adapter +{ + public function processMessage( + Message $message, + callable $messageCallback, + callable $successCallback, + callable $errorCallback, + ): void { + $this->process($message, $messageCallback, $successCallback, $errorCallback); + } + + public function start(): self + { + return $this; + } + + public function stop(): self + { + return $this; + } + + public function workerStart(callable $callback): self + { + return $this; + } + + public function workerStop(callable $callback): self + { + return $this; + } +} diff --git a/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php index e6beb486d..7b5b06557 100644 --- a/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php +++ b/packages/queue/tests/Queue/E2E/Adapter/ReliableSwooleTest.php @@ -5,6 +5,7 @@ namespace Tests\E2E\Adapter; use PHPUnit\Framework\Attributes\DataProvider; +use PHPUnit\Framework\Attributes\Group; use PHPUnit\Framework\TestCase; use Swoole\Coroutine; use Swoole\Coroutine\Channel; @@ -111,6 +112,65 @@ function (Message $message) use ($adapter, &$processed, &$token): void { $this->assertSame(0, $this->stat('processing')); } + #[Group('redis-restart')] + public function testSharedCommandConnectionRecoversAfterRealRedisRestart(): void + { + $this->prepare(visibility: 2, heartbeat: 1, scan: 1, batch: 10); + $broker = $this->broker(); + $broker->enqueue($this->queue, ['restart' => true]); + $message = $broker->receive($this->queue, 1); + $this->assertInstanceOf(Message::class, $message); + $this->assertTrue($broker->extend($this->queue, $message)); + $this->redis->zAdd($this->key('processing'), 0, $message->getPid()); + $redisStopped = false; + + try { + $this->redis->close(); + $this->composeRedis('stop'); + $redisStopped = true; + + $started = microtime(true); + try { + $broker->extend($this->queue, $message); + $this->fail('Heartbeat must report that Redis stayed unavailable past its bounded retries.'); + } catch (\RedisException) { + } + $this->assertLessThan(10.0, microtime(true) - $started); + + try { + $broker->commit($this->queue, $message); + $this->fail('An acknowledgement attempted during the outage must remain uncertain.'); + } catch (\RedisException) { + } + + $this->composeRedis('start'); + $redisStopped = false; + $this->reconnectRedis(); + + $claims = $broker->expired($this->queue, 10); + $this->assertCount(1, $claims); + $replacement = $broker->reclaim($this->queue, $claims[0]); + $this->assertInstanceOf(Message::class, $replacement); + + $recovered = $broker->receive($this->queue, 2); + $this->assertInstanceOf(Message::class, $recovered); + $this->assertNotSame($message->getPid(), $recovered->getPid()); + $broker->commit($this->queue, $recovered); + } finally { + if ($redisStopped) { + $this->composeRedis('start'); + } + if (!isset($this->redis) || !$this->redis->isConnected()) { + $this->reconnectRedis(); + } + } + + $this->assertSame(1, $this->stat('reclaimed')); + $this->assertSame(1, $this->stat('success')); + $this->assertSame(0, $this->stat('failed')); + $this->assertSame(0, $this->stat('processing')); + } + public function testConcurrentRecoveryLoopsProduceOneReplacement(): void { $this->prepare(); @@ -148,17 +208,23 @@ public function testRecoveryDrainsConsecutiveBoundedBatchesInOneScan(): void { $this->prepare(scan: 1, batch: 2); $broker = $this->broker(); - for ($index = 0; $index < 5; $index++) { - $broker->enqueue($this->queue, ['index' => $index]); + $expired = []; + foreach (['A', 'B', 'C'] as $score => $order) { + $broker->enqueue($this->queue, ['order' => $order]); $message = $broker->receive($this->queue, 1); $this->assertInstanceOf(Message::class, $message); - $this->redis->zAdd($this->key('processing'), 0, $message->getPid()); + $expired[$order] = $message->getPid(); + $this->redis->zAdd($this->key('processing'), $score + 1, $message->getPid()); } - $processed = 0; + foreach (['D', 'E'] as $order) { + $broker->enqueue($this->queue, ['order' => $order]); + } + $consumer = new RecordingRecoverableConsumer($this->broker()); + $processed = []; - Coroutine\run(function () use (&$processed): void { + Coroutine\run(function () use ($consumer, &$processed): void { $adapter = new Swoole( - $this->broker(), + $consumer, 1, $this->queue->name, self::NAMESPACE, @@ -166,8 +232,9 @@ public function testRecoveryDrainsConsecutiveBoundedBatchesInOneScan(): void reliable: $this->queue->reliable, ); $adapter->consume( - function () use ($adapter, &$processed): void { - if (++$processed === 5) { + function (Message $message) use ($adapter, &$processed): void { + $processed[] = $message->getPayload()['order']; + if (\count($processed) === 5) { $adapter->stop(); } }, @@ -176,8 +243,9 @@ function () use ($adapter, &$processed): void { ); }); - $this->assertSame(5, $processed); - $this->assertSame(5, $this->stat('reclaimed')); + $this->assertSame([$expired['C'], $expired['B'], $expired['A']], $consumer->reclaimOrder); + $this->assertSame(['A', 'B', 'C', 'D', 'E'], $processed); + $this->assertSame(3, $this->stat('reclaimed')); $this->assertSame(5, $this->stat('success')); $this->assertSame(0, $this->stat('processing')); $this->assertSame(0, $this->redis->lLen($this->key('queue'))); @@ -414,6 +482,55 @@ private function runCreateFailure(string $stage): array return ['code' => $code, 'output' => $output, 'error' => $error]; } + private function composeRedis(string $action): void + { + $process = proc_open( + ['docker', 'compose', '-f', \dirname(__DIR__, 4) . '/docker-compose.yml', $action, 'redis'], + [ + 0 => ['pipe', 'r'], + 1 => ['pipe', 'w'], + 2 => ['pipe', 'w'], + ], + $pipes, + \dirname(__DIR__, 4), + ); + if (!\is_resource($process)) { + throw new \RuntimeException("Failed to run docker compose {$action} redis."); + } + fclose($pipes[0]); + $output = stream_get_contents($pipes[1]); + $error = stream_get_contents($pipes[2]); + fclose($pipes[1]); + fclose($pipes[2]); + $code = proc_close($process); + if ($code !== 0) { + throw new \RuntimeException("docker compose {$action} redis failed: {$output}{$error}"); + } + } + + private function reconnectRedis(): void + { + $deadline = microtime(true) + 10.0; + do { + $redis = new \Redis(); + try { + $redis->connect('127.0.0.1', self::PORT, 0.2); + $redis->ping(); + $this->redis = $redis; + + return; + } catch (\RedisException) { + try { + $redis->close(); + } catch (\Throwable) { + } + usleep(100_000); + } + } while (microtime(true) < $deadline); + + throw new \RuntimeException('Redis did not become ready after restart.'); + } + private function cleanup(): void { $this->redis->del([ @@ -518,3 +635,62 @@ public function releaseSecondClaim(): void $this->release->push(true); } } + +final class RecordingRecoverableConsumer implements Consumer, Recoverable +{ + private int $reclaims = 0; + + /** @var list */ + public array $reclaimOrder = []; + + public function __construct(private readonly RedisBroker $consumer) {} + + public function receive(Queue $queue, int $timeout): ?Message + { + $deadline = microtime(true) + 3.0; + while ($this->reclaims < 3 && microtime(true) < $deadline) { + Coroutine::sleep(0.01); + } + if ($this->reclaims < 3) { + throw new \RuntimeException('Recovery did not reclaim all expired messages before the bounded wait elapsed.'); + } + + return $this->consumer->receive($queue, $timeout); + } + + public function commit(Queue $queue, Message $message): void + { + $this->consumer->commit($queue, $message); + } + + public function reject(Queue $queue, Message $message): void + { + $this->consumer->reject($queue, $message); + } + + public function close(): void + { + $this->consumer->close(); + } + + public function extend(Queue $queue, Message $message): bool + { + return $this->consumer->extend($queue, $message); + } + + public function expired(Queue $queue, int $limit): array + { + return $this->consumer->expired($queue, $limit); + } + + public function reclaim(Queue $queue, Claim $claim): ?Message + { + $message = $this->consumer->reclaim($queue, $claim); + if ($message instanceof Message) { + $this->reclaimOrder[] = $claim->pid; + $this->reclaims++; + } + + return $message; + } +} diff --git a/packages/queue/tests/e2e.sh b/packages/queue/tests/e2e.sh index d10185a74..6fcea02df 100755 --- a/packages/queue/tests/e2e.sh +++ b/packages/queue/tests/e2e.sh @@ -4,6 +4,10 @@ set -e cd "$(dirname "$0")/.." +# This test intentionally restarts Redis. Run it before long-lived workers so +# their sockets cannot be invalidated for the rest of the suite. +phpunit --testsuite e2e --group redis-restart + php tests/Queue/servers/Swoole/worker.php & SWOOLE=$! php tests/Queue/servers/SwooleRedisCluster/worker.php & CLUSTER=$! php tests/Queue/servers/Workerman/worker.php start & WORKERMAN=$! @@ -16,4 +20,4 @@ trap cleanup EXIT INT TERM sleep 3 -phpunit --testsuite e2e +phpunit --testsuite e2e --exclude-group redis-restart