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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
81 changes: 58 additions & 23 deletions src/Worker/WorkerProcess.php
Original file line number Diff line number Diff line change
Expand Up @@ -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<array{stream: string, id: string, payload: string}>
*/
private array $prefetched = [];

/** Limits — set once in run(), checked by timers to know when to cancel. */
private int $maxJobs = 10_000;

Expand Down Expand Up @@ -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<array{stream: string, id: string, payload: string}>
*/
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;
}

/**
Expand Down Expand Up @@ -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) {
Expand Down
83 changes: 83 additions & 0 deletions tests/Unit/Worker/WorkerPrefetchTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
<?php

declare(strict_types=1);

use Webpatser\Torque\Worker\WorkerProcess;

use function Fledge\Async\Redis\createRedisClient;

/**
* A multi-stream XREADGROUP with COUNT 1 delivers up to one entry per stream.
* Every delivered entry must reach a Fiber; the ones beyond the first are
* buffered and served before the next read instead of being stranded in the
* PEL until retry_after (scrpr 2026-08-28: slowSync Persist parts waited
* 30 minutes each while the default stream was busy).
*/
function makeWorker(): WorkerProcess
{
return new WorkerProcess(['redis' => ['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([]);
});
Loading