Visitar URL original
feat(queue): bound the Redis broker's automatic reap sweep by levivannoort · Pull Request #14243 · appwrite/appwrite · GitHub
Skip to content

feat(queue): bound the Redis broker's automatic reap sweep - #14243

Open
levivannoort wants to merge 3 commits into
mainfrom
fix/queue-reap-limits
Open

levivannoort wants to merge 3 commits into
mainfrom
fix/queue-reap-limits

Conversation

@levivannoort

@levivannoort levivannoort commented Oct 7, 2026 •

Copy link
Copy Markdown
Member

What does this PR do?

Lets consumers limit what the Redis broker's automatic recovery sweep replays.

maintain() calls reap() with no maxAttempts or newerThan, and the constructor only exposes reapAfter. So the sweep requeues any claim older than reapAfter that has no live heartbeat, however old it is and however many times it has already been requeued, and consumers have no way to change that.

This hurts when upgrading from queue 1.x. Its processing list was never read back, so stale claims and their payloads (which never expire) piled up for months. On the first deploy of 6.x, the sweep replayed all of them at 1,000 per minute. Messages in an old payload shape failed. Ones that still decoded ran months-old work again. A message that crashes its worker every time would also loop forever, once per reapAfter.

The change adds two optional constructor arguments that pass through to reap():

new Redis($receive, $commands, reapMaxAttempts: 3, reapMaxAge: 172_800);
  • reapMaxAttempts: a claim that has already been requeued this many times goes to the dead list instead of being requeued again.
  • reapMaxAge: a claim published longer ago than this many seconds goes to the dead list instead of being replayed. To have any effect it has to be larger than reapAfter. Otherwise every claim the sweep picks up goes to dead. The age is measured from the claim's latest publish. A requeue republishes the message, which restarts the clock, the same way newerThan already works on reap() and retry(). So reapMaxAge catches a backlog that was never requeued, and reapMaxAttempts is what stops a crash loop.

Both default to null, so existing users see no change.

Test Plan

Two new E2E tests in RedisBrokerRecoveryTest drive maintain() against a real Redis:

  • testMaintainParksClaimsOlderThanTheMaxAge: of two stranded claims, the one older than the maximum age goes to dead and the recent one is requeued.
  • testMaintainParksClaimsAtTheMaxAttempts: a claim is requeued on its first stranding and goes to dead once it reaches the attempt limit.

Both fail with the pass-through removed. RedisBrokerRecoveryTest passes (22 tests) and bin/monorepo check queue is clean.

Related PRs and Issues

  • Found while upgrading a consumer from queue 1.6 to 6.1. The upgrade was rolled back until the sweep can be bounded.

Checklist

  • Have you read the Contributing Guidelines on issues?
  • If the PR includes a change to an API's metadata (desc, label, params, etc.), does it also include updated API specs and example docs? N/A

maintain() called reap() with no maxAttempts or newerThan, and the
constructor only exposed reapAfter, so consumers could not stop the sweep
from replaying arbitrarily old claims or looping poison messages.
reapMaxAttempts and reapMaxAge pass through to reap(); both default to
null, keeping the current behaviour.
@tenki-reviewer

tenki-reviewer Bot commented Oct 7, 2026 •

Copy link
Copy Markdown

Review complete. 🟡 1 medium

📍 Findings outside the diff (1) — 🟡 1 medium — defects on lines GitHub can't attach comments to

🟡 Medium — reapMaxAge never ages out repeatedly requeued claims · Redis.php:484–491 · unchanged line

// packages/queue/src/Broker/Redis.php
484	            $dead = ($maxAttempts !== null && $job->getAttempts() >= $maxAttempts)
485	                || ($newerThan !== null && $job->getTimestamp() < $now - $newerThan);
486	            $moved = $this->script($this->commands, 'reclaim', [
487	                $ownerKey, "{$queue->namespace}.claims.{$queue->name}.{$pid}",
488	                "{$queue->namespace}.jobs.{$queue->name}.{$pid}", $processing,
489	                "{$queue->namespace}.stats.{$queue->name}.processing",
490	                "{$queue->namespace}." . ($dead ? 'dead' : 'queue') . ".{$queue->name}",
491	            ], [\is_string($owner) ? $owner : '', $pid, $dead ? '' : $this->retryPayload($queue, $job), $queue->jobTtl]);

The reapMaxAge comparison at packages/queue/src/Broker/Redis.php:485 reads $job->getTimestamp(), but every requeue through retryPayload() (Redis.php:520) resets timestamp to time(). A claim stranded by a worker crash before settlement is requeued, re-claimed, stranded again, and each cycle stamps a fresh timestamp, so it can never satisfy $job->getTimestamp() < $now - $newerThan and loops forever under reapMaxAge alone. The new dead-letter age guard silently does not bound its own requeue path; only reapMaxAttempts can catch these.


