From e32c8700e8e2b6c605a10d21d047a634f66b3694 Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Thu, 24 Sep 2026 21:11:48 +0800 Subject: [PATCH] fix(mqtt): drop max-retried entries after the resend iteration completes (#289) resend() iterated unconfirmedPacketIds.values() with a for-each loop and called confirm(...) inside the loop for entries past maxResendTimes. confirm() removes entries from the same map through a different iterator, so the outer iterator throws ConcurrentModificationException once another entry remains after the dropped head. The exception escapes the scheduled task, so flush(true) and scheduleResend() never run: the retransmission chain dies and the session silently stops delivering messages. Collect the entries to drop and process them after the iteration. Regression test: resendDroppingHeadAbortsDropOfRemaining - two QoS1 messages, acknowledge neither (the existing partial-ack test leaves a single in-flight entry, so it can never trigger the CME). Control experiment: fails before the fix (1 dropped), passes after (2 dropped). Full TransientSessionHandlerTest (64) and PersistentSessionHandlerTest (23) pass. --- .../mqtt/handler/MQTTSessionHandler.java | 12 +++-- .../v3/MQTT3TransientSessionHandlerTest.java | 48 +++++++++++++++++++ 2 files changed, 57 insertions(+), 3 deletions(-) diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java index 97c859075..b1967fd0d 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java @@ -77,6 +77,7 @@ import io.netty.handler.codec.mqtt.MqttTopicSubscription; import io.netty.handler.codec.mqtt.MqttUnsubscribeMessage; import java.time.Duration; +import java.util.ArrayList; import java.util.HashSet; import java.util.Iterator; import java.util.LinkedHashMap; @@ -1286,6 +1287,7 @@ private void scheduleResend() { private void resend() { long now = sessionCtx.nanoTime(); boolean flush = false; + List toDrop = new ArrayList<>(); for (ConfirmingMessage confirmingMsg : unconfirmedPacketIds.values()) { if (confirmingMsg.sentCount <= settings.maxResendTimes) { if (ctx.channel().isWritable()) { @@ -1306,14 +1308,18 @@ private void resend() { break; } } else { - reportDropConfirmableMsgEvent(confirmingMsg.message, DropReason.MaxRetried); - confirm(confirmingMsg, false); - receiveQuota.onErrorSignal(now); + toDrop.add(confirmingMsg); } } if (flush) { flush(true); } + // drop after the iteration completes: confirm() structurally modifies unconfirmedPacketIds + for (ConfirmingMessage dropped : toDrop) { + reportDropConfirmableMsgEvent(dropped.message, DropReason.MaxRetried); + confirm(dropped, false); + receiveQuota.onErrorSignal(now); + } if (!unconfirmedPacketIds.isEmpty()) { scheduleResend(); } diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3TransientSessionHandlerTest.java b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3TransientSessionHandlerTest.java index 552adc4bf..69ac60aa6 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3TransientSessionHandlerTest.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3TransientSessionHandlerTest.java @@ -1730,4 +1730,52 @@ public void handlerAdded(ChannelHandlerContext ctx) throws Exception { verify(eventCollector, atLeast(1)).report(argThat(e -> e instanceof QoS1Confirmed c && !c.delivered())); } + + @Test + public void resendDroppingHeadAbortsDropOfRemaining() { + when(settingProvider.provide(eq(ResendTimeoutSeconds), anyString())).thenReturn(1); + when(settingProvider.provide(eq(MaxResendTimes), anyString())).thenReturn(1); + when(settingProvider.provide(eq(ReceivingMaximum), anyString())).thenReturn(2); + + channel.pipeline().removeLast(); + channel.pipeline().addLast(new ChannelDuplexHandler() { + @Override + public void handlerAdded(ChannelHandlerContext ctx) throws Exception { + super.handlerAdded(ctx); + ctx.pipeline().addLast( + MQTT3TransientSessionHandler.builder().settings(new TenantSettings(tenantId, settingProvider)) + .tenantMeter(tenantMeter).oomCondition(oomCondition).userSessionId(userSessionId(clientInfo)) + .keepAliveTimeSeconds(120).clientInfo(clientInfo).willMessage(null).ctx(ctx).build()); + ctx.pipeline().remove(this); + } + }); + transientSessionHandler = (MQTT3TransientSessionHandler) channel.pipeline().last(); + + mockCheckPermission(true); + mockDistMatch(true); + transientSessionHandler.subscribe(System.nanoTime(), topicFilter, QoS.AT_LEAST_ONCE); + channel.runPendingTasks(); + ArgumentCaptor longCaptor = ArgumentCaptor.forClass(Long.class); + verify(localDistService).match(anyLong(), eq(topicFilter), longCaptor.capture(), any()); + + transientSessionHandler.publish(s2cMessageList(topic, 2, QoS.AT_LEAST_ONCE), + Collections.singleton(new IMQTTTransientSession.MatchedTopicFilter(topicFilter, longCaptor.getValue()))); + channel.runPendingTasks(); + + assertNotNull(channel.readOutbound()); + assertNotNull(channel.readOutbound()); + // deliberately never ack: both stay in flight, so the head is dropped while another entry remains + + for (int i = 0; i < 6; i++) { + testTicker.advanceTimeBy(2, TimeUnit.SECONDS); + channel.advanceTimeBy(2, TimeUnit.SECONDS); + channel.runScheduledPendingTasks(); + channel.runPendingTasks(); + channel.flushOutbound(); + } + + // both in-flight messages have exceeded maxResendTimes and must be dropped + verify(eventCollector, times(2)).report(argThat(e -> + e instanceof QoS1Dropped d && d.reason() == DropReason.MaxRetried)); + } }