From f96d54f4a54cb64379d581bf4dde3095f77eab67 Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Thu, 24 Sep 2026 00:41:48 +0800 Subject: [PATCH 1/2] fix(mqtt5): apply isDup flag in buildMqttPubMessage (three builder paths) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Re-delivered PUBLISH packets were built with DUP=0, causing standards-compliant clients (e.g. HiveMQ MQTT client) to treat the re-delivery as a protocol error and disconnect — message loss in server-side bridge scenarios. Per MQTT-3.3.1-1, the DUP flag MUST be set on re-delivery. The isDup parameter was accepted but never applied in any of the three builder paths (topic-alias first-use / topic-alias reuse / no-alias). Fixes #283 Signed-off-by: DanXie <517964478@qq.com> --- .../apache/bifromq/mqtt/handler/v5/MQTT5ProtocolHelper.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/v5/MQTT5ProtocolHelper.java b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/v5/MQTT5ProtocolHelper.java index 830671bf5..6cddbd1be 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/v5/MQTT5ProtocolHelper.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/v5/MQTT5ProtocolHelper.java @@ -474,10 +474,10 @@ public MqttPublishMessage buildMqttPubMessage(int packetId, RoutedMessage messag senderTopicAliasManager.tryAlias(message.topic()); if (aliasCreationResult.isPresent()) { if (aliasCreationResult.get().isFirstTime()) { - return MQTT5MessageBuilders.pub().packetId(packetId).setupAlias(true) + return MQTT5MessageBuilders.pub().packetId(packetId).setupAlias(true).dup(isDup) .topicAlias(aliasCreationResult.get().alias()).message(message).build(); } else { - return MQTT5MessageBuilders.pub().packetId(packetId).topicAlias(aliasCreationResult.get().alias()) + return MQTT5MessageBuilders.pub().packetId(packetId).dup(isDup).topicAlias(aliasCreationResult.get().alias()) .message(message).build(); } } @@ -491,6 +491,7 @@ public MqttPublishMessage buildMqttPubMessage(int packetId, RoutedMessage messag message.hlc()); return MQTT5MessageBuilders.pub() .packetId(packetId) + .dup(isDup) .message(message) .extraUserProps(extraUserProps) .build(); From 66d1d19b896786f009b3d2368ddd5be4a9a8df9e Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Thu, 24 Sep 2026 02:27:56 +0800 Subject: [PATCH 2/2] test(mqtt5): assert DUP=1 on QoS1 re-delivery Adds qos1RedeliveryCarriesDupFlag to TransientSessionHandlerTest: publish one QoS1 message, keep it in flight, advance past ResendTimeoutSeconds, and assert the re-delivered PUBLISH carries DUP=1 per [MQTT-3.3.1-1] while going through the topic-alias reuse builder path. The handler is rebuilt with stubbed settings because TenantSettings reads its values once at construction. Verified against the unfixed MQTT5ProtocolHelper the new test fails with 'DUP=1 expected true but found false'; with the fix the full TransientSessionHandlerTest (52 tests) passes. --- .../v5/TransientSessionHandlerTest.java | 73 +++++++++++++++++++ 1 file changed, 73 insertions(+) diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v5/TransientSessionHandlerTest.java b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v5/TransientSessionHandlerTest.java index dc697f023..1837cbedd 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v5/TransientSessionHandlerTest.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v5/TransientSessionHandlerTest.java @@ -62,9 +62,11 @@ import static org.apache.bifromq.plugin.eventcollector.EventType.SUB_ACKED; import static org.apache.bifromq.plugin.eventcollector.EventType.UNSUB_ACKED; import static org.apache.bifromq.plugin.eventcollector.EventType.UNSUB_ACTION_DISALLOW; +import static org.apache.bifromq.plugin.settingprovider.Setting.MaxResendTimes; import static org.apache.bifromq.plugin.settingprovider.Setting.MsgPubPerSec; import static org.apache.bifromq.plugin.settingprovider.Setting.ReceivingMaximum; import static org.apache.bifromq.plugin.settingprovider.Setting.RetainEnabled; +import static org.apache.bifromq.plugin.settingprovider.Setting.ResendTimeoutSeconds; import static org.apache.bifromq.retain.rpc.proto.RetainReply.Result.CLEARED; import static org.apache.bifromq.retain.rpc.proto.RetainReply.Result.ERROR; import static org.apache.bifromq.retain.rpc.proto.RetainReply.Result.RETAINED; @@ -82,6 +84,7 @@ import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; @@ -1077,6 +1080,76 @@ public void qos1PubAndAck() { QOS1_CONFIRMED, QOS1_CONFIRMED); } + @Test + public void qos1RedeliveryCarriesDupFlag() { + when(settingProvider.provide(eq(ResendTimeoutSeconds), anyString())).thenReturn(1); + when(settingProvider.provide(eq(MaxResendTimes), anyString())).thenReturn(5); + + // rebuild the handler so TenantSettings picks up the stubbed settings + // (its fields are final and read once at construction) + MqttProperties mqttProperties = new MqttProperties(); + mqttProperties.add(new MqttProperties.IntegerProperty(TOPIC_ALIAS_MAXIMUM.value(), 10)); + channel.pipeline().removeLast(); + channel.pipeline().addLast(new ChannelDuplexHandler() { + @Override + public void handlerAdded(ChannelHandlerContext ctx) throws Exception { + super.handlerAdded(ctx); + ctx.pipeline().addLast(MQTT5TransientSessionHandler.builder() + .settings(new TenantSettings(tenantId, settingProvider)) + .tenantMeter(tenantMeter) + .oomCondition(oomCondition) + .connMsg(MqttMessageBuilders.connect() + .protocolVersion(MqttVersion.MQTT_5) + .properties(mqttProperties) + .build()) + .userSessionId(userSessionId(clientInfo)) + .keepAliveTimeSeconds(120) + .clientInfo(clientInfo) + .willMessage(null).ctx(ctx) + .build()); + ctx.pipeline().remove(this); + } + }); + transientSessionHandler = (MQTT5TransientSessionHandler) 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(s2cMQTT5MessageList(topic, 1, QoS.AT_LEAST_ONCE), + Collections.singleton(new IMQTTTransientSession.MatchedTopicFilter(topicFilter, longCaptor.getValue()))); + channel.runPendingTasks(); + + // first delivery: not a re-delivery, DUP must be 0 + MqttPublishMessage first = channel.readOutbound(); + assertNotNull(first); + assertFalse(first.fixedHeader().isDup()); + int packetId = first.variableHeader().packetId(); + // no PUBACK — keep the message in flight so the resend timer fires + + testTicker.advanceTimeBy(2, TimeUnit.SECONDS); + channel.advanceTimeBy(2, TimeUnit.SECONDS); + channel.runScheduledPendingTasks(); + channel.runPendingTasks(); + channel.flushOutbound(); + + // re-delivery of the same packet id: DUP must be 1 per [MQTT-3.3.1-1] + MqttPublishMessage redelivered = channel.readOutbound(); + assertNotNull(redelivered); + assertEquals(redelivered.variableHeader().packetId(), packetId); + assertTrue(redelivered.fixedHeader().isDup(), + "re-delivered PUBLISH must carry DUP=1 (MQTT-3.3.1-1)"); + // redelivery goes through the topic-alias reuse builder path + assertEquals(redelivered.variableHeader().properties() + .getProperties(TOPIC_ALIAS.value()).get(0).value(), 1); + + channel.writeInbound(MQTTMessageUtils.pubAckMessage(packetId)); + channel.runPendingTasks(); + } + @Test public void qoS2PubAndRel() { mockCheckPermission(true);