This PR wires two optional constructor parameters through the Redis broker's recovery logic so reaping can park poison messages after a configurable attempt threshold and skip claims newer than a configurable age, with new E2E tests covering both paths.

Files Change
packages/queue/src/Broker/Redis.php Adds reapMaxAttempts/reapMaxAge constructor params and passes them into reap() from maintain()
packages/queue/tests/E2E/RedisBrokerRecoveryTest.php New E2E tests for attempt-based parking and age-based reap filtering

One medium-severity issue remains: reapMaxAge only filters by the original claim timestamp, so a claim that keeps being requeued never crosses the age boundary and can loop indefinitely instead of being parked.

Reviewed commit: 828cdc3

@hansi-codes

hansi-codes Bot commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

🟢 Tier S · Ready to merge

The incremental change only clarifies an existing comment and introduces no behavioral defect.

Adds optional maximum-attempt and maximum-age bounds to Redis broker recovery and applies them during automatic maintenance sweeps, parking exhausted claims in the dead queue. Adds real-Redis E2E coverage for age- and attempt-based parking while preserving recovery of eligible claims.

Latest changes: The newest commit clarifies that the age bound resets on requeue, while the attempt bound limits repeated handler crashes.

Verdict New comments Fixed Still open
✅ Approved 0 0 0
📂 Walkthrough · 2
File Change
packages/queue/src/Broker/Redis.php Adds configurable bounds to automatic recovery and clarifies how they affect repeated crashes.
packages/queue/tests/E2E/RedisBrokerRecoveryTest.php Exercises automatic parking of claims that exceed age or attempt bounds.

Reviewed the commits since f582b55 · Details · Comment @hansi-codes review to re-run, or mention @hansi-codes with a question.

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 Tier S · Looks good to merge. Summary

@github-actions

github-actions Bot commented Oct 7, 2026

Copy link
Copy Markdown

Security rules

No new WARNING or ERROR findings from security rules.

32 existing findings tracked in .semgrep/baseline.json
  • php.appwrite.guest-write-without-abuse-limit (17)
  • php.appwrite.permissive-write-permission (8)
  • php.appwrite.secret-compare-timing (5)
  • php.appwrite.weak-secret-env-default (2)

Posted by Checks / Rules. Re-runs update this comment in place. Rule details and baseline: .semgrep/README.md.

@github-actions

github-actions Bot commented Oct 7, 2026 •

Copy link
Copy Markdown

✨ Benchmark results

Comparing main (before) → fix/queue-reap-limits (after).

Metric Before After Change
🚀 Requests/sec 205.3 214.76 ⚪ +4.6%
⏱️ Latency P50 85.26 ms 81.55 ms ⚪ -4.3%
⏱️ Latency P95 194.82 ms 188.1 ms ⚪ -3.4%
Per-scenario breakdown & investigation details

Metrics below reflect the current branch (after). Δ P95 compares against the base.

Scenario P50 (ms) P95 (ms) Requests RPS Δ P95 (ms)
API total 81.55 188.1 13,623 214.76 -6.72
Account 155.63 288.48 717 11.96 -4.52
TablesDB 78.66 145.25 7,409 119.45 -17.72
Storage 74.69 163.68 3,585 59.45 -0.26
Functions 118.06 224.08 1,912 32.35 +1.65

Top API waits (after)

API request Max wait (ms)
account.prefs.update 445.6
account.name.update 418.43
functions.create 410.4
functions.variables.update 383.7
tablesdb.rows.create 374.81

@levivannoort

Copy link
Copy Markdown
Member Author

@tenki-reviewer

@tenki-reviewer

tenki-reviewer Bot commented Oct 8, 2026 •

Copy link
Copy Markdown

Review complete. 🟡 2 medium

📍 Findings outside the diff (2) — 🟡 2 medium — defects on lines GitHub can't attach comments to

🟡 Medium — Reservation recovery bypasses reapMaxAttempts/reapMaxAge · Redis.php:278–281 · unchanged line

// packages/queue/src/Broker/Redis.php
278	            $expired = $this->script($this->commands, 'expired', [$registry], [self::REAP_LIMIT]);
279	            if ($expired !== []) {
280	                $this->script($this->commands, 'recover', [$registry, "{$queue->namespace}.queue.{$queue->name}", ...$expired], [self::REAP_LIMIT]);
281	            }

The new reapMaxAttempts/reapMaxAge dead-queue gates count $job->getAttempts() and the envelope timestamp, but both are only advanced or checked on the retry/reclaim paths via retryPayload() (packages/queue/src/Broker/Redis.php:484-485, 522). Expired reservations recovered by maintain() go through the separate expired/recover Lua path (Redis.php:278-280), where recover.lua pushes the raw envelope back to the ready queue unchanged — attempts are never incremented and no age check runs. A worker that repeatedly crashes between reserve and claim therefore re-delivers the same message forever without it ever being parked on the dead queue, defeating the guard these parameters were added to provide. The same reliance on getAttempts() appears in retry() at Redis.php:402.


