CB-527/CB-528: bound AMQP prefetch and confirm publishes #83

Closed
agent wants to merge 0 commits from worker/cb527-997d99-3 into main
Member

Closes #10 and #11.

CB-527 (prefetch) — the AMQP consumer 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 via broker.prefetch (default 32, AmqpReplyInbox.DEFAULT_PREFETCH).

CB-528 (enforce publish) — publish now runs on a dedicated confirm-mode channel, separate from the consume/ack channel, with mandatory=true and a return listener. An unroutable or unconfirmed publish now throws IllegalStateException instead of silently vanishing. Because a broker Return always precedes its matching Confirm, the confirm callback checks a per-message returned flag at confirm time rather than trusting an ack alone. The ack path is untouched (separate channel/lock) and never waits on a publish confirm.

Tests: mvn clean install (807 tests, unpiped) is green. The 3 new @Tag("contract") tests (prefetchBoundsTheHeldBacklog, unroutablePublishReportsFailureNotSilentSuccess, confirmedPublishDeliversNormally) plus the 5 existing ones all pass under mvn test -Pcontract against a Testcontainers RabbitMQ broker (8/8, BUILD SUCCESS).

Closes #10 and #11. **CB-527 (prefetch)** — the AMQP consumer 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 via broker.prefetch (default 32, AmqpReplyInbox.DEFAULT_PREFETCH). **CB-528 (enforce publish)** — publish now runs on a dedicated confirm-mode channel, separate from the consume/ack channel, with mandatory=true and a return listener. An unroutable or unconfirmed publish now throws IllegalStateException instead of silently vanishing. Because a broker Return always precedes its matching Confirm, the confirm callback checks a per-message returned flag at confirm time rather than trusting an ack alone. The ack path is untouched (separate channel/lock) and never waits on a publish confirm. Tests: mvn clean install (807 tests, unpiped) is green. The 3 new @Tag("contract") tests (prefetchBoundsTheHeldBacklog, unroutablePublishReportsFailureNotSilentSuccess, confirmedPublishDeliversNormally) plus the 5 existing ones all pass under mvn test -Pcontract against a Testcontainers RabbitMQ broker (8/8, BUILD SUCCESS).
agent added 1 commit 2026-08-16 16:56:00 +02:00
CB-527/CB-528: bound AMQP prefetch and confirm publishes before claiming durable
CI / build (pull_request) Successful in 1m3s
CI / contract (pull_request) Successful in 1m15s
4fa6553db5
CB-527: basicQos(prefetch) on the consume channel before basicConsume, configurable
via broker.prefetch (default 32), so an undrained inbox backlog stays on the broker
instead of growing the JVM heap without limit.

CB-528: publish moves to its own confirm-mode channel with mandatory=true and a
return listener, so an unroutable or unconfirmed reply now throws instead of
vanishing silently. The confirm callback checks the per-message returned flag
(set by the return listener, which the broker always fires before the matching
confirm) so an acked-but-returned publish is still reported as a failure. The ack
path stays on its own channel/lock and never waits on a publish confirm.
Owner

Lead review — not merging yet. One real gap, fix round dispatched.

What I verified myself

  • Independent mvn -f bridged/pom.xml clean install on this branch, unpiped: Tests run: 807, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS, exit 0.
  • The multiple flag is handled correctly — resolveConfirm iterates headMap(seq, true) and checks each entry's own returned flag, so a mixed batch where some messages were returned and some were not resolves per message. This was the thing most likely to be wrong and it is right.
  • Issue #11's core requirement holds: MessageService.reply (MessageService.java:276) calls inbox.publish(...) with no try/catch, so a failed publish propagates instead of returning a false true. Neither BridgeMcp.reply nor BridgedApp.replyMessage swallows it on the way out.

One caveat on the test count. main was already at 807. Every new test here is @Tag("contract"), which the default profile excludes — so CI proves nothing about this change. The author did run the contract suite locally under Testcontainers (8 tests, green) and said so plainly, which is the right way to report it. Worth knowing that the green tick on this PR is not evidence.

The finding — a reviewer caught a real race

failPendingPublishesOnRecovery() walks and clears pendingBySeq / pendingByMsgId without holding publishChannelLock, while publish() holds that lock across the seq + map-put + basicPublish.

The interleaving: the connection recovers, the channel is usable again, handleRecovery starts sweeping. A publish() runs concurrently, registers its Pending, and succeeds on the recovered channel. The still-running sweep then reaches that new entry and fails it.

This inverts the whole point of CB-528. The ticket exists so a failed publish is never reported as success. This is the mirror image — a successful publish reported as failed. And it is worse than a stray error, because MessageService.reply mints a fresh UUID per call, so the retry carries a new msgId and dedup cannot catch it. The reply lands twice.

