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(); 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);