🟡 Medium — Requeue resets timestamp, so reapMaxAge never accumulates · Redis.php:515–524 · unchanged line

// packages/queue/src/Broker/Redis.php
515	    private function retryPayload(Queue $queue, Message $job): string
516	    {
517	        $payload = [
518	            'pid' => uniqid(more_entropy: true),
519	            'queue' => $queue->name,
520	            'timestamp' => time(),
521	            'payload' => $job->getPayload(),
522	            'attempts' => $job->getAttempts() + 1,
523	        ];
524	        return $this->codec->encode($payload);

The new reapMaxAge gate is documented as parking claims "published longer ago than this" (Redis.php:63-65), but the only timestamp it can measure is reset on every recovery: retryPayload() writes a fresh timestamp => time() whenever a claim is requeued (packages/queue/src/Broker/Redis.php:515-524). A poison message that keeps stranding with a cycle shorter than reapMaxAge gets its clock reset each sweep, so its age never crosses the gate and — with reapMaxAttempts left null — the sweep resurrects it forever. The bound configured to stop infinite replay silently provides no protection in the most common stranding cadence.

🧹 Nitpicks (1) — 🟢 1 low
  • 🟢 Destructure receive() result without short-array guard (RedisBrokerRecoveryTest.php:307) — The new maxAge test destructures [$stale] = $broker->receive($this->queue, 0, 2) and immediately calls $stale->getPid(), without the [0] ?? null plus assertInstanceOf guard every other test in this suite uses (packages/queue/tests/E2E/RedisBrokerRecoveryTest.php:307-308 vs lines 292-293).

This PR extends the utopia queue Redis broker so maintain() reaps claims that exceed a configurable attempt count or publish age and parks them on a dead queue instead of requeueing, with a new E2E test covering the maxAttempts branch.

The gates are enforced only in reap() and retry() on the claim-reclaim path, while maintain()'s expired-reservation recovery pushes stranded envelopes straight back to the ready queue without touching attempts or age, and retryPayload() resets the envelope timestamp on every requeue, so a poison message can loop forever without ever being parked.

Files Change
packages/queue/src/Broker/Redis.php Adds reapMaxAttempts/reapMaxAge constructor params forwarded to reap(), dead-parks over-limit claims, and resets the envelope timestamp in retryPayload().
packages/queue/tests/E2E/RedisBrokerRecoveryTest.php Adds an E2E test asserting claims older than reapMaxAge land on the dead queue, using bare destructuring for one receive() call.

Note: one low-severity lapsed-heartbeat dead-parking edge case was verified at confidence 75 and falls just below the reporting threshold; it concerns a live worker whose heartbeat lapsed past CLAIM_TTL having its in-flight message terminal-parked when the new bounds are configured.

Reviewed commit: 828cdc3

@levivannoort

Copy link
Copy Markdown
Member Author

@tenki-reviewer thanks. Addressed in f582b55:

  • Requeues reset the reapMaxAge clock: correct, and it's intended. retryPayload() republishes the message, and the reapAfter staleness check relies on that fresh timestamp. newerThan already behaves this way on reap() and retry(). I've updated the constructor comment and the PR description to say the age counts from the latest publish and that reapMaxAttempts is what stops a crash loop. reapMaxAge is for a backlog that was never requeued.
  • Reservation recovery skips the limits: I'm leaving this as is. receive() reserves and claims inside the same call, so a reservation only exists before receive() returns, and the message has never reached a handler. It can't be what crashes the worker, so there's nothing to count as an attempt. Envelopes that fail to decode already go to poison at claim time.
  • Destructuring nit: fixed. The test now uses the same [0] ?? null plus assertInstanceOf guard as the rest of the suite.

@tenki-reviewer

tenki-reviewer Bot commented Oct 8, 2026 •

Copy link
Copy Markdown

Review complete. 🟡 1 medium

💬 Inline comments (1)

  • 🟡 Reservation recovery bypasses the new dead-letter bounds — Redis.php:281

This PR extends the Redis broker crash-recovery path so messages that repeatedly fail recovery or sit unclaimed too long are parked in the dead queue instead of being replayed forever. Two new optional constructor params on Broker\Redis thread through maintain() into reap(), preserving the previous unbounded behavior when left null.

Files Change
packages/queue/src/Broker/Redis.php Adds reapMaxAttempts/reapMaxAge params and applies them in the reap dead-parking gate.
packages/queue/tests/E2E/RedisBrokerRecoveryTest.php Adds E2E tests covering max-attempt and max-age dead-lettering.

Reviewed commit: f582b55

@tenki-reviewer tenki-reviewer Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Adds configurable dead-letter bounds (reapMaxAttempts, reapMaxAge) to the Redis queue broker reservation reaping, with two new E2E recovery tests.

Key findings

  • 🟡 Reservation recovery bypasses the new dead-letter bounds — Redis.php:281

Comment thread packages/queue/src/Broker/Redis.php

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants