diff --git a/CHANGELOG.md b/CHANGELOG.md index 4c7aaba..d2d65ba 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,11 @@ All notable changes to Torque will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [0.16.1] - 2026-08-28 + +### Fixed +- **Multi-stream `XREADGROUP` stranded every message but the first.** A worker polls all eligible streams in one `XREADGROUP ... COUNT 1 STREAMS a b c`, and Redis returns up to `COUNT` entries *per stream*. `parseXreadgroupResponse()` returned only the first stream's entry, so the other entries were delivered to the consumer (in its PEL) and never handed to a Fiber; they waited until `retry_after` let another worker steal them (30 minutes on scrpr's `slowSync`, 2 hours on `slowScrpr`) and, if that worker was busy on a hot stream, were stranded again. Every delivered entry is now parsed and queued in a per-worker prefetch buffer that `readNextMessage()` drains before issuing the next read (same for the periodic own-PEL read). No config or data changes. + ## [0.16.0] - 2026-08-28 ### Added diff --git a/src/Worker/WorkerProcess.php b/src/Worker/WorkerProcess.php index e55a7b0..3ebd0f7 100644 --- a/src/Worker/WorkerProcess.php +++ b/src/Worker/WorkerProcess.php @@ -87,6 +87,17 @@ final class WorkerProcess */ private array $streamActive = []; + /** + * Messages already delivered to this consumer by a multi-stream XREADGROUP + * but not yet handed to a Fiber. Redis returns up to COUNT entries per + * stream, so one read can deliver more than one message; anything beyond + * the first is queued here and served before the next XREADGROUP, otherwise + * it would sit in the PEL until retry_after lets another worker steal it. + * + * @var list + */ + private array $prefetched = []; + /** Limits — set once in run(), checked by timers to know when to cancel. */ private int $maxJobs = 10_000; @@ -618,46 +629,66 @@ public function releaseConsumer( */ private function parseXreadgroupResponse(mixed $result): ?array { - if ($result === null) { + $messages = $this->parseXreadgroupMessages($result); + + if ($messages === []) { return null; } + $first = array_shift($messages); + + foreach ($messages as $message) { + $this->prefetched[] = $message; + } + + return $first; + } + + /** + * Every message in an XREADGROUP reply, in stream order. + * + * @return list + */ + private function parseXreadgroupMessages(mixed $result): array + { + if ($result === null || ! is_array($result)) { + return []; + } + + $parsed = []; + foreach ($result as $streamData) { if ($streamData === null) { continue; } $streamKey = (string) $streamData[0]; - $messages = $streamData[1] ?? []; - if ($messages === []) { - continue; - } + foreach ($streamData[1] ?? [] as $message) { + $messageId = (string) $message[0]; + $fields = $message[1]; - $message = $messages[0]; - $messageId = (string) $message[0]; - $fields = $message[1]; + $payload = null; + for ($i = 0, $count = count($fields); $i < $count; $i += 2) { + if ((string) $fields[$i] === 'payload') { + $payload = (string) $fields[$i + 1]; + break; + } + } - $payload = null; - for ($i = 0, $count = count($fields); $i < $count; $i += 2) { - if ((string) $fields[$i] === 'payload') { - $payload = (string) $fields[$i + 1]; - break; + if ($payload === null) { + continue; } - } - if ($payload === null) { - continue; + $parsed[] = [ + 'stream' => $streamKey, + 'id' => $messageId, + 'payload' => $payload, + ]; } - - return [ - 'stream' => $streamKey, - 'id' => $messageId, - 'payload' => $payload, - ]; } - return null; + return $parsed; } /** @@ -966,6 +997,10 @@ private function readNextMessage( string $consumerGroup, \Closure $buildStreamKey, ): ?array { + if ($this->prefetched !== []) { + return array_shift($this->prefetched); + } + $args = ['GROUP', $consumerGroup, $this->consumerId, 'COUNT', '1', 'STREAMS']; foreach ($queues as $queue) { diff --git a/tests/Unit/Worker/WorkerPrefetchTest.php b/tests/Unit/Worker/WorkerPrefetchTest.php new file mode 100644 index 0000000..c1a7864 --- /dev/null +++ b/tests/Unit/Worker/WorkerPrefetchTest.php @@ -0,0 +1,83 @@ + ['uri' => 'redis://127.0.0.1:6379'], 'queues' => ['default', 'slowSync']]); +} + +function parseResponse(WorkerProcess $worker, mixed $result): ?array +{ + return (new ReflectionMethod(WorkerProcess::class, 'parseXreadgroupResponse'))->invoke($worker, $result); +} + +function prefetched(WorkerProcess $worker): array +{ + return (new ReflectionProperty(WorkerProcess::class, 'prefetched'))->getValue($worker); +} + +it('returns the first delivered message and buffers the rest', function () { + $worker = makeWorker(); + + $first = parseResponse($worker, [ + ['torque:default', [['1-0', ['payload', 'a']]]], + ['torque:slowSync', [['2-0', ['payload', 'b']], ['3-0', ['payload', 'c']]]], + ]); + + expect($first)->toBe(['stream' => 'torque:default', 'id' => '1-0', 'payload' => 'a']) + ->and(prefetched($worker))->toBe([ + ['stream' => 'torque:slowSync', 'id' => '2-0', 'payload' => 'b'], + ['stream' => 'torque:slowSync', 'id' => '3-0', 'payload' => 'c'], + ]); +}); + +it('skips null streams and entries without a payload field', function () { + $worker = makeWorker(); + + $first = parseResponse($worker, [ + null, + ['torque:default', [['1-0', ['other', 'x']]]], + ['torque:slowSync', [['2-0', ['payload', 'b']]]], + ]); + + expect($first)->toBe(['stream' => 'torque:slowSync', 'id' => '2-0', 'payload' => 'b']) + ->and(prefetched($worker))->toBe([]); +}); + +it('returns null for an empty reply and leaves the buffer untouched', function () { + $worker = makeWorker(); + + expect(parseResponse($worker, null))->toBeNull() + ->and(parseResponse($worker, []))->toBeNull() + ->and(prefetched($worker))->toBe([]); +}); + +it('serves buffered messages before reading from Redis again', function () { + $worker = makeWorker(); + parseResponse($worker, [ + ['torque:default', [['1-0', ['payload', 'a']]]], + ['torque:slowSync', [['2-0', ['payload', 'b']]]], + ]); + + // With a non-empty buffer readNextMessage returns before touching Redis, + // so a client on a port nothing listens on proves the read is skipped. + $redis = createRedisClient('redis://127.0.0.1:1'); + + $next = (new ReflectionMethod(WorkerProcess::class, 'readNextMessage')) + ->invoke($worker, $redis, ['default', 'slowSync'], 'torque:', 'torque', fn (string $q) => 'torque:'.$q); + + expect($next)->toBe(['stream' => 'torque:slowSync', 'id' => '2-0', 'payload' => 'b']) + ->and(prefetched($worker))->toBe([]); +});