@ -18,8 +18,9 @@ package org.thingsboard.mqtt;
import io.netty.channel.EventLoop ;
import io.netty.channel.EventLoop ;
import io.netty.handler.codec.mqtt.MqttFixedHeader ;
import io.netty.handler.codec.mqtt.MqttFixedHeader ;
import io.netty.handler.codec.mqtt.MqttMessage ;
import io.netty.handler.codec.mqtt.MqttMessage ;
import io.netty.util.concurrent.ScheduledFuture ;
import io.netty.handler.codec.mqtt.MqttMessageType ;
import io.netty.handler.codec.mqtt.MqttMessageType ;
import io.netty.handler.codec.mqtt.MqttQoS ;
import io.netty.util.concurrent.ScheduledFuture ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.TimeUnit ;
import java.util.function.BiConsumer ;
import java.util.function.BiConsumer ;
@ -45,7 +46,10 @@ final class RetransmissionHandler<T extends MqttMessage> {
private void startTimer ( EventLoop eventLoop ) {
private void startTimer ( EventLoop eventLoop ) {
this . timer = eventLoop . schedule ( ( ) - > {
this . timer = eventLoop . schedule ( ( ) - > {
this . timeout + = 5 ;
this . timeout + = 5 ;
boolean isDup = this . originalMessage . fixedHeader ( ) . messageType ( ) = = MqttMessageType . PUBLISH ? true : this . originalMessage . fixedHeader ( ) . isDup ( ) ;
boolean isDup = this . originalMessage . fixedHeader ( ) . isDup ( ) ;
if ( this . originalMessage . fixedHeader ( ) . messageType ( ) = = MqttMessageType . PUBLISH & & this . originalMessage . fixedHeader ( ) . qosLevel ( ) ! = MqttQoS . AT_MOST_ONCE ) {
isDup = true ;
}
MqttFixedHeader fixedHeader = new MqttFixedHeader ( this . originalMessage . fixedHeader ( ) . messageType ( ) , isDup , this . originalMessage . fixedHeader ( ) . qosLevel ( ) , this . originalMessage . fixedHeader ( ) . isRetain ( ) , this . originalMessage . fixedHeader ( ) . remainingLength ( ) ) ;
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 ) ;
handler . accept ( fixedHeader , originalMessage ) ;
startTimer ( eventLoop ) ;
startTimer ( eventLoop ) ;