From dda3124835bf559118abe71a4566216e3d1388c3 Mon Sep 17 00:00:00 2001 From: blackstar-baba <535650957@qq.com> Date: Sun, 24 May 2020 23:22:51 +0800 Subject: [PATCH] fix bug: Resending the message causes the client to interrupt MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit In the mqtt 3.1 protocol or mqtt 3.1.1 protocol, except for messages of type pubish, the isDup of other messages must be zero,Sending a subscription message multiple times causes the client to reconnect,therefore, resending the subscription message with isDup true will cause the server to think that the message is abnormal, thus interrupting the connection with the client --- .../main/java/org/thingsboard/mqtt/RetransmissionHandler.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/RetransmissionHandler.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/RetransmissionHandler.java index 7d8d063d89..8646d28cc8 100644 --- a/netty-mqtt/src/main/java/org/thingsboard/mqtt/RetransmissionHandler.java +++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/RetransmissionHandler.java @@ -19,6 +19,7 @@ import io.netty.channel.EventLoop; import io.netty.handler.codec.mqtt.MqttFixedHeader; import io.netty.handler.codec.mqtt.MqttMessage; import io.netty.util.concurrent.ScheduledFuture; +import io.netty.handler.codec.mqtt.MqttMessageType; import java.util.concurrent.TimeUnit; import java.util.function.BiConsumer; @@ -44,7 +45,8 @@ final class RetransmissionHandler { private void startTimer(EventLoop eventLoop){ this.timer = eventLoop.schedule(() -> { this.timeout += 5; - MqttFixedHeader fixedHeader = new MqttFixedHeader(this.originalMessage.fixedHeader().messageType(), true, this.originalMessage.fixedHeader().qosLevel(), this.originalMessage.fixedHeader().isRetain(), this.originalMessage.fixedHeader().remainingLength()); + boolean isDup = this.originalMessage.fixedHeader().messageType() == MqttMessageType.PUBLISH ? true : this.originalMessage.fixedHeader().isDup(); + MqttFixedHeader fixedHeader = new MqttFixedHeader(this.originalMessage.fixedHeader().messageType(), isDup, this.originalMessage.fixedHeader().qosLevel(), this.originalMessage.fixedHeader().isRetain(), this.originalMessage.fixedHeader().remainingLength()); handler.accept(fixedHeader, originalMessage); startTimer(eventLoop); }, timeout, TimeUnit.SECONDS);