CB-527: basicQos prefetch on the AMQP consumer — the queue drains into gateway heap #10

Closed
opened 2026-08-10 17:35:07 +02:00 by kevin · 1 comment
Owner

Found by the 2026-08-10 adversarial design review of CB-308 / wiki chapter 10, but it is a v1.0.x as-built gap, independent of federation.

Problem

AmqpReplyInbox starts its manual-ack consumer without ever calling basicQos. With no prefetch bound, the broker pushes every message into the JVM immediately, so:

  • the durable queue sits near-empty at all times — the real backlog lives in the in-memory held map, where nothing caps it (unbounded heap growth if a primary never drains);
  • any future queue-level control (x-max-length, per-message TTL, the CB-308 inbox caps) would guard an empty queue and never fire;
  • broker-side queue-depth metrics read healthy while the gateway holds the actual backlog.

Fix

Set a prefetch bound (channel.basicQos(n), n small — e.g. 16–64, configurable) before basicConsume, so unacked messages beyond the window stay on the queue, where durability, caps, and observability actually apply.

Acceptance

  • Contract test (@Tag("contract"), runs in CI against the real broker): publish >n messages to a held inbox without draining → broker reports queue depth ≥ (published − n); after drain+ack, depth returns to 0.
  • Existing consume-and-hold semantics unchanged: dedup by msgId, ack-after-drain, cross-restart redelivery (CB-307 Stage 2 test still green).
  • mvn clean install green.

🤖 Generated with Claude Code

https://claude.ai/code/session_013ZGgxLQ2VpwZhEYoru8rkf

Found by the 2026-08-10 adversarial design review of CB-308 / wiki chapter 10, but it is a **v1.0.x as-built gap**, independent of federation. ## Problem `AmqpReplyInbox` starts its manual-ack consumer without ever calling `basicQos`. With no prefetch bound, the broker pushes every message into the JVM immediately, so: - the durable queue sits near-empty at all times — the real backlog lives in the in-memory `held` map, where **nothing caps it** (unbounded heap growth if a primary never drains); - any future queue-level control (`x-max-length`, per-message TTL, the CB-308 inbox caps) would guard an empty queue and never fire; - broker-side queue-depth metrics read healthy while the gateway holds the actual backlog. ## Fix Set a prefetch bound (`channel.basicQos(n)`, n small — e.g. 16–64, configurable) before `basicConsume`, so unacked messages beyond the window stay **on the queue**, where durability, caps, and observability actually apply. ## Acceptance - Contract test (`@Tag("contract")`, runs in CI against the real broker): publish >n messages to a held inbox without draining → broker reports queue depth ≥ (published − n); after drain+ack, depth returns to 0. - Existing consume-and-hold semantics unchanged: dedup by msgId, ack-after-drain, cross-restart redelivery (CB-307 Stage 2 test still green). - `mvn clean install` green. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_013ZGgxLQ2VpwZhEYoru8rkf
ltms added this to the 1.1 — single-host close-out milestone 2026-08-16 16:49:38 +02:00
Owner

Delivered and merged to main as part of #88 (which carried #83).

AmqpReplyInbox now calls basicQos before basicConsume, so the held backlog per owned target is bounded instead of the broker pushing the whole queue into the JVM. Configurable as broker.prefetch, default 32 (AmqpReplyInbox.DEFAULT_PREFETCH).

Verified by the lead against a real broker, not taken on the worker's word: mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest → Tests run: 8, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS.

The acceptance criterion above asked for a @Tag("contract") test running in CI against the real broker. That is satisfied: .gitea/workflows/ci.yml has a dedicated contract job (CB-521, lines 64–97) running exactly this test against a rabbitmq:3.13 service container. The default job excludes the group because it has no broker; the second job supplies one.

Note the new knob is not yet documented in bridged.example.yaml — that whole broker: section is missing from it. Tracked as #85 (CB-597).

Delivered and merged to `main` as part of #88 (which carried #83). `AmqpReplyInbox` now calls `basicQos` before `basicConsume`, so the held backlog per owned target is bounded instead of the broker pushing the whole queue into the JVM. Configurable as `broker.prefetch`, default 32 (`AmqpReplyInbox.DEFAULT_PREFETCH`). Verified by the lead against a real broker, not taken on the worker's word: `mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest` → Tests run: 8, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS. The acceptance criterion above asked for a `@Tag("contract")` test running in CI against the real broker. That is satisfied: `.gitea/workflows/ci.yml` has a dedicated `contract` job (CB-521, lines 64–97) running exactly this test against a `rabbitmq:3.13` service container. The default job excludes the group because it has no broker; the second job supplies one. Note the new knob is not yet documented in `bridged.example.yaml` — that whole `broker:` section is missing from it. Tracked as #85 (CB-597).
ltms closed this issue 2026-08-16 17:30:33 +02:00
Sign in to join this conversation.
2 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#10