Skip to content

[BUG] resend() removes entries while iterating unconfirmedPacketIds, killing the QoS1/2 retransmission chain #289

Description

@ImDanXie

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:

  1. MaxResendTimes=1, ResendTimeoutSeconds=1, ReceivingMaximum=2
  2. Publish two QoS1 messages
  3. 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
  4. 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.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions