Environment
- BiFroMQ
4.0.0
- Verified against
main @ ba2c34e0 — still present on the latest main at the time of filing
- Module:
bifromq-mqtt
- Files:
MQTTSessionHandler.java:1292 (iteration), :1311 (call), :898-904 (structural removal)
Summary
resend() iterates unconfirmedPacketIds.values() with a for-each loop, and inside the loop calls confirm(...) for every entry that has exceeded maxResendTimes. confirm() removes entries from the same map through a different iterator, so the outer iterator's modCount becomes stale and the next next() throws ConcurrentModificationException.
The exception escapes the scheduled task, which means the code after the loop — flush(true) and scheduleResend() — never runs. The retransmission chain dies permanently: the remaining in-flight entries are neither retransmitted nor dropped, unconfirmedPacketIds stays non-empty, clientReceiveQuota() stays at 0, and the session stops receiving messages altogether.
Root cause
// :190 unconfirmedPacketIds is a plain LinkedHashMap
private final LinkedHashMap<Integer, ConfirmingMessage> unconfirmedPacketIds = new LinkedHashMap<>();
// :1292 resend() — for-each over values()
for (ConfirmingMessage confirmingMsg : unconfirmedPacketIds.values()) {
if (confirmingMsg.sentCount <= settings.maxResendTimes) {
... // retransmit
} else {
reportDropConfirmableMsgEvent(confirmingMsg.message, DropReason.MaxRetried);
confirm(confirmingMsg, false); // :1311 -> removes from the map
receiveQuota.onErrorSignal(now);
}
}
// :898-904 confirm() — a *different* iterator removes entries
confirmingMsg.setAcked();
Iterator<Integer> packetIdItr = unconfirmedPacketIds.keySet().iterator();
while (packetIdItr.hasNext()) {
int packetId = packetIdItr.next();
ConfirmingMessage head = unconfirmedPacketIds.get(packetId);
if (head.acked) {
packetIdItr.remove(); // structural modification
...
} else {
break;
}
}
Note that all five other confirm() call sites (:1115, :1122, :1152, :1161, :1166) deliberately defer through ctx.executor().execute(...) to avoid exactly this — only the one inside resend() is a direct call.
Steps to reproduce
Using the project's own TestNG harness, no custom client needed:
MaxResendTimes=1, ResendTimeoutSeconds=1, ReceivingMaximum=2
- Publish two QoS1 messages
- Acknowledge neither — the crucial difference from the existing
qoS1ResendKeepsRunningAfterPartialAck, which acks the first message and therefore leaves only one entry in the map when the drop happens, so its iterator never has a next element
- Advance the clock 6 × 2s to pass
maxResendTimes
Observed
Wanted 2 times:
-> at ...MQTT3TransientSessionHandlerTest
.resendDroppingHeadAbortsDropOfRemaining(MQTT3TransientSessionHandlerTest.java:1778)
But was 1 time:
-> at ...MQTTSessionHandler.reportDropConfirmableMsgEvent(MQTTSessionHandler.java:1264)
Both messages are past maxResendTimes and should be dropped; only one is, because the loop aborts immediately after dropping the head.
Control experiment
To rule out the alternative explanation ("the second message simply wasn't eligible for dropping"), the same assertion was run before and after the fix:
|
Result |
| before fix |
FAIL — 1 dropped |
| after fix |
PASS — 2 dropped |
This confirms the second message was eligible, and the only reason it was not dropped is the aborted iteration.
Impact
Default configuration (MaxResendTimes=3, ResendTimeoutSeconds=10) reaches the drop path after roughly 40 seconds of a subscriber being online but not acknowledging. From then on the session silently stops delivering messages entirely.
Suggested fix
Collect the entries to drop and process them after the iteration completes:
List<ConfirmingMessage> toDrop = new ArrayList<>();
for (ConfirmingMessage confirmingMsg : unconfirmedPacketIds.values()) {
if (confirmingMsg.sentCount <= settings.maxResendTimes) {
... // unchanged
} else {
toDrop.add(confirmingMsg);
}
}
if (flush) {
flush(true);
}
// confirm() structurally modifies unconfirmedPacketIds, so drop after the iteration
for (ConfirmingMessage dropped : toDrop) {
reportDropConfirmableMsgEvent(dropped.message, DropReason.MaxRetried);
confirm(dropped, false);
receiveQuota.onErrorSignal(now);
}
Regression test added: MQTT3TransientSessionHandlerTest.resendDroppingHeadAbortsDropOfRemaining. With the fix applied, MQTT3TransientSessionHandlerTest + MQTT3PersistentSessionHandlerTest (87 tests) pass.
Environment
4.0.0main@ba2c34e0— still present on the latestmainat the time of filingbifromq-mqttMQTTSessionHandler.java:1292(iteration),:1311(call),:898-904(structural removal)Summary
resend()iteratesunconfirmedPacketIds.values()with a for-each loop, and inside the loop callsconfirm(...)for every entry that has exceededmaxResendTimes.confirm()removes entries from the same map through a different iterator, so the outer iterator'smodCountbecomes stale and the nextnext()throwsConcurrentModificationException.The exception escapes the scheduled task, which means the code after the loop —
flush(true)andscheduleResend()— never runs. The retransmission chain dies permanently: the remaining in-flight entries are neither retransmitted nor dropped,unconfirmedPacketIdsstays non-empty,clientReceiveQuota()stays at 0, and the session stops receiving messages altogether.Root cause
Note that all five other
confirm()call sites (:1115,:1122,:1152,:1161,:1166) deliberately defer throughctx.executor().execute(...)to avoid exactly this — only the one insideresend()is a direct call.Steps to reproduce
Using the project's own TestNG harness, no custom client needed:
MaxResendTimes=1,ResendTimeoutSeconds=1,ReceivingMaximum=2qoS1ResendKeepsRunningAfterPartialAck, which acks the first message and therefore leaves only one entry in the map when the drop happens, so its iterator never has a next elementmaxResendTimesObserved
Both messages are past
maxResendTimesand should be dropped; only one is, because the loop aborts immediately after dropping the head.Control experiment
To rule out the alternative explanation ("the second message simply wasn't eligible for dropping"), the same assertion was run before and after the fix:
This confirms the second message was eligible, and the only reason it was not dropped is the aborted iteration.
Impact
Default configuration (
MaxResendTimes=3,ResendTimeoutSeconds=10) reaches the drop path after roughly 40 seconds of a subscriber being online but not acknowledging. From then on the session silently stops delivering messages entirely.Suggested fix
Collect the entries to drop and process them after the iteration completes:
Regression test added:
MQTT3TransientSessionHandlerTest.resendDroppingHeadAbortsDropOfRemaining. With the fix applied,MQTT3TransientSessionHandlerTest+MQTT3PersistentSessionHandlerTest(87 tests) pass.