diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
index 26c82a33de..61d9586095 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
@@ -32,6 +32,7 @@ import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.DeviceStateManager;
import org.thingsboard.rule.engine.api.MailService;
+import org.thingsboard.rule.engine.api.MqttClientSettings;
import org.thingsboard.rule.engine.api.NotificationCenter;
import org.thingsboard.rule.engine.api.SmsService;
import org.thingsboard.rule.engine.api.notification.SlackService;
@@ -639,6 +640,10 @@ public class ActorSystemContext {
@Getter
private long cfCalculationResultTimeout;
+ @Autowired
+ @Getter
+ private MqttClientSettings mqttClientSettings;
+
@Getter
@Setter
private TbActorSystem actorSystem;
diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
index 033e10ca9a..3fb28aee38 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
@@ -23,14 +23,15 @@ import org.bouncycastle.util.Arrays;
import org.thingsboard.common.util.DebugModeUtil;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor;
+import org.thingsboard.rule.engine.api.DeviceStateManager;
import org.thingsboard.rule.engine.api.MailService;
+import org.thingsboard.rule.engine.api.MqttClientSettings;
import org.thingsboard.rule.engine.api.NotificationCenter;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
import org.thingsboard.rule.engine.api.RuleEngineApiUsageStateService;
import org.thingsboard.rule.engine.api.RuleEngineAssetProfileCache;
import org.thingsboard.rule.engine.api.RuleEngineCalculatedFieldQueueService;
import org.thingsboard.rule.engine.api.RuleEngineDeviceProfileCache;
-import org.thingsboard.rule.engine.api.DeviceStateManager;
import org.thingsboard.rule.engine.api.RuleEngineRpcService;
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.rule.engine.api.ScriptEngine;
@@ -1010,13 +1011,17 @@ public class DefaultTbContext implements TbContext {
return mainCtx.getAuditLogService();
}
+ @Override
+ public MqttClientSettings getMqttClientSettings() {
+ return mainCtx.getMqttClientSettings();
+ }
+
private TbMsgMetaData getActionMetaData(RuleNodeId ruleNodeId) {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("ruleNodeId", ruleNodeId.toString());
return metaData;
}
-
@Override
public void schedule(Runnable runnable, long delay, TimeUnit timeUnit) {
mainCtx.getScheduler().schedule(runnable, delay, timeUnit);
diff --git a/application/src/main/java/org/thingsboard/server/config/mqtt/MqttClientRetransmissionSettingsComponent.java b/application/src/main/java/org/thingsboard/server/config/mqtt/MqttClientRetransmissionSettingsComponent.java
new file mode 100644
index 0000000000..33e9358d2b
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/config/mqtt/MqttClientRetransmissionSettingsComponent.java
@@ -0,0 +1,37 @@
+/**
+ * Copyright © 2016-2025 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.config.mqtt;
+
+import jakarta.validation.constraints.PositiveOrZero;
+import lombok.Data;
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.validation.annotation.Validated;
+
+@Data
+@Validated
+@Configuration
+@ConfigurationProperties(prefix = "mqtt.client.retransmission")
+public class MqttClientRetransmissionSettingsComponent {
+
+ @PositiveOrZero
+ private int maxAttempts;
+ @PositiveOrZero
+ private long initialDelayMillis;
+ @PositiveOrZero
+ private double jitterFactor;
+
+}
diff --git a/application/src/main/java/org/thingsboard/server/config/mqtt/MqttClientSettingsComponent.java b/application/src/main/java/org/thingsboard/server/config/mqtt/MqttClientSettingsComponent.java
new file mode 100644
index 0000000000..25df212925
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/config/mqtt/MqttClientSettingsComponent.java
@@ -0,0 +1,47 @@
+/**
+ * Copyright © 2016-2025 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.config.mqtt;
+
+import lombok.EqualsAndHashCode;
+import lombok.RequiredArgsConstructor;
+import lombok.ToString;
+import org.springframework.context.annotation.Configuration;
+import org.thingsboard.rule.engine.api.MqttClientSettings;
+
+@ToString
+@EqualsAndHashCode
+@Configuration
+@RequiredArgsConstructor
+public class MqttClientSettingsComponent implements MqttClientSettings {
+
+ private final MqttClientRetransmissionSettingsComponent retransmissionSettingsComponent;
+
+ @Override
+ public int getRetransmissionMaxAttempts() {
+ return retransmissionSettingsComponent.getMaxAttempts();
+ }
+
+ @Override
+ public long getRetransmissionInitialDelayMillis() {
+ return retransmissionSettingsComponent.getInitialDelayMillis();
+ }
+
+ @Override
+ public double getRetransmissionJitterFactor() {
+ return retransmissionSettingsComponent.getJitterFactor();
+ }
+
+}
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index d1910399f0..df33a92857 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -1952,3 +1952,27 @@ mobileApp:
googlePlayLink: "${TB_MOBILE_APP_GOOGLE_PLAY_LINK:https://play.google.com/store/apps/details?id=org.thingsboard.demo.app}"
# Link to App Store for Thingsboard Live mobile application
appStoreLink: "${TB_MOBILE_APP_APP_STORE_LINK:https://apps.apple.com/us/app/thingsboard-live/id1594355695}"
+
+mqtt:
+ # MQTT client configuration parameters
+ client:
+ # Parameters that control the retransmission mechanism.
+ # This mechanism only applies to the handling of MQTT Publish, Subscribe, Unsubscribe and Pubrel messages.
+ # With the updated default settings:
+ # - After sending the message, wait approximately 5000 ms (± jitter) for the 1st attempt.
+ # - The 2nd attempt will occur after roughly 5000 * 2 = 10,000 ms (± jitter).
+ # - The 3rd attempt will occur after roughly 5000 * 4 = 20,000 ms (± jitter).
+ # - The 4th "attempt" will not actually perform a retransmission.
+ # Instead, the system will detect that the maximum number of attempts has been reached and drop the pending message.
+ retransmission:
+ # Maximum number of retransmission attempts allowed.
+ # If the attempt count exceeds this value, retransmissions will stop and the pending message will be dropped.
+ max_attempts: "${TB_MQTT_CLIENT_RETRANSMISSION_MAX_ATTEMPTS:3}"
+ # Base delay (in milliseconds) before the first retransmission attempt, measured from the moment the message is sent.
+ # Subsequent delays are calculated using exponential backoff.
+ # This base delay is also used as the reference value for applying jitter.
+ initial_delay_millis: "${TB_MQTT_CLIENT_RETRANSMISSION_INITIAL_DELAY_MILLIS:5000}"
+ # Jitter factor applied to the calculated retransmission delay.
+ # The actual delay is randomized within a range defined by multiplying the base delay by a factor between (1 - jitter_factor) and (1 + jitter_factor).
+ # For example, a jitter_factor of 0.15 means the actual delay may vary by up to ±15% of the base delay.
+ jitter_factor: "${TB_MQTT_CLIENT_RETRANSMISSION_JITTER_FACTOR:0.15}"
diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java
index f2e721bdcb..e00b0828c9 100644
--- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java
+++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java
@@ -55,8 +55,8 @@ public class ContainerTestSuite {
private static final String TB_JS_EXECUTOR_LOG_REGEXP = ".*template started.*";
private static final Duration CONTAINER_STARTUP_TIMEOUT = Duration.ofSeconds(400);
- private DockerComposeContainer> testContainer;
- private ThingsBoardDbInstaller installTb;
+ private DockerComposeContainer> testContainer;
+ private ThingsBoardDbInstaller installTb;
private boolean isActive;
private static ContainerTestSuite containerTestSuite;
@@ -194,7 +194,7 @@ public class ContainerTestSuite {
setActive(true);
} catch (Exception e) {
log.error("Failed to create test container", e);
- fail("Failed to create test container");
+ fail("Failed to create test container", e);
}
}
@@ -263,7 +263,7 @@ public class ContainerTestSuite {
log.info("Trying to delete temp dir {}", targetDir);
FileUtils.deleteDirectory(new File(targetDir));
} catch (IOException e) {
- log.error("Can't delete temp directory " + targetDir, e);
+ log.error("Can't delete temp directory {}", targetDir, e);
}
}
@@ -286,8 +286,8 @@ public class ContainerTestSuite {
FileUtils.writeStringToFile(file, outputContent, StandardCharsets.UTF_8);
assertThat(FileUtils.readFileToString(file, StandardCharsets.UTF_8), is(outputContent));
} catch (IOException e) {
- log.error("failed to update file " + sourceFilename, e);
- fail("failed to update file");
+ log.error("failed to update file {}", sourceFilename, e);
+ fail("failed to update file", e);
}
}
diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java
index ebdfb4e3c9..1b1ed9bf0f 100644
--- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java
+++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java
@@ -42,7 +42,6 @@ import org.thingsboard.mqtt.MqttClient;
import org.thingsboard.mqtt.MqttClientCallback;
import org.thingsboard.mqtt.MqttClientConfig;
import org.thingsboard.mqtt.MqttHandler;
-import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.DeviceProfileProvisionType;
@@ -82,7 +81,6 @@ import java.util.concurrent.TimeoutException;
import static org.assertj.core.api.Assertions.assertThat;
import static org.testng.Assert.assertNotNull;
import static org.testng.Assert.fail;
-import static org.thingsboard.server.common.data.DataConstants.DEVICE;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE;
import static org.thingsboard.server.msa.prototypes.DevicePrototypes.defaultDevicePrototype;
@@ -301,7 +299,7 @@ public class MqttClientTest extends AbstractContainerTest {
assertThat(Objects.requireNonNull(requestFromServer).getMessage()).isEqualTo("{\"method\":\"getValue\",\"params\":true}");
- Integer requestId = Integer.valueOf(Objects.requireNonNull(requestFromServer).getTopic().substring("v1/devices/me/rpc/request/".length()));
+ int requestId = Integer.parseInt(Objects.requireNonNull(requestFromServer).getTopic().substring("v1/devices/me/rpc/request/".length()));
JsonObject clientResponse = new JsonObject();
clientResponse.addProperty("response", "someResponse");
// Send a response to the server's RPC request
@@ -340,7 +338,7 @@ public class MqttClientTest extends AbstractContainerTest {
assertThat(Objects.requireNonNull(requestFromServer).getMessage()).isEqualTo("{\"method\":\"getValue\",\"params\":true}");
- Integer requestId = Integer.valueOf(Objects.requireNonNull(requestFromServer).getTopic().substring("v1/devices/me/rpc/request/".length()));
+ int requestId = Integer.parseInt(Objects.requireNonNull(requestFromServer).getTopic().substring("v1/devices/me/rpc/request/".length()));
JsonObject clientResponse = new JsonObject();
clientResponse.addProperty("response", "someResponse");
// Send a response to the server's RPC request
@@ -520,13 +518,13 @@ public class MqttClientTest extends AbstractContainerTest {
mqttClient.on("/provision/response", listener, MqttQoS.AT_LEAST_ONCE).get(3 * timeoutMultiplier, TimeUnit.SECONDS);
TimeUnit.SECONDS.sleep(2 * timeoutMultiplier);
assertThat(subAckResult[0]).isNotNull();
- assertThat(MqttReasonCodes.SubAck.GRANTED_QOS_1.equals(subAckResult[0]));
+ assertThat(MqttReasonCodes.SubAck.GRANTED_QOS_1).isEqualTo(subAckResult[0]);
subAckResult[0] = null;
mqttClient.on("v1/devices/me/attributes", listener, MqttQoS.AT_LEAST_ONCE).get(3 * timeoutMultiplier, TimeUnit.SECONDS);
TimeUnit.SECONDS.sleep(2 * timeoutMultiplier);
assertThat(subAckResult[0]).isNotNull();
- assertThat(MqttReasonCodes.SubAck.TOPIC_FILTER_INVALID.equals(subAckResult[0]));
+ assertThat(MqttReasonCodes.SubAck.TOPIC_FILTER_INVALID).isEqualTo(subAckResult[0]);
testRestClient.deleteDeviceIfExists(device.getId());
updateDeviceProfileWithProvisioningStrategy(deviceProfile, DeviceProfileProvisionType.DISABLED);
@@ -596,7 +594,7 @@ public class MqttClientTest extends AbstractContainerTest {
.await()
.alias("Check device disconnect.")
.atMost(TIMEOUT*timeoutMultiplier, TimeUnit.SECONDS)
- .until(() -> returnCodeByteValue.size() > 0);
+ .until(() -> !returnCodeByteValue.isEmpty());
assertThat(returnCodeByteValueSecondClient).isEmpty();
assertThat(returnCodeByteValue).isNotEmpty();
@@ -663,7 +661,7 @@ public class MqttClientTest extends AbstractContainerTest {
.stream()
.filter(RuleChain::isRoot)
.findFirst();
- if (!defaultRuleChain.isPresent()) {
+ if (defaultRuleChain.isEmpty()) {
fail("Root rule chain wasn't found");
}
return defaultRuleChain.get().getId();
@@ -717,6 +715,7 @@ public class MqttClientTest extends AbstractContainerTest {
clientConfig.setClientId("MQTT client from test");
clientConfig.setUsername(username);
clientConfig.setProtocolVersion(mqttVersion);
+ clientConfig.setRetransmissionConfig(new MqttClientConfig.RetransmissionConfig(3, 5000L, 0.15d)); // same as defaults in thingsboard.yml as of time of this writing
MqttClient mqttClient = MqttClient.create(clientConfig, listener, handlerExecutor);
if (connect) {
mqttClient.connect(TRANSPORT_HOST, TRANSPORT_PORT).get();
diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java
index cc587fbbd5..fd9d4f557d 100644
--- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java
+++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java
@@ -39,7 +39,6 @@ import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.mqtt.MqttClient;
import org.thingsboard.mqtt.MqttClientConfig;
import org.thingsboard.mqtt.MqttHandler;
-import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.DeviceId;
@@ -65,7 +64,6 @@ import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
-import static org.thingsboard.server.common.data.DataConstants.DEVICE;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE;
import static org.thingsboard.server.msa.prototypes.DevicePrototypes.defaultGatewayPrototype;
@@ -76,7 +74,6 @@ public class MqttGatewayClientTest extends AbstractContainerTest {
private MqttClient mqttClient;
private Device createdDevice;
private MqttMessageListener listener;
- private JsonParser jsonParser = new JsonParser();
AbstractListeningExecutor handlerExecutor;
@@ -100,7 +97,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest {
}
@AfterMethod
- public void removeGateway() {
+ public void removeGateway() {
testRestClient.deleteDeviceIfExists(this.gatewayDevice.getId());
testRestClient.deleteDeviceIfExists(this.createdDevice.getId());
this.listener = null;
@@ -197,7 +194,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest {
mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(requestData.toString().getBytes())).get();
event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS);
- JsonObject responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject();
+ JsonObject responseData = JsonParser.parseString(Objects.requireNonNull(event).getMessage()).getAsJsonObject();
assertThat(responseData.has("value")).isTrue();
assertThat(responseData.get("value").getAsString()).isEqualTo(sharedAttributes.get("attr1").getAsString());
@@ -213,7 +210,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest {
mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get();
mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(requestData.toString().getBytes())).get();
event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS);
- responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject();
+ responseData = JsonParser.parseString(Objects.requireNonNull(event).getMessage()).getAsJsonObject();
assertThat(responseData.has("values")).isTrue();
assertThat(responseData.get("values").getAsJsonObject().get("attr1").getAsString()).isEqualTo(sharedAttributes.get("attr1").getAsString());
@@ -231,7 +228,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest {
mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get();
mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(requestData.toString().getBytes())).get();
event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS);
- responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject();
+ responseData = JsonParser.parseString(Objects.requireNonNull(event).getMessage()).getAsJsonObject();
assertThat(responseData.has("values")).isTrue();
assertThat(responseData.get("values").getAsJsonObject().get("attr1").getAsString()).isEqualTo(sharedAttributes.get("attr1").getAsString());
@@ -390,7 +387,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest {
mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(gatewayAttributesRequest.toString().getBytes())).get();
MqttEvent clientAttributeEvent = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS);
assertThat(clientAttributeEvent).isNotNull();
- JsonObject responseMessage = new JsonParser().parse(Objects.requireNonNull(clientAttributeEvent).getMessage()).getAsJsonObject();
+ JsonObject responseMessage = JsonParser.parseString(Objects.requireNonNull(clientAttributeEvent).getMessage()).getAsJsonObject();
assertThat(responseMessage.get("id").getAsInt()).isEqualTo(messageId);
assertThat(responseMessage.get("device").getAsString()).isEqualTo(createdDevice.getName());
@@ -427,6 +424,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest {
clientConfig.setOwnerId(getOwnerId());
clientConfig.setClientId("MQTT client from test");
clientConfig.setUsername(deviceCredentials.getCredentialsId());
+ clientConfig.setRetransmissionConfig(new MqttClientConfig.RetransmissionConfig(3, 5000L, 0.15d)); // same as defaults in thingsboard.yml as of time of this writing
MqttClient mqttClient = MqttClient.create(clientConfig, listener, handlerExecutor);
mqttClient.connect("localhost", 1883).get();
return mqttClient;
diff --git a/netty-mqtt/pom.xml b/netty-mqtt/pom.xml
index 7e0a00eee6..8ce2c68310 100644
--- a/netty-mqtt/pom.xml
+++ b/netty-mqtt/pom.xml
@@ -87,6 +87,26 @@
awaitility
test
+
+ org.testcontainers
+ testcontainers
+ test
+
+
+ org.testcontainers
+ junit-jupiter
+ test
+
+
+ software.xdev
+ testcontainers-junit4-mock
+ test
+
+
+ org.testcontainers
+ hivemq
+ test
+
diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MaxRetransmissionsReachedException.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MaxRetransmissionsReachedException.java
new file mode 100644
index 0000000000..3d483dd541
--- /dev/null
+++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MaxRetransmissionsReachedException.java
@@ -0,0 +1,24 @@
+/**
+ * Copyright © 2016-2025 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.mqtt;
+
+public class MaxRetransmissionsReachedException extends RuntimeException {
+
+ public MaxRetransmissionsReachedException(String message) {
+ super(message);
+ }
+
+}
diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttChannelHandler.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttChannelHandler.java
index ad976c848a..9686a2b1d7 100644
--- a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttChannelHandler.java
+++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttChannelHandler.java
@@ -57,7 +57,7 @@ final class MqttChannelHandler extends SimpleChannelInboundHandler
}
@Override
- protected void channelRead0(ChannelHandlerContext ctx, MqttMessage msg) throws Exception {
+ protected void channelRead0(ChannelHandlerContext ctx, MqttMessage msg) {
if (msg.decoderResult().isSuccess()) {
switch (msg.fixedHeader().messageType()) {
case CONNACK:
@@ -120,6 +120,7 @@ final class MqttChannelHandler extends SimpleChannelInboundHandler
this.client.getClientConfig().getUsername(),
this.client.getClientConfig().getPassword() != null ? this.client.getClientConfig().getPassword().getBytes(CharsetUtil.UTF_8) : null
);
+ log.debug("{} Sending CONNECT", client.getClientConfig().getOwnerId());
ctx.channel().writeAndFlush(new MqttConnectMessage(fixedHeader, variableHeader, payload));
}
@@ -173,6 +174,7 @@ final class MqttChannelHandler extends SimpleChannelInboundHandler
}
private void handleConack(Channel channel, MqttConnAckMessage message) {
+ log.debug("{} Handling CONNACK", client.getClientConfig().getOwnerId());
switch (message.variableHeader().connectReturnCode()) {
case CONNECTION_ACCEPTED:
this.connectFuture.setSuccess(new MqttConnectResult(true, MqttConnectReturnCode.CONNECTION_ACCEPTED, channel.closeFuture()));
@@ -219,9 +221,9 @@ final class MqttChannelHandler extends SimpleChannelInboundHandler
}
pendingSubscription.onSubackReceived();
for (MqttPendingSubscription.MqttPendingHandler handler : pendingSubscription.getHandlers()) {
- MqttSubscription subscription = new MqttSubscription(pendingSubscription.getTopic(), handler.getHandler(), handler.isOnce());
+ MqttSubscription subscription = new MqttSubscription(pendingSubscription.getTopic(), handler.handler(), handler.once());
this.client.getSubscriptions().put(pendingSubscription.getTopic(), subscription);
- this.client.getHandlerToSubscription().put(handler.getHandler(), subscription);
+ this.client.getHandlerToSubscription().put(handler.handler(), subscription);
}
this.client.getPendingSubscribeTopics().remove(pendingSubscription.getTopic());
@@ -282,17 +284,16 @@ final class MqttChannelHandler extends SimpleChannelInboundHandler
}
private void handlePuback(MqttPubAckMessage message) {
- MqttPendingPublish pendingPublish = this.client.getPendingPublishes().get(message.variableHeader().messageId());
- if (pendingPublish == null) {
- return;
- }
- pendingPublish.getFuture().setSuccess(null);
- pendingPublish.onPubackReceived();
- this.client.getPendingPublishes().remove(message.variableHeader().messageId());
- pendingPublish.getPayload().release();
- if (this.client.getCallback() != null) {
- this.client.getCallback().onPubAck(message);
- }
+ log.trace("{} Handling PUBACK", client.getClientConfig().getOwnerId());
+ client.getPendingPublishes().computeIfPresent(message.variableHeader().messageId(), (__, pendingPublish) -> {
+ pendingPublish.getFuture().setSuccess(null);
+ pendingPublish.onPubackReceived();
+ pendingPublish.getPayload().release();
+ if (client.getCallback() != null) {
+ client.getCallback().onPubAck(message);
+ }
+ return null;
+ });
}
private void handlePubrec(Channel channel, MqttMessage message) {
@@ -335,6 +336,7 @@ final class MqttChannelHandler extends SimpleChannelInboundHandler
}
private void handleDisconnect(MqttMessage message) {
+ log.debug("{} Handling DISCONNECT", client.getClientConfig().getOwnerId());
if (this.client.getCallback() != null) {
this.client.getCallback().onDisconnect(message);
}
diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClient.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClient.java
index db0459e08a..4d845320e8 100644
--- a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClient.java
+++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClient.java
@@ -184,7 +184,7 @@ public interface MqttClient {
* @param config The config object to use while looking for settings
* @param defaultHandler The handler for incoming messages that do not match any topic subscriptions
*/
- static MqttClient create(MqttClientConfig config, MqttHandler defaultHandler, ListeningExecutor handlerExecutor){
+ static MqttClient create(MqttClientConfig config, MqttHandler defaultHandler, ListeningExecutor handlerExecutor) {
return new MqttClientImpl(config, defaultHandler, handlerExecutor);
}
diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientConfig.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientConfig.java
index 41df077d71..24feb3e58e 100644
--- a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientConfig.java
+++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientConfig.java
@@ -47,6 +47,26 @@ public final class MqttClientConfig {
private long reconnectDelay = 1L;
private int maxBytesInMessage = 8092;
+ @Getter
+ @Setter
+ private RetransmissionConfig retransmissionConfig;
+
+ public record RetransmissionConfig(int maxAttempts, long initialDelayMillis, double jitterFactor) {
+
+ public RetransmissionConfig {
+ if (maxAttempts < 0) {
+ throw new IllegalArgumentException("Max retransmission attempts (maxAttempts) must be zero or greater, but was " + maxAttempts);
+ }
+ if (initialDelayMillis < 0) {
+ throw new IllegalArgumentException("Initial retransmission delay (initialDelayMillis) must be zero or greater, but was " + initialDelayMillis);
+ }
+ if (jitterFactor < 0) {
+ throw new IllegalArgumentException("Jitter factor (jitterFactor) must be zero or greater, but was " + jitterFactor);
+ }
+ }
+
+ }
+
public MqttClientConfig() {
this(null);
}
diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java
index 47eae565dc..ee07752db3 100644
--- a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java
+++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java
@@ -17,6 +17,7 @@ package org.thingsboard.mqtt;
import com.google.common.collect.HashMultimap;
import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Sets;
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.channel.Channel;
@@ -384,8 +385,33 @@ final class MqttClientImpl implements MqttClient {
MqttFixedHeader fixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, false, qos, retain, 0);
MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader(topic, getNewMessageId().messageId());
MqttPublishMessage message = new MqttPublishMessage(fixedHeader, variableHeader, payload);
- MqttPendingPublish pendingPublish = new MqttPendingPublish(variableHeader.packetId(), future,
- payload.retain(), message, qos, () -> !pendingPublishes.containsKey(variableHeader.packetId()));
+
+ final var pendingPublish = MqttPendingPublish.builder()
+ .messageId(variableHeader.packetId())
+ .future(future)
+ .payload(payload.retain())
+ .message(message)
+ .qos(qos)
+ .ownerId(clientConfig.getOwnerId())
+ .retransmissionConfig(clientConfig.getRetransmissionConfig())
+ .pendingOperation(new PendingOperation() {
+ @Override
+ public boolean isCancelled() {
+ return !pendingPublishes.containsKey(variableHeader.packetId());
+ }
+
+ @Override
+ public void onMaxRetransmissionAttemptsReached() {
+ pendingPublishes.computeIfPresent(variableHeader.packetId(), (__, pendingPublish) -> {
+ var message = "Unable to deliver publish message due to max retransmission attempts (%s) being reached for client '%s' on topic '%s' (message ID: %d)"
+ .formatted(clientConfig.getRetransmissionConfig().maxAttempts(), clientConfig.getClientId(), topic, variableHeader.packetId());
+ pendingPublish.getFuture().tryFailure(new MaxRetransmissionsReachedException(message));
+ pendingPublish.getPayload().release();
+ return null;
+ });
+ }
+ }).build();
+
this.pendingPublishes.put(pendingPublish.getMessageId(), pendingPublish);
ChannelFuture channelFuture = this.sendAndFlushPacket(message);
@@ -499,9 +525,30 @@ final class MqttClientImpl implements MqttClient {
MqttSubscribePayload payload = new MqttSubscribePayload(Collections.singletonList(subscription));
MqttSubscribeMessage message = new MqttSubscribeMessage(fixedHeader, variableHeader, payload);
- final MqttPendingSubscription pendingSubscription = new MqttPendingSubscription(future, topic, message,
- () -> !pendingSubscriptions.containsKey(variableHeader.messageId()));
- pendingSubscription.addHandler(handler, once);
+ final var pendingSubscription = MqttPendingSubscription.builder()
+ .future(future)
+ .topic(topic)
+ .handlers(Sets.newHashSet(new MqttPendingSubscription.MqttPendingHandler(handler, once)))
+ .subscribeMessage(message)
+ .ownerId(clientConfig.getOwnerId())
+ .retransmissionConfig(clientConfig.getRetransmissionConfig())
+ .pendingOperation(new PendingOperation() {
+ @Override
+ public boolean isCancelled() {
+ return !pendingSubscriptions.containsKey(variableHeader.messageId());
+ }
+
+ @Override
+ public void onMaxRetransmissionAttemptsReached() {
+ pendingSubscriptions.computeIfPresent(variableHeader.messageId(), (__, pendingSubscription) -> {
+ var message = "Unable to deliver subscribe message due to max retransmission attempts (%s) being reached for client '%s' on topic '%s' (message ID: %d)"
+ .formatted(clientConfig.getRetransmissionConfig().maxAttempts(), clientConfig.getClientId(), topic, variableHeader.messageId());
+ pendingSubscription.getFuture().tryFailure(new MaxRetransmissionsReachedException(message));
+ return null;
+ });
+ }
+ }).build();
+
this.pendingSubscriptions.put(variableHeader.messageId(), pendingSubscription);
this.pendingSubscribeTopics.add(topic);
pendingSubscription.setSent(this.sendAndFlushPacket(message) != null); //If not sent, we will send it when the connection is opened
@@ -518,8 +565,29 @@ final class MqttClientImpl implements MqttClient {
MqttUnsubscribePayload payload = new MqttUnsubscribePayload(Collections.singletonList(topic));
MqttUnsubscribeMessage message = new MqttUnsubscribeMessage(fixedHeader, variableHeader, payload);
- MqttPendingUnsubscription pendingUnsubscription = new MqttPendingUnsubscription(promise, topic, message,
- () -> !pendingServerUnsubscribes.containsKey(variableHeader.messageId()));
+ final var pendingUnsubscription = MqttPendingUnsubscription.builder()
+ .future(promise)
+ .topic(topic)
+ .unsubscribeMessage(message)
+ .ownerId(clientConfig.getOwnerId())
+ .retransmissionConfig(clientConfig.getRetransmissionConfig())
+ .pendingOperation(new PendingOperation() {
+ @Override
+ public boolean isCancelled() {
+ return !pendingServerUnsubscribes.containsKey(variableHeader.messageId());
+ }
+
+ @Override
+ public void onMaxRetransmissionAttemptsReached() {
+ pendingServerUnsubscribes.computeIfPresent(variableHeader.messageId(), (__, pendingUnsubscription) -> {
+ var message = "Unable to deliver unsubscribe message due to max retransmission attempts (%s) being reached for client '%s' on topic '%s' (message ID: %d)"
+ .formatted(clientConfig.getRetransmissionConfig().maxAttempts(), clientConfig.getClientId(), topic, variableHeader.messageId());
+ pendingUnsubscription.getFuture().tryFailure(new MaxRetransmissionsReachedException(message));
+ return null;
+ });
+ }
+ }).build();
+
this.pendingServerUnsubscribes.put(variableHeader.messageId(), pendingUnsubscription);
pendingUnsubscription.startRetransmissionTimer(this.eventLoop.next(), this::sendAndFlushPacket);
diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttConnectResult.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttConnectResult.java
index 911bc1d395..67757d2a7a 100644
--- a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttConnectResult.java
+++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttConnectResult.java
@@ -17,7 +17,9 @@ package org.thingsboard.mqtt;
import io.netty.channel.ChannelFuture;
import io.netty.handler.codec.mqtt.MqttConnectReturnCode;
+import lombok.ToString;
+@ToString
@SuppressWarnings({"WeakerAccess", "unused"})
public final class MqttConnectResult {
@@ -42,4 +44,5 @@ public final class MqttConnectResult {
public ChannelFuture getCloseFuture() {
return closeFuture;
}
+
}
diff --git a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttPendingPublish.java b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttPendingPublish.java
index e8c3ef35f7..1846bdb12b 100644
--- a/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttPendingPublish.java
+++ b/netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttPendingPublish.java
@@ -21,9 +21,13 @@ import io.netty.handler.codec.mqtt.MqttMessage;
import io.netty.handler.codec.mqtt.MqttPublishMessage;
import io.netty.handler.codec.mqtt.MqttQoS;
import io.netty.util.concurrent.Promise;
+import lombok.AccessLevel;
+import lombok.Getter;
+import lombok.Setter;
import java.util.function.Consumer;
+@Getter(AccessLevel.PACKAGE)
final class MqttPendingPublish {
private final int messageId;
@@ -32,80 +36,126 @@ final class MqttPendingPublish {
private final MqttPublishMessage message;
private final MqttQoS qos;
+ @Getter(AccessLevel.NONE)
private final RetransmissionHandler publishRetransmissionHandler;
+ @Getter(AccessLevel.NONE)
private final RetransmissionHandler pubrelRetransmissionHandler;
+ @Setter(AccessLevel.PACKAGE)
private boolean sent = false;
- MqttPendingPublish(int messageId, Promise future, ByteBuf payload, MqttPublishMessage message, MqttQoS qos, PendingOperation operation) {
+ private MqttPendingPublish(
+ int messageId,
+ Promise future,
+ ByteBuf payload,
+ MqttPublishMessage message,
+ MqttQoS qos,
+ String ownerId,
+ MqttClientConfig.RetransmissionConfig retransmissionConfig,
+ PendingOperation pendingOperation
+ ) {
this.messageId = messageId;
this.future = future;
this.payload = payload;
this.message = message;
this.qos = qos;
- this.publishRetransmissionHandler = new RetransmissionHandler<>(operation);
- this.publishRetransmissionHandler.setOriginalMessage(message);
- this.pubrelRetransmissionHandler = new RetransmissionHandler<>(operation);
- }
-
- int getMessageId() {
- return messageId;
- }
-
- Promise getFuture() {
- return future;
- }
-
- ByteBuf getPayload() {
- return payload;
- }
-
- boolean isSent() {
- return sent;
- }
-
- void setSent(boolean sent) {
- this.sent = sent;
- }
-
- MqttPublishMessage getMessage() {
- return message;
- }
-
- MqttQoS getQos() {
- return qos;
+ publishRetransmissionHandler = new RetransmissionHandler<>(retransmissionConfig, pendingOperation, ownerId);
+ publishRetransmissionHandler.setOriginalMessage(message);
+ pubrelRetransmissionHandler = new RetransmissionHandler<>(retransmissionConfig, pendingOperation, ownerId);
}
void startPublishRetransmissionTimer(EventLoop eventLoop, Consumer