Fix dispatched: guard the sweep with publishChannelLock. I checked the deadlock question first — publish() waits for its confirm outside the lock, so the sweep can only ever wait for an in-flight basicPublish to return. It also gives the right semantics rather than merely a safe one: a publish holding the lock when recovery fires really did go out on the old channel, so failing it is correct.

Also in the fix round: close() will fail in-flight publishes promptly instead of letting them sit out the full 10s timeout, and the msgId-uniqueness assumption in pendingByMsgId gets a comment. That last one is not getting a guard — it is unreachable from the only production caller, and building for it would be unearned complexity.

Everything else checked out

Prefetch before basicConsume, the dedicated publish channel, mandatory plus the return listener, return-before-confirm ordering, and no map leaks on any path — success, IOException, timeout, ack, nack, recovery — all correct, using the identity-checked remove(key, value) form throughout.

Good work. The structure is right; it is one missing lock away from merging.

Noted for later, out of scope here: the stale line in docs/CB-307-Reliable-Delivery.md telling readers not to add a broker: block. Now filed as part of #85 (CB-597).

**Lead review — not merging yet. One real gap, fix round dispatched.** ## What I verified myself - Independent `mvn -f bridged/pom.xml clean install` on this branch, unpiped: **`Tests run: 807, Failures: 0, Errors: 0, Skipped: 0`, BUILD SUCCESS, exit 0.** - The `multiple` flag is handled correctly — `resolveConfirm` iterates `headMap(seq, true)` and checks each entry's own `returned` flag, so a mixed batch where some messages were returned and some were not resolves per message. This was the thing most likely to be wrong and it is right. - Issue #11's core requirement holds: `MessageService.reply` (`MessageService.java:276`) calls `inbox.publish(...)` with **no try/catch**, so a failed publish propagates instead of returning a false `true`. Neither `BridgeMcp.reply` nor `BridgedApp.replyMessage` swallows it on the way out. **One caveat on the test count.** `main` was already at 807. Every new test here is `@Tag("contract")`, which the default profile excludes — so **CI proves nothing about this change**. The author did run the contract suite locally under Testcontainers (8 tests, green) and said so plainly, which is the right way to report it. Worth knowing that the green tick on this PR is not evidence. ## The finding — a reviewer caught a real race `failPendingPublishesOnRecovery()` walks and clears `pendingBySeq` / `pendingByMsgId` **without holding `publishChannelLock`**, while `publish()` holds that lock across the seq + map-put + `basicPublish`. The interleaving: the connection recovers, the channel is usable again, `handleRecovery` starts sweeping. A `publish()` runs concurrently, registers its `Pending`, and succeeds on the **recovered** channel. The still-running sweep then reaches that new entry and fails it. **This inverts the whole point of CB-528.** The ticket exists so a failed publish is never reported as success. This is the mirror image — a *successful* publish reported as failed. And it is worse than a stray error, because `MessageService.reply` mints a fresh `UUID` per call, so the retry carries a **new `msgId`** and dedup cannot catch it. The reply lands twice. Fix dispatched: guard the sweep with `publishChannelLock`. I checked the deadlock question first — `publish()` waits for its confirm **outside** the lock, so the sweep can only ever wait for an in-flight `basicPublish` to return. It also gives the right semantics rather than merely a safe one: a publish holding the lock when recovery fires really did go out on the old channel, so failing it is correct. Also in the fix round: `close()` will fail in-flight publishes promptly instead of letting them sit out the full 10s timeout, and the `msgId`-uniqueness assumption in `pendingByMsgId` gets a comment. That last one is **not** getting a guard — it is unreachable from the only production caller, and building for it would be unearned complexity. ## Everything else checked out Prefetch before `basicConsume`, the dedicated publish channel, `mandatory` plus the return listener, return-before-confirm ordering, and no map leaks on any path — success, IOException, timeout, ack, nack, recovery — all correct, using the identity-checked `remove(key, value)` form throughout. Good work. The structure is right; it is one missing lock away from merging. Noted for later, out of scope here: the stale line in `docs/CB-307-Reliable-Delivery.md` telling readers not to add a `broker:` block. Now filed as part of **#85 (CB-597)**.
Owner

Merged as part of #88, which stacked the review fix on top of this branch. This branch's head 4fa6553 is now an ancestor of origin/main — confirmed with git merge-base --is-ancestor, not assumed. Closing as delivered, not as rejected. The work in here is good; it just shipped under a different PR number.

Merged as part of #88, which stacked the review fix on top of this branch. This branch's head `4fa6553` is now an ancestor of `origin/main` — confirmed with `git merge-base --is-ancestor`, not assumed. Closing as delivered, not as rejected. The work in here is good; it just shipped under a different PR number.
ltms closed this pull request 2026-08-16 17:27:39 +02:00
Some checks are pending
CI / build (pull_request) Successful in 1m3s
CI / contract (pull_request) Successful in 1m15s

Pull request closed

Sign in to join this conversation.