CB-528: close the recovery race in AmqpReplyInbox #88

Merged
ltms merged 2 commits from worker/cb528-recovery-race into main 2026-08-16 17:27:22 +02:00
Member

Follow-up to PR #83 (CB-527/CB-528). An independent reviewer found the recovery sweep racing a concurrent publish.

Finding 1 (the race): failPendingPublishesOnRecovery walked and cleared pendingBySeq/pendingByMsgId without holding publishChannelLock. A publish() that registered while the sweep was still iterating could be failed even though it went out successfully on the already-recovered channel — a successful publish reported as failed. Since MessageService.reply mints a fresh msgId per retry, dedup-by-msgId cannot catch the resulting duplicate reply.

Fix: guard failPendingPublishesOnRecovery with publishChannelLock, the same lock publish() holds for its seq/map-put/basicPublish (it awaits the confirm outside the lock, so no deadlock — the sweep can only ever wait for an in-flight basicPublish call to return, never a broker round trip).

Finding 2: close() now fails in-flight publishes immediately with a clear message instead of leaving them to idle out the 10s CONFIRM_TIMEOUT_MS.

Finding 3: documented (not guarded) — pendingByMsgId assumes msgId uniqueness per in-flight publish; not reachable today since the only caller generates a fresh UUID per call.

Test: AmqpReplyInboxRecoveryRaceTest drives failPendingPublishesOnRecovery and a real publish() against each other directly (Proxy-backed fake AMQP channels, no live broker reconnect needed) with a large in-flight backlog to make the race window observable. Confirmed it fails without the fix (reply "fresh" wrongly failed as "connection recovered mid-publish") and passes reliably with it. Also covers close()'s fail-fast behavior.

mvn -f bridged/pom.xml clean install: Tests run: 809, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS.
mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest: Tests run: 8, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS.

Follow-up to PR #83 (CB-527/CB-528). An independent reviewer found the recovery sweep racing a concurrent publish. **Finding 1 (the race):** `failPendingPublishesOnRecovery` walked and cleared `pendingBySeq`/`pendingByMsgId` without holding `publishChannelLock`. A `publish()` that registered while the sweep was still iterating could be failed even though it went out successfully on the already-recovered channel — a *successful* publish reported as failed. Since `MessageService.reply` mints a fresh `msgId` per retry, dedup-by-`msgId` cannot catch the resulting duplicate reply. Fix: guard `failPendingPublishesOnRecovery` with `publishChannelLock`, the same lock `publish()` holds for its seq/map-put/`basicPublish` (it awaits the confirm outside the lock, so no deadlock — the sweep can only ever wait for an in-flight `basicPublish` call to return, never a broker round trip). **Finding 2:** `close()` now fails in-flight publishes immediately with a clear message instead of leaving them to idle out the 10s `CONFIRM_TIMEOUT_MS`. **Finding 3:** documented (not guarded) — `pendingByMsgId` assumes `msgId` uniqueness per in-flight publish; not reachable today since the only caller generates a fresh UUID per call. **Test:** `AmqpReplyInboxRecoveryRaceTest` drives `failPendingPublishesOnRecovery` and a real `publish()` against each other directly (Proxy-backed fake AMQP channels, no live broker reconnect needed) with a large in-flight backlog to make the race window observable. Confirmed it fails without the fix (reply "fresh" wrongly failed as "connection recovered mid-publish") and passes reliably with it. Also covers `close()`'s fail-fast behavior. `mvn -f bridged/pom.xml clean install`: Tests run: 809, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS. `mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest`: Tests run: 8, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS.
agent added 2 commits 2026-08-16 17:25:04 +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.
CB-528: close the recovery race in AmqpReplyInbox
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Failing after 1m20s
c553d795d8
failPendingPublishesOnRecovery walked and cleared pendingBySeq/pendingByMsgId
without holding publishChannelLock, so a publish() that registered while the
sweep was still iterating could be failed even though it published on the
already-recovered channel — a successful publish reported as failed, and
since MessageService.reply mints a fresh msgId per retry, dedup can't catch
the resulting duplicate. Guard the sweep with publishChannelLock: publish()
only holds it for the seq/map-put/basicPublish, so the sweep can only ever
wait for an in-flight basicPublish to return, never a broker round trip.

Also make close() fail in-flight publishes immediately with a clear message
instead of leaving them to idle out the 10s confirm timeout, and record the
(currently unreachable) msgId-uniqueness assumption pendingByMsgId relies on.

AmqpReplyInboxRecoveryRaceTest drives the sweep and a real publish() against
each other directly (no live broker reconnect) using Proxy-backed fake AMQP
channels and a large in-flight backlog to make the race window observable;
confirmed it fails without the guard (reply "fresh" wrongly failed as
"connection recovered mid-publish") and passes with it.
ltms merged commit 03286a589b into main 2026-08-16 17:27:22 +02:00
Owner

Correction to my merge message. I wrote there that "#83's own new tests are contract-tagged, so the default suite and CI prove nothing about them." The first half is right; the claim about CI is wrong.

.gitea/workflows/ci.yml has a dedicated contract job (CB-521, lines 64–97) that runs mvn -B -Pcontract test -Dtest=AmqpReplyInboxContractTest against a real rabbitmq:3.13 service container, with AMQP_URI pointed at it. So CI does cover these tests — the split is deliberate and documented in the workflow's own comments: the default job excludes the contract group because it has no broker, and the second job supplies one.

The merge stands, and my local contract run was still a real check rather than a redundant one — it confirmed the fix before the merge rather than after. But the reasoning I gave for it was wrong, and anyone reading that message would conclude this repo has a CI gap it does not have.

**Correction to my merge message.** I wrote there that "#83's own new tests are contract-tagged, so the default suite and CI prove nothing about them." The first half is right; the claim about CI is wrong. `.gitea/workflows/ci.yml` has a dedicated `contract` job (CB-521, lines 64–97) that runs `mvn -B -Pcontract test -Dtest=AmqpReplyInboxContractTest` against a real `rabbitmq:3.13` service container, with `AMQP_URI` pointed at it. So CI does cover these tests — the split is deliberate and documented in the workflow's own comments: the default job excludes the `contract` group because it has no broker, and the second job supplies one. The merge stands, and my local contract run was still a real check rather than a redundant one — it confirmed the fix before the merge rather than after. But the reasoning I gave for it was wrong, and anyone reading that message would conclude this repo has a CI gap it does not have.
Sign in to join this conversation.