diff --git a/application/pom.xml b/application/pom.xml
index 4d89099d5b..47d69d34dd 100644
--- a/application/pom.xml
+++ b/application/pom.xml
@@ -140,6 +140,10 @@
org.eclipse.paho
org.eclipse.paho.client.mqttv3
+
+ org.eclipse.paho
+ org.eclipse.paho.mqttv5.client
+
org.cassandraunit
cassandra-unit
diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java
index b5d256bfa8..5f4395fb47 100644
--- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java
@@ -51,7 +51,7 @@ import java.util.UUID;
public class AlarmEdgeProcessor extends BaseEdgeProcessor {
public ListenableFuture processAlarmFromEdge(TenantId tenantId, AlarmUpdateMsg alarmUpdateMsg) {
- log.trace("[{}] onAlarmUpdate [{}]", tenantId, alarmUpdateMsg);
+ log.trace("[{}] processAlarmFromEdge [{}]", tenantId, alarmUpdateMsg);
EntityId originatorId = getAlarmOriginator(tenantId, alarmUpdateMsg.getOriginatorName(),
EntityType.valueOf(alarmUpdateMsg.getOriginatorType()));
if (originatorId == null) {
diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java
index 81e7970611..02aca5165c 100644
--- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java
@@ -82,7 +82,7 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
private static final ReentrantLock deviceCreationLock = new ReentrantLock();
public ListenableFuture processDeviceFromEdge(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) {
- log.trace("[{}] onDeviceUpdate [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName());
+ log.trace("[{}] processDeviceFromEdge [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName());
switch (deviceUpdateMsg.getMsgType()) {
case ENTITY_CREATED_RPC_MESSAGE:
String deviceName = deviceUpdateMsg.getName();
@@ -155,7 +155,7 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
}
public ListenableFuture processDeviceCredentialsFromEdge(TenantId tenantId, DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg) {
- log.debug("Executing onDeviceCredentialsUpdate, deviceCredentialsUpdateMsg [{}]", deviceCredentialsUpdateMsg);
+ log.debug("[{}] Executing processDeviceCredentialsFromEdge, deviceCredentialsUpdateMsg [{}]", tenantId, deviceCredentialsUpdateMsg);
DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsUpdateMsg.getDeviceIdMSB(), deviceCredentialsUpdateMsg.getDeviceIdLSB()));
ListenableFuture deviceFuture = deviceService.findDeviceByIdAsync(tenantId, deviceId);
return Futures.transform(deviceFuture, device -> {
@@ -201,9 +201,7 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
device.setCustomerId(getCustomerId(deviceUpdateMsg));
Optional deviceDataOpt =
dataDecodingEncodingService.decode(deviceUpdateMsg.getDeviceDataBytes().toByteArray());
- if (deviceDataOpt.isPresent()) {
- device.setDeviceData(deviceDataOpt.get());
- }
+ deviceDataOpt.ifPresent(device::setDeviceData);
Device savedDevice = deviceService.saveDevice(device);
tbClusterService.onDeviceUpdated(savedDevice, device, false);
return saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null);
@@ -462,7 +460,6 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
}
private DownlinkMsg convertRpcCallEventToDownlink(EdgeEvent edgeEvent) {
- log.trace("Executing convertRpcCallEventToDownlink, edgeEvent [{}]", edgeEvent);
return DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceRpcCallMsg(deviceMsgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody()))
diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java
index 2998001a0a..8b07d771d6 100644
--- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java
@@ -58,7 +58,7 @@ import java.util.UUID;
public class RelationEdgeProcessor extends BaseEdgeProcessor {
public ListenableFuture processRelationFromEdge(TenantId tenantId, RelationUpdateMsg relationUpdateMsg) {
- log.trace("[{}] onRelationUpdate [{}]", tenantId, relationUpdateMsg);
+ log.trace("[{}] processRelationFromEdge [{}]", tenantId, relationUpdateMsg);
try {
EntityRelation entityRelation = new EntityRelation();
diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java
index 4be7ddb1c8..2f72d96770 100644
--- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java
+++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java
@@ -45,6 +45,7 @@ import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.id.WidgetsBundleId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.DataType;
+import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntityRelationsQuery;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
@@ -52,13 +53,10 @@ import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.relation.RelationsSearchParameters;
import org.thingsboard.server.common.data.widget.WidgetType;
import org.thingsboard.server.common.data.widget.WidgetsBundle;
-import org.thingsboard.server.dao.asset.AssetProfileService;
-import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.attributes.AttributesService;
-import org.thingsboard.server.dao.device.DeviceProfileService;
-import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.edge.EdgeEventService;
import org.thingsboard.server.dao.relation.RelationService;
+import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.widget.WidgetTypeService;
import org.thingsboard.server.dao.widget.WidgetsBundleService;
import org.thingsboard.server.gen.edge.v1.AttributesRequestMsg;
@@ -93,24 +91,15 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
private AttributesService attributesService;
@Autowired
- private RelationService relationService;
-
+ private TimeseriesService timeseriesService;
+
@Autowired
- private DeviceService deviceService;
-
- @Autowired
- private AssetService assetService;
+ private RelationService relationService;
@Lazy
@Autowired
private TbEntityViewService entityViewService;
- @Autowired
- private DeviceProfileService deviceProfileService;
-
- @Autowired
- private AssetProfileService assetProfileService;
-
@Autowired
private WidgetsBundleService widgetsBundleService;
@@ -141,77 +130,89 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
EntityId entityId = EntityIdFactory.getByTypeAndUuid(
EntityType.valueOf(attributesRequestMsg.getEntityType()),
new UUID(attributesRequestMsg.getEntityIdMSB(), attributesRequestMsg.getEntityIdLSB()));
- final EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(entityId.getEntityType());
- if (type == null) {
+ final EdgeEventType entityType = EdgeUtils.getEdgeEventTypeByEntityType(entityId.getEntityType());
+ if (entityType == null) {
log.warn("[{}] Type doesn't supported {}", tenantId, entityId.getEntityType());
return Futures.immediateFuture(null);
}
- SettableFuture futureToSet = SettableFuture.create();
String scope = attributesRequestMsg.getScope();
ListenableFuture> findAttrFuture = attributesService.findAll(tenantId, entityId, scope);
- Futures.addCallback(findAttrFuture, new FutureCallback<>() {
- @Override
- public void onSuccess(@Nullable List ssAttributes) {
- if (ssAttributes == null || ssAttributes.isEmpty()) {
- log.trace("[{}][{}] No attributes found for entity {} [{}]", tenantId,
- edge.getName(),
- entityId.getEntityType(),
- entityId.getId());
- futureToSet.set(null);
- return;
- }
-
- try {
- Map entityData = new HashMap<>();
- ObjectNode attributes = JacksonUtil.OBJECT_MAPPER.createObjectNode();
- for (AttributeKvEntry attr : ssAttributes) {
- if (DefaultDeviceStateService.PERSISTENT_ATTRIBUTES.contains(attr.getKey())
- && !DefaultDeviceStateService.INACTIVITY_TIMEOUT.equals(attr.getKey())) {
- continue;
- }
- if (attr.getDataType() == DataType.BOOLEAN && attr.getBooleanValue().isPresent()) {
- attributes.put(attr.getKey(), attr.getBooleanValue().get());
- } else if (attr.getDataType() == DataType.DOUBLE && attr.getDoubleValue().isPresent()) {
- attributes.put(attr.getKey(), attr.getDoubleValue().get());
- } else if (attr.getDataType() == DataType.LONG && attr.getLongValue().isPresent()) {
- attributes.put(attr.getKey(), attr.getLongValue().get());
- } else {
- attributes.put(attr.getKey(), attr.getValueAsString());
- }
- }
- entityData.put("kv", attributes);
- entityData.put("scope", scope);
- JsonNode body = JacksonUtil.OBJECT_MAPPER.valueToTree(entityData);
- log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body);
- ListenableFuture future = saveEdgeEvent(tenantId, edge.getId(), type, EdgeEventActionType.ATTRIBUTES_UPDATED, entityId, body);
- Futures.addCallback(future, new FutureCallback<>() {
- @Override
- public void onSuccess(@Nullable Void unused) {
- futureToSet.set(null);
- }
+ return Futures.transformAsync(findAttrFuture, ssAttributes -> {
+ if (ssAttributes == null || ssAttributes.isEmpty()) {
+ log.trace("[{}][{}] No attributes found for entity {} [{}]", tenantId,
+ edge.getName(),
+ entityId.getEntityType(),
+ entityId.getId());
+ return Futures.immediateFuture(null);
+ }
+ return processEntityAttributesAndAddToEdgeQueue(tenantId, entityId, edge, entityType, scope, ssAttributes, attributesRequestMsg);
+ }, dbCallbackExecutorService);
+ }
- @Override
- public void onFailure(Throwable throwable) {
- String errMsg = String.format("[%s] Failed to save edge event [%s]", edge.getId(), attributesRequestMsg);
- log.error(errMsg, throwable);
- futureToSet.setException(new RuntimeException(errMsg, throwable));
- }
- }, dbCallbackExecutorService);
- } catch (Exception e) {
- String errMsg = String.format("[%s] Failed to save attribute updates to the edge [%s]", edge.getId(), attributesRequestMsg);
- log.error(errMsg, e);
- futureToSet.setException(new RuntimeException(errMsg, e));
+ private ListenableFuture processEntityAttributesAndAddToEdgeQueue(TenantId tenantId, EntityId entityId, Edge edge,
+ EdgeEventType entityType, String scope, List ssAttributes,
+ AttributesRequestMsg attributesRequestMsg) {
+ try {
+ Map entityData = new HashMap<>();
+ ObjectNode attributes = JacksonUtil.OBJECT_MAPPER.createObjectNode();
+ for (AttributeKvEntry attr : ssAttributes) {
+ if (DefaultDeviceStateService.PERSISTENT_ATTRIBUTES.contains(attr.getKey())
+ && !DefaultDeviceStateService.INACTIVITY_TIMEOUT.equals(attr.getKey())) {
+ continue;
+ }
+ if (attr.getDataType() == DataType.BOOLEAN && attr.getBooleanValue().isPresent()) {
+ attributes.put(attr.getKey(), attr.getBooleanValue().get());
+ } else if (attr.getDataType() == DataType.DOUBLE && attr.getDoubleValue().isPresent()) {
+ attributes.put(attr.getKey(), attr.getDoubleValue().get());
+ } else if (attr.getDataType() == DataType.LONG && attr.getLongValue().isPresent()) {
+ attributes.put(attr.getKey(), attr.getLongValue().get());
+ } else {
+ attributes.put(attr.getKey(), attr.getValueAsString());
}
}
-
- @Override
- public void onFailure(Throwable t) {
- String errMsg = String.format("[%s] Can't find attributes [%s]", edge.getId(), attributesRequestMsg);
- log.error(errMsg, t);
- futureToSet.setException(new RuntimeException(errMsg, t));
+ ListenableFuture future;
+ if (attributes.size() > 0) {
+ entityData.put("kv", attributes);
+ entityData.put("scope", scope);
+ JsonNode body = JacksonUtil.OBJECT_MAPPER.valueToTree(entityData);
+ log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body);
+ future = saveEdgeEvent(tenantId, edge.getId(), entityType, EdgeEventActionType.ATTRIBUTES_UPDATED, entityId, body);
+ } else {
+ future = Futures.immediateFuture(null);
}
+ return Futures.transformAsync(future, v -> processLatestTimeseriesAndAddToEdgeQueue(tenantId, entityId, edge, entityType), dbCallbackExecutorService);
+ } catch (Exception e) {
+ String errMsg = String.format("[%s] Failed to save attribute updates to the edge [%s]", edge.getId(), attributesRequestMsg);
+ log.error(errMsg, e);
+ return Futures.immediateFailedFuture(new RuntimeException(errMsg, e));
+ }
+ }
+
+ private ListenableFuture processLatestTimeseriesAndAddToEdgeQueue(TenantId tenantId, EntityId entityId, Edge edge,
+ EdgeEventType entityType) {
+ ListenableFuture> getAllLatestFuture = timeseriesService.findAllLatest(tenantId, entityId);
+ return Futures.transformAsync(getAllLatestFuture, tsKvEntries -> {
+ if (tsKvEntries == null || tsKvEntries.isEmpty()) {
+ log.trace("[{}][{}] No timeseries found for entity {} [{}]", tenantId,
+ edge.getName(),
+ entityId.getEntityType(),
+ entityId.getId());
+ return Futures.immediateFuture(null);
+ }
+ List> futures = new ArrayList<>();
+ for (TsKvEntry tsKvEntry : tsKvEntries) {
+ if (DefaultDeviceStateService.PERSISTENT_ATTRIBUTES.contains(tsKvEntry.getKey())) {
+ continue;
+ }
+ ObjectNode entityBody = JacksonUtil.OBJECT_MAPPER.createObjectNode();
+ ObjectNode ts = JacksonUtil.OBJECT_MAPPER.createObjectNode();
+ ts.put(tsKvEntry.getKey(), tsKvEntry.getValueAsString());
+ entityBody.set("data", ts);
+ entityBody.put("ts", tsKvEntry.getTs());
+ futures.add(saveEdgeEvent(tenantId, edge.getId(), entityType, EdgeEventActionType.TIMESERIES_UPDATED, entityId, JacksonUtil.valueToTree(entityBody)));
+ }
+ return Futures.transform(Futures.allAsList(futures), v -> null, dbCallbackExecutorService);
}, dbCallbackExecutorService);
- return futureToSet;
}
@Override
diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
index ee34bd7f23..9532bc3044 100644
--- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
@@ -561,7 +561,7 @@ public class DefaultDataUpdateService implements DataUpdateService {
while (hasNext) {
for (Alarm alarm : alarms.getData()) {
if (alarm.getCustomerId() == null && alarm.getOriginator() != null) {
- alarm.setCustomerId(entityService.fetchEntityCustomerId(tenantId, alarm.getOriginator()));
+ alarm.setCustomerId(entityService.fetchEntityCustomerId(tenantId, alarm.getOriginator()).get());
alarmDao.save(tenantId, alarm);
}
if (processed.incrementAndGet() % 1000 == 0) {
diff --git a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
index 02406da08a..5314d4872f 100644
--- a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
+++ b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
@@ -18,7 +18,6 @@ package org.thingsboard.server.edge;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
-import com.google.protobuf.AbstractMessage;
import com.google.protobuf.InvalidProtocolBufferException;
import com.google.protobuf.MessageLite;
import org.junit.After;
@@ -38,9 +37,6 @@ import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.alarm.AlarmSeverity;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.AssetProfile;
-import org.thingsboard.server.common.data.device.data.DefaultDeviceConfiguration;
-import org.thingsboard.server.common.data.device.data.DeviceData;
-import org.thingsboard.server.common.data.device.data.MqttDeviceTransportConfiguration;
import org.thingsboard.server.common.data.device.profile.AlarmCondition;
import org.thingsboard.server.common.data.device.profile.AlarmConditionFilter;
import org.thingsboard.server.common.data.device.profile.AlarmConditionFilterKey;
@@ -447,11 +443,6 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
}
protected Device saveDeviceOnCloudAndVerifyDeliveryToEdge() throws Exception {
- // create ota package
- edgeImitator.expectMessageAmount(1);
- OtaPackageInfo firmwareOtaPackageInfo = saveOtaPackageInfo(thermostatDeviceProfile.getId());
- Assert.assertTrue(edgeImitator.waitForMessages());
-
// create device and assign to edge
Device savedDevice = saveDevice(StringUtils.randomAlphanumeric(15), thermostatDeviceProfile.getName());
edgeImitator.expectMessageAmount(2); // device and device profile messages
@@ -471,38 +462,6 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType());
Assert.assertEquals(thermostatDeviceProfile.getUuidId().getMostSignificantBits(), deviceProfileUpdateMsg.getIdMSB());
Assert.assertEquals(thermostatDeviceProfile.getUuidId().getLeastSignificantBits(), deviceProfileUpdateMsg.getIdLSB());
-
- // update device
- edgeImitator.expectMessageAmount(1);
- savedDevice.setFirmwareId(firmwareOtaPackageInfo.getId());
-
- DeviceData deviceData = new DeviceData();
- deviceData.setConfiguration(new DefaultDeviceConfiguration());
- MqttDeviceTransportConfiguration transportConfiguration = new MqttDeviceTransportConfiguration();
- transportConfiguration.getProperties().put("topic", "tb_rule_engine.thermostat");
- deviceData.setTransportConfiguration(transportConfiguration);
- savedDevice.setDeviceData(deviceData);
-
- savedDevice = doPost("/api/device", savedDevice, Device.class);
- Assert.assertTrue(edgeImitator.waitForMessages());
- AbstractMessage latestMessage = edgeImitator.getLatestMessage();
- Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg);
- deviceUpdateMsg = (DeviceUpdateMsg) latestMessage;
- Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType());
- Assert.assertEquals(savedDevice.getUuidId().getMostSignificantBits(), deviceUpdateMsg.getIdMSB());
- Assert.assertEquals(savedDevice.getUuidId().getLeastSignificantBits(), deviceUpdateMsg.getIdLSB());
- Assert.assertEquals(savedDevice.getName(), deviceUpdateMsg.getName());
- Assert.assertEquals(savedDevice.getType(), deviceUpdateMsg.getType());
- Assert.assertEquals(firmwareOtaPackageInfo.getUuidId().getMostSignificantBits(), deviceUpdateMsg.getFirmwareIdMSB());
- Assert.assertEquals(firmwareOtaPackageInfo.getUuidId().getLeastSignificantBits(), deviceUpdateMsg.getFirmwareIdLSB());
- Optional deviceDataOpt =
- dataDecodingEncodingService.decode(deviceUpdateMsg.getDeviceDataBytes().toByteArray());
- Assert.assertTrue(deviceDataOpt.isPresent());
- deviceData = deviceDataOpt.get();
- Assert.assertTrue(deviceData.getTransportConfiguration() instanceof MqttDeviceTransportConfiguration);
- MqttDeviceTransportConfiguration mqttDeviceTransportConfiguration =
- (MqttDeviceTransportConfiguration) deviceData.getTransportConfiguration();
- Assert.assertEquals("tb_rule_engine.thermostat", mqttDeviceTransportConfiguration.getProperties().get("topic"));
return savedDevice;
}
diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java
index e5fe165a1e..4af8fede14 100644
--- a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java
+++ b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java
@@ -31,8 +31,12 @@ import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
+import org.thingsboard.server.common.data.OtaPackageInfo;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TenantProfile;
+import org.thingsboard.server.common.data.device.data.DefaultDeviceConfiguration;
+import org.thingsboard.server.common.data.device.data.DeviceData;
+import org.thingsboard.server.common.data.device.data.MqttDeviceTransportConfiguration;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
@@ -57,8 +61,8 @@ import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.gen.edge.v1.UplinkMsg;
import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import java.util.List;
import java.util.Map;
@@ -170,6 +174,8 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
// create device and assign to edge; update device
Device savedDevice = saveDeviceOnCloudAndVerifyDeliveryToEdge();
+ verifyUpdateFirmwareIdAndDeviceData(savedDevice);
+
// update device credentials - ACCESS_TOKEN
edgeImitator.expectMessageAmount(1);
DeviceCredentials deviceCredentials =
@@ -204,6 +210,45 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
Assert.assertEquals(deviceCredentials.getCredentialsValue(), deviceCredentialsUpdateMsg.getCredentialsValue());
}
+ private void verifyUpdateFirmwareIdAndDeviceData(Device savedDevice) throws InterruptedException {
+ // create ota package
+ edgeImitator.expectMessageAmount(1);
+ OtaPackageInfo firmwareOtaPackageInfo = saveOtaPackageInfo(thermostatDeviceProfile.getId());
+ Assert.assertTrue(edgeImitator.waitForMessages());
+
+ // update device
+ edgeImitator.expectMessageAmount(1);
+ savedDevice.setFirmwareId(firmwareOtaPackageInfo.getId());
+
+ DeviceData deviceData = new DeviceData();
+ deviceData.setConfiguration(new DefaultDeviceConfiguration());
+ MqttDeviceTransportConfiguration transportConfiguration = new MqttDeviceTransportConfiguration();
+ transportConfiguration.getProperties().put("topic", "tb_rule_engine.thermostat");
+ deviceData.setTransportConfiguration(transportConfiguration);
+ savedDevice.setDeviceData(deviceData);
+
+ savedDevice = doPost("/api/device", savedDevice, Device.class);
+ Assert.assertTrue(edgeImitator.waitForMessages());
+ AbstractMessage latestMessage = edgeImitator.getLatestMessage();
+ Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg);
+ DeviceUpdateMsg deviceUpdateMsg = (DeviceUpdateMsg) latestMessage;
+ Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType());
+ Assert.assertEquals(savedDevice.getUuidId().getMostSignificantBits(), deviceUpdateMsg.getIdMSB());
+ Assert.assertEquals(savedDevice.getUuidId().getLeastSignificantBits(), deviceUpdateMsg.getIdLSB());
+ Assert.assertEquals(savedDevice.getName(), deviceUpdateMsg.getName());
+ Assert.assertEquals(savedDevice.getType(), deviceUpdateMsg.getType());
+ Assert.assertEquals(firmwareOtaPackageInfo.getUuidId().getMostSignificantBits(), deviceUpdateMsg.getFirmwareIdMSB());
+ Assert.assertEquals(firmwareOtaPackageInfo.getUuidId().getLeastSignificantBits(), deviceUpdateMsg.getFirmwareIdLSB());
+ Optional deviceDataOpt =
+ dataDecodingEncodingService.decode(deviceUpdateMsg.getDeviceDataBytes().toByteArray());
+ Assert.assertTrue(deviceDataOpt.isPresent());
+ deviceData = deviceDataOpt.get();
+ Assert.assertTrue(deviceData.getTransportConfiguration() instanceof MqttDeviceTransportConfiguration);
+ MqttDeviceTransportConfiguration mqttDeviceTransportConfiguration =
+ (MqttDeviceTransportConfiguration) deviceData.getTransportConfiguration();
+ Assert.assertEquals("tb_rule_engine.thermostat", mqttDeviceTransportConfiguration.getProperties().get("topic"));
+ }
+
@Test
public void testDeviceReachedMaximumAllowedOnCloud() throws Exception {
// update tenant profile configuration
@@ -323,6 +368,9 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
"inactivityTimeout", "3600000");
sendAttributesRequestAndVerify(device, DataConstants.SHARED_SCOPE, "{\"key2\":\"value2\"}",
"key2", "value2");
+
+ doDelete("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/" + DataConstants.SERVER_SCOPE, "keys","key1, inactivityTimeout");
+ doDelete("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/" + DataConstants.SHARED_SCOPE, "keys", "key2");
}
@Test
@@ -640,4 +688,53 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
client.disconnect();
}
+
+ @Test
+ public void testVerifyDeliveryOfLatestTimeseriesOnAttributesRequest() throws Exception {
+ Device device = findDeviceByName("Edge Device 1");
+
+ JsonNode timeseriesData = mapper.readTree("{\"temperature\":25}");
+
+ doPost("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE,
+ timeseriesData);
+
+ // Wait before device timeseries saved to database before requesting them from edge
+ Awaitility.await()
+ .atMost(10, TimeUnit.SECONDS)
+ .until(() -> {
+ String urlTemplate = "/api/plugins/telemetry/DEVICE/" + device.getId() + "/keys/timeseries";
+ List actualKeys = doGetAsyncTyped(urlTemplate, new TypeReference<>() {});
+ return actualKeys != null && !actualKeys.isEmpty() && actualKeys.contains("temperature");
+ });
+
+ UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
+ AttributesRequestMsg.Builder attributesRequestMsgBuilder = AttributesRequestMsg.newBuilder();
+ attributesRequestMsgBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits());
+ attributesRequestMsgBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits());
+ attributesRequestMsgBuilder.setEntityType(EntityType.DEVICE.name());
+ attributesRequestMsgBuilder.setScope(DataConstants.SERVER_SCOPE);
+ uplinkMsgBuilder.addAttributesRequestMsg(attributesRequestMsgBuilder.build());
+
+ edgeImitator.expectResponsesAmount(1);
+ edgeImitator.expectMessageAmount(1);
+ edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
+ Assert.assertTrue(edgeImitator.waitForResponses());
+ Assert.assertTrue(edgeImitator.waitForMessages());
+
+ AbstractMessage latestMessage = edgeImitator.getLatestMessage();
+ Assert.assertTrue(latestMessage instanceof EntityDataProto);
+ EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage;
+ Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB());
+ Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB());
+ Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
+ Assert.assertTrue(latestEntityDataMsg.hasPostTelemetryMsg());
+
+ TransportProtos.PostTelemetryMsg timeseriesUpdatedMsg = latestEntityDataMsg.getPostTelemetryMsg();
+ Assert.assertEquals(1, timeseriesUpdatedMsg.getTsKvListList().size());
+ TransportProtos.TsKvListProto tsKvListProto = timeseriesUpdatedMsg.getTsKvListList().get(0);
+ Assert.assertEquals(1, tsKvListProto.getKvList().size());
+ TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKvList().get(0);
+ Assert.assertEquals(25, keyValueProto.getLongV());
+ Assert.assertEquals("temperature", keyValueProto.getKey());
+ }
}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestCallback.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestCallback.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestCallback.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestCallback.java
index 18815d1c38..8c2bc5d488 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestCallback.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestCallback.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt;
+package org.thingsboard.server.transport.mqtt.mqttv3;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestClient.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestClient.java
similarity index 97%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestClient.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestClient.java
index 96846f8678..a0aa1905eb 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestClient.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestClient.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt;
+package org.thingsboard.server.transport.mqtt.mqttv3;
import io.netty.handler.codec.mqtt.MqttQoS;
import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
@@ -104,6 +104,10 @@ public class MqttTestClient {
return client.subscribe(topic, qoS.value());
}
+ public boolean isConnected() {
+ return client.isConnected();
+ }
+
public void enableManualAcks() {
client.setManualAcks(true);
}
@@ -112,6 +116,10 @@ public class MqttTestClient {
client.messageArrivedComplete(mqttMessage.getId(), mqttMessage.getQos());
}
+ private MqttAsyncClient createClient() throws MqttException {
+ return createClient(null);
+ }
+
private MqttAsyncClient createClient(String clientId) throws MqttException {
if (StringUtils.isEmpty(clientId)) {
clientId = MqttAsyncClient.generateClientId();
@@ -119,8 +127,4 @@ public class MqttTestClient {
return new MqttAsyncClient(MQTT_URL, clientId, new MemoryPersistence());
}
- private MqttAsyncClient createClient() throws MqttException {
- return createClient(null);
- }
-
}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/AbstractMqttAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java
similarity index 99%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/AbstractMqttAttributesIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java
index 56d9ef5e9e..8811767326 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/AbstractMqttAttributesIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes;
import com.github.os72.protobuf.dynamic.DynamicSchema;
import com.google.protobuf.Descriptors;
@@ -22,7 +22,6 @@ import com.google.protobuf.InvalidProtocolBufferException;
import com.squareup.wire.schema.internal.parser.ProtoFileElement;
import io.netty.handler.codec.mqtt.MqttQoS;
import lombok.extern.slf4j.Slf4j;
-import org.springframework.test.context.TestPropertySource;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DynamicProtoUtils;
@@ -41,8 +40,8 @@ import org.thingsboard.server.gen.transport.TransportApiProtos;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import java.util.ArrayList;
import java.util.List;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java
similarity index 97%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java
index 2231cb89eb..e9337b5446 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.request;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.request;
import lombok.extern.slf4j.Slf4j;
import org.junit.Test;
@@ -22,7 +22,7 @@ import org.thingsboard.server.common.data.device.profile.MqttTopics;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
import java.util.ArrayList;
import java.util.List;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestIntegrationTest.java
similarity index 93%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestIntegrationTest.java
index c4da6b3184..a9dc90f3c7 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.request;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.request;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
@@ -21,7 +21,7 @@ import org.junit.Test;
import org.thingsboard.server.common.data.device.profile.MqttTopics;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
@Slf4j
@DaoSqlTest
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestJsonIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestJsonIntegrationTest.java
similarity index 93%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestJsonIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestJsonIntegrationTest.java
index a27c9391fd..d8e2f4428e 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestJsonIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestJsonIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.request;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.request;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
@@ -22,7 +22,7 @@ import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.common.data.device.profile.MqttTopics;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
@Slf4j
@DaoSqlTest
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestProtoIntegrationTest.java
similarity index 96%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestProtoIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestProtoIntegrationTest.java
index d330cdbbaf..ef8533e582 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestProtoIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestProtoIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.request;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.request;
import lombok.extern.slf4j.Slf4j;
import org.junit.Test;
@@ -22,7 +22,7 @@ import org.thingsboard.server.common.data.device.profile.MqttTopics;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
import java.util.ArrayList;
import java.util.List;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java
similarity index 96%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java
index 05253cd969..3332b962e2 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java
@@ -13,14 +13,14 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.updates;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.updates;
import lombok.extern.slf4j.Slf4j;
import org.junit.Test;
import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_PROTO_TOPIC;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesIntegrationTest.java
similarity index 93%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesIntegrationTest.java
index 01a33c5b5c..830c40aa91 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.updates;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.updates;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
@@ -21,7 +21,7 @@ import org.junit.Test;
import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_TOPIC;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesJsonIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesJsonIntegrationTest.java
similarity index 93%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesJsonIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesJsonIntegrationTest.java
index f35ce1e084..834a01fb03 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesJsonIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesJsonIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.updates;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.updates;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
@@ -21,7 +21,7 @@ import org.junit.Test;
import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_TOPIC;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesProtoIntegrationTest.java
similarity index 94%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesProtoIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesProtoIntegrationTest.java
index dc7e92a128..91c302f322 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesProtoIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesProtoIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.attributes.updates;
+package org.thingsboard.server.transport.mqtt.mqttv3.attributes.updates;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
@@ -21,7 +21,7 @@ import org.junit.Test;
import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
-import org.thingsboard.server.transport.mqtt.attributes.AbstractMqttAttributesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.attributes.AbstractMqttAttributesIntegrationTest;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_SHORT_PROTO_TOPIC;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimBackwardCompatibilityDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimBackwardCompatibilityDeviceTest.java
similarity index 97%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimBackwardCompatibilityDeviceTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimBackwardCompatibilityDeviceTest.java
index e7670b5170..00a77c2e2b 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimBackwardCompatibilityDeviceTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimBackwardCompatibilityDeviceTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.claim;
+package org.thingsboard.server.transport.mqtt.mqttv3.claim;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimDeviceTest.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimDeviceTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimDeviceTest.java
index df3ac12e45..c46b78b77d 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimDeviceTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimDeviceTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.claim;
+package org.thingsboard.server.transport.mqtt.mqttv3.claim;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
@@ -28,7 +28,7 @@ import org.thingsboard.server.dao.device.claim.ClaimResult;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.transport.TransportApiProtos;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import static org.junit.Assert.assertEquals;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimJsonDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimJsonDeviceTest.java
similarity index 97%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimJsonDeviceTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimJsonDeviceTest.java
index 05d4974f9d..3c90eda7ed 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimJsonDeviceTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimJsonDeviceTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.claim;
+package org.thingsboard.server.transport.mqtt.mqttv3.claim;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimProtoDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimProtoDeviceTest.java
similarity index 96%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimProtoDeviceTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimProtoDeviceTest.java
index 89c8c5f231..0bfc057806 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/claim/MqttClaimProtoDeviceTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/claim/MqttClaimProtoDeviceTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.claim;
+package org.thingsboard.server.transport.mqtt.mqttv3.claim;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
@@ -21,7 +21,7 @@ import org.junit.Test;
import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.transport.TransportApiProtos;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
@Slf4j
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/client/AbstractMqttClientConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/client/AbstractMqttClientConnectionTest.java
new file mode 100644
index 0000000000..634458c07d
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/client/AbstractMqttClientConnectionTest.java
@@ -0,0 +1,50 @@
+/**
+ * Copyright © 2016-2022 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.transport.mqtt.mqttv3.client;
+
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.junit.Assert;
+import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
+
+public abstract class AbstractMqttClientConnectionTest extends AbstractMqttIntegrationTest {
+
+ protected void processClientWithCorrectAccessTokenTest() throws Exception {
+ MqttTestClient client = new MqttTestClient();
+ client.connectAndWait(accessToken);
+ Assert.assertTrue(client.isConnected());
+ client.disconnect();
+ }
+
+ protected void processClientWithWrongAccessTokenTest() throws Exception {
+ MqttTestClient client = new MqttTestClient();
+ try {
+ client.connectAndWait("wrongAccessToken");
+ } catch (MqttException e) {
+ Assert.assertEquals(MqttException.REASON_CODE_FAILED_AUTHENTICATION, e.getReasonCode());
+ }
+ }
+
+ protected void processClientWithWrongClientIdAndEmptyUsernamePasswordTest() throws Exception {
+ MqttTestClient client = new MqttTestClient("unknownClientId");
+ try {
+ client.connectAndWait();
+ } catch (MqttException e) {
+ Assert.assertEquals(MqttException.REASON_CODE_INVALID_CLIENT_ID, e.getReasonCode());
+ }
+ }
+
+}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/client/MqttClientConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/client/MqttClientConnectionTest.java
new file mode 100644
index 0000000000..5e6e6e0188
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/client/MqttClientConnectionTest.java
@@ -0,0 +1,48 @@
+/**
+ * Copyright © 2016-2022 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.transport.mqtt.mqttv3.client;
+
+import org.junit.Before;
+import org.junit.Test;
+import org.thingsboard.server.dao.service.DaoSqlTest;
+import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
+
+@DaoSqlTest
+public class MqttClientConnectionTest extends AbstractMqttClientConnectionTest {
+
+ @Before
+ public void beforeTest() throws Exception {
+ MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder()
+ .deviceName("Test MqttV5 client device")
+ .build();
+ processBeforeTest(configProperties);
+ }
+
+ @Test
+ public void testClientWithCorrectAccessToken() throws Exception {
+ processClientWithCorrectAccessTokenTest();
+ }
+
+ @Test
+ public void testClientWithWrongAccessToken() throws Exception {
+ processClientWithWrongAccessTokenTest();
+ }
+
+ @Test
+ public void testClientWithWrongClientIdAndEmptyUsernamePassword() throws Exception {
+ processClientWithWrongClientIdAndEmptyUsernamePasswordTest();
+ }
+}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/credentials/BasicMqttCredentialsTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/credentials/BasicMqttCredentialsTest.java
similarity index 93%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/credentials/BasicMqttCredentialsTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/credentials/BasicMqttCredentialsTest.java
index 0884f1457c..af463e320e 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/credentials/BasicMqttCredentialsTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/credentials/BasicMqttCredentialsTest.java
@@ -13,10 +13,11 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.credentials;
+package org.thingsboard.server.transport.mqtt.mqttv3.credentials;
import com.fasterxml.jackson.core.type.TypeReference;
-import org.eclipse.paho.client.mqttv3.MqttSecurityException;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.common.util.JacksonUtil;
@@ -27,7 +28,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import java.util.Arrays;
import java.util.HashSet;
@@ -114,11 +115,16 @@ public class BasicMqttCredentialsTest extends AbstractMqttIntegrationTest {
testTelemetryIsDelivered(accessToken2Device, mqttTestClient5);
}
- @Test(expected = MqttSecurityException.class)
+ @Test(expected = MqttException.class)
public void testCorrectClientIdAndUserNameButWrongPassword() throws Exception {
// Not correct. Correct clientId and username, but wrong password
MqttTestClient mqttTestClient = new MqttTestClient(CLIENT_ID);
- mqttTestClient.connectAndWait(USER_NAME3, "WRONG PASSWORD");
+ try {
+ mqttTestClient.connectAndWait(USER_NAME3, "WRONG PASSWORD");
+ Assert.fail(); // This should not happens, because we have a wrong password
+ } catch (MqttException e) {
+ Assert.assertEquals(4, e.getReasonCode()); // 4 - Reason code for bad username or password in MQTT v3
+ }
testTelemetryIsNotDelivered(clientIdAndUserNameAndPasswordDevice3, mqttTestClient);
}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/provision/MqttProvisionJsonDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionJsonDeviceTest.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/provision/MqttProvisionJsonDeviceTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionJsonDeviceTest.java
index 0a06f99209..deb2685bf7 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/provision/MqttProvisionJsonDeviceTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionJsonDeviceTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.provision;
+package org.thingsboard.server.transport.mqtt.mqttv3.provision;
import com.fasterxml.jackson.databind.JsonNode;
import io.netty.handler.codec.mqtt.MqttQoS;
@@ -33,8 +33,8 @@ import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.device.provision.ProvisionResponseStatus;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import java.util.concurrent.TimeUnit;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/provision/MqttProvisionProtoDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionProtoDeviceTest.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/provision/MqttProvisionProtoDeviceTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionProtoDeviceTest.java
index 9e61b69ebf..ce994a3057 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/provision/MqttProvisionProtoDeviceTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionProtoDeviceTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.provision;
+package org.thingsboard.server.transport.mqtt.mqttv3.provision;
import io.netty.handler.codec.mqtt.MqttQoS;
import lombok.extern.slf4j.Slf4j;
@@ -41,8 +41,8 @@ import org.thingsboard.server.gen.transport.TransportProtos.ValidateBasicMqttCre
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceTokenRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509CertRequestMsg;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import java.util.concurrent.TimeUnit;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java
similarity index 99%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java
index 2cdfb76645..483c2bb350 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.rpc;
+package org.thingsboard.server.transport.mqtt.mqttv3.rpc;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
@@ -39,8 +39,8 @@ import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadCo
import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration;
import org.thingsboard.server.gen.transport.TransportApiProtos;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import java.util.ArrayList;
import java.util.List;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcBackwardCompatibilityIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcBackwardCompatibilityIntegrationTest.java
similarity index 99%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcBackwardCompatibilityIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcBackwardCompatibilityIntegrationTest.java
index 5fb123666f..d8a896a83e 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcBackwardCompatibilityIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcBackwardCompatibilityIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.rpc;
+package org.thingsboard.server.transport.mqtt.mqttv3.rpc;
import lombok.extern.slf4j.Slf4j;
import org.junit.Test;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcDefaultIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcDefaultIntegrationTest.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcDefaultIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcDefaultIntegrationTest.java
index 95e3b7fde3..042e2808b3 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcDefaultIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcDefaultIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.rpc;
+package org.thingsboard.server.transport.mqtt.mqttv3.rpc;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import lombok.extern.slf4j.Slf4j;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcJsonIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcJsonIntegrationTest.java
similarity index 96%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcJsonIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcJsonIntegrationTest.java
index b3c3a994d1..508ce02fcd 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcJsonIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcJsonIntegrationTest.java
@@ -13,14 +13,14 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.rpc;
+package org.thingsboard.server.transport.mqtt.mqttv3.rpc;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.dao.service.DaoSqlTest;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_RPC_REQUESTS_SUB_SHORT_JSON_TOPIC;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcProtoIntegrationTest.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcProtoIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcProtoIntegrationTest.java
index e888f49ff4..0286c20f3f 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/MqttServerSideRpcProtoIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/MqttServerSideRpcProtoIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.rpc;
+package org.thingsboard.server.transport.mqtt.mqttv3.rpc;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesIntegrationTest.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesIntegrationTest.java
index 1d986d0453..04cbbe9a30 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.attributes;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.attributes;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
@@ -24,7 +24,7 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import java.util.Arrays;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesJsonIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesJsonIntegrationTest.java
similarity index 96%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesJsonIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesJsonIntegrationTest.java
index f77885c0fd..17e0edb14c 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesJsonIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesJsonIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.attributes;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.attributes;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesProtoIntegrationTest.java
similarity index 99%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesProtoIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesProtoIntegrationTest.java
index 7b2fa1143d..da25f39def 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/attributes/MqttAttributesProtoIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/attributes/MqttAttributesProtoIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.attributes;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.attributes;
import com.github.os72.protobuf.dynamic.DynamicSchema;
import com.google.protobuf.Descriptors;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
similarity index 98%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
index 5bb0d28a21..0a58804e92 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries;
import com.fasterxml.jackson.core.type.TypeReference;
import io.netty.handler.codec.mqtt.MqttQoS;
@@ -23,8 +23,8 @@ import org.junit.Test;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import java.util.Arrays;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesJsonIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesJsonIntegrationTest.java
similarity index 97%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesJsonIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesJsonIntegrationTest.java
index 4a7037f6e8..dc94f53f38 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesJsonIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesJsonIntegrationTest.java
@@ -13,14 +13,14 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.server.common.data.TransportPayloadType;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import java.util.Arrays;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java
similarity index 99%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java
index 9c92e3f4f0..0b67f0e48a 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesProtoIntegrationTest.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries;
import com.github.os72.protobuf.dynamic.DynamicSchema;
import com.google.protobuf.Descriptors;
@@ -31,8 +31,8 @@ import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadCo
import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration;
import org.thingsboard.server.gen.transport.TransportApiProtos;
import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.transport.mqtt.MqttTestCallback;
-import org.thingsboard.server.transport.mqtt.MqttTestClient;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
import java.util.Arrays;
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlIntegrationTest.java
similarity index 80%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlIntegrationTest.java
index 7006638eb1..1d9baa5b08 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlIntegrationTest.java
@@ -13,10 +13,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries.nosql;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.nosql;
import org.thingsboard.server.dao.service.DaoNoSqlTest;
-import org.thingsboard.server.transport.mqtt.telemetry.timeseries.AbstractMqttTimeseriesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.AbstractMqttTimeseriesIntegrationTest;
@DaoNoSqlTest
public class MqttTimeseriesNoSqlIntegrationTest extends AbstractMqttTimeseriesIntegrationTest {
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlJsonIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlJsonIntegrationTest.java
similarity index 80%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlJsonIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlJsonIntegrationTest.java
index 8e6e2cf679..999e5717b4 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlJsonIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlJsonIntegrationTest.java
@@ -13,10 +13,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries.nosql;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.nosql;
import org.thingsboard.server.dao.service.DaoNoSqlTest;
-import org.thingsboard.server.transport.mqtt.telemetry.timeseries.AbstractMqttTimeseriesJsonIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.AbstractMqttTimeseriesJsonIntegrationTest;
@DaoNoSqlTest
public class MqttTimeseriesNoSqlJsonIntegrationTest extends AbstractMqttTimeseriesJsonIntegrationTest {
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlProtoIntegrationTest.java
similarity index 80%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlProtoIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlProtoIntegrationTest.java
index e91d5982d7..06f2bff3b5 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/nosql/MqttTimeseriesNoSqlProtoIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/nosql/MqttTimeseriesNoSqlProtoIntegrationTest.java
@@ -13,10 +13,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries.nosql;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.nosql;
import org.thingsboard.server.dao.service.DaoNoSqlTest;
-import org.thingsboard.server.transport.mqtt.telemetry.timeseries.AbstractMqttTimeseriesProtoIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.AbstractMqttTimeseriesProtoIntegrationTest;
@DaoNoSqlTest
public class MqttTimeseriesNoSqlProtoIntegrationTest extends AbstractMqttTimeseriesProtoIntegrationTest {
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlIntegrationTest.java
similarity index 80%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlIntegrationTest.java
index de9a596d43..e317debabc 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlIntegrationTest.java
@@ -13,10 +13,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries.sql;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.sql;
import org.thingsboard.server.dao.service.DaoSqlTest;
-import org.thingsboard.server.transport.mqtt.telemetry.timeseries.AbstractMqttTimeseriesIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.AbstractMqttTimeseriesIntegrationTest;
@DaoSqlTest
public class MqttTimeseriesSqlIntegrationTest extends AbstractMqttTimeseriesIntegrationTest {
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlJsonIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlJsonIntegrationTest.java
similarity index 80%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlJsonIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlJsonIntegrationTest.java
index 323dd751f1..bf991c0c06 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlJsonIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlJsonIntegrationTest.java
@@ -13,10 +13,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries.sql;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.sql;
import org.thingsboard.server.dao.service.DaoSqlTest;
-import org.thingsboard.server.transport.mqtt.telemetry.timeseries.AbstractMqttTimeseriesJsonIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.AbstractMqttTimeseriesJsonIntegrationTest;
@DaoSqlTest
public class MqttTimeseriesSqlJsonIntegrationTest extends AbstractMqttTimeseriesJsonIntegrationTest {
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlProtoIntegrationTest.java
similarity index 80%
rename from application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlProtoIntegrationTest.java
rename to application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlProtoIntegrationTest.java
index ebfc1d49cf..7e17d30844 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/telemetry/timeseries/sql/MqttTimeseriesSqlProtoIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/sql/MqttTimeseriesSqlProtoIntegrationTest.java
@@ -13,10 +13,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.transport.mqtt.telemetry.timeseries.sql;
+package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.sql;
import org.thingsboard.server.dao.service.DaoSqlTest;
-import org.thingsboard.server.transport.mqtt.telemetry.timeseries.AbstractMqttTimeseriesProtoIntegrationTest;
+import org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries.AbstractMqttTimeseriesProtoIntegrationTest;
@DaoSqlTest
public class MqttTimeseriesSqlProtoIntegrationTest extends AbstractMqttTimeseriesProtoIntegrationTest {
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/AbstractMqttV5Test.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/AbstractMqttV5Test.java
new file mode 100644
index 0000000000..4714184f0e
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/AbstractMqttV5Test.java
@@ -0,0 +1,21 @@
+/**
+ * Copyright © 2016-2022 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.transport.mqtt.mqttv5;
+
+import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
+
+public abstract class AbstractMqttV5Test extends AbstractMqttIntegrationTest {
+}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestCallback.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestCallback.java
new file mode 100644
index 0000000000..7979936659
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestCallback.java
@@ -0,0 +1,110 @@
+/**
+ * Copyright © 2016-2022 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.transport.mqtt.mqttv5;
+
+import lombok.Data;
+import lombok.extern.slf4j.Slf4j;
+import org.eclipse.paho.mqttv5.client.IMqttToken;
+import org.eclipse.paho.mqttv5.client.MqttCallback;
+import org.eclipse.paho.mqttv5.client.MqttDisconnectResponse;
+import org.eclipse.paho.mqttv5.common.MqttException;
+import org.eclipse.paho.mqttv5.common.MqttMessage;
+import org.eclipse.paho.mqttv5.common.packet.MqttProperties;
+import org.eclipse.paho.mqttv5.common.packet.MqttWireMessage;
+
+import java.util.concurrent.CountDownLatch;
+
+@Slf4j
+@Data
+public class MqttV5TestCallback implements MqttCallback {
+
+ protected CountDownLatch subscribeLatch;
+ protected final CountDownLatch deliveryLatch;
+ protected int qoS;
+ protected byte[] payloadBytes;
+ protected String awaitSubTopic;
+ protected boolean pubAckReceived;
+ protected MqttMessage lastReceivedMessage;
+
+ public MqttV5TestCallback() {
+ this.subscribeLatch = new CountDownLatch(1);
+ this.deliveryLatch = new CountDownLatch(1);
+ }
+
+ public MqttV5TestCallback(int subscribeCount) {
+ this.subscribeLatch = new CountDownLatch(subscribeCount);
+ this.deliveryLatch = new CountDownLatch(1);
+ }
+
+ public MqttV5TestCallback(String awaitSubTopic) {
+ this.subscribeLatch = new CountDownLatch(1);
+ this.deliveryLatch = new CountDownLatch(1);
+ this.awaitSubTopic = awaitSubTopic;
+ }
+
+ @Override
+ public void disconnected(MqttDisconnectResponse mqttDisconnectResponse) {
+ if (mqttDisconnectResponse.getException() != null) {
+ log.warn("connectionLost: ", mqttDisconnectResponse.getException());
+ deliveryLatch.countDown();
+ }
+ log.warn("Disconnected with reason: {}", mqttDisconnectResponse.getReasonString());
+ }
+
+ @Override
+ public void mqttErrorOccurred(MqttException e) {
+ log.warn("Error occurred:", e);
+ }
+
+ @Override
+ public void messageArrived(String requestTopic, MqttMessage mqttMessage) {
+ lastReceivedMessage = mqttMessage;
+ if (awaitSubTopic == null) {
+ log.warn("messageArrived on topic: {}", requestTopic);
+ qoS = mqttMessage.getQos();
+ payloadBytes = mqttMessage.getPayload();
+ subscribeLatch.countDown();
+ } else {
+ messageArrivedOnAwaitSubTopic(requestTopic, mqttMessage);
+ }
+ }
+
+ protected void messageArrivedOnAwaitSubTopic(String requestTopic, MqttMessage mqttMessage) {
+ log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic);
+ if (awaitSubTopic.equals(requestTopic)) {
+ qoS = mqttMessage.getQos();
+ payloadBytes = mqttMessage.getPayload();
+ subscribeLatch.countDown();
+ }
+ }
+
+ @Override
+ public void deliveryComplete(IMqttToken iMqttToken) {
+ log.warn("delivery complete: {}", iMqttToken.getResponse());
+ pubAckReceived = iMqttToken.getResponse().getType() == MqttWireMessage.MESSAGE_TYPE_PUBACK;
+ deliveryLatch.countDown();
+ }
+
+ @Override
+ public void connectComplete(boolean reconnect, String serverURI) {
+ log.warn("Connect completed: reconnect - {}, serverURI - {}", reconnect, serverURI);
+ }
+
+ @Override
+ public void authPacketArrived(int reasonCode, MqttProperties mqttProperties) {
+ log.warn("Auth package received: reasonCode - {}, mqtt properties - {}", reasonCode, mqttProperties);
+ }
+}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java
new file mode 100644
index 0000000000..0671dd6f95
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java
@@ -0,0 +1,175 @@
+/**
+ * Copyright © 2016-2022 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.transport.mqtt.mqttv5;
+
+import io.netty.handler.codec.mqtt.MqttQoS;
+import org.eclipse.paho.mqttv5.client.IMqttToken;
+import org.eclipse.paho.mqttv5.client.MqttAsyncClient;
+import org.eclipse.paho.mqttv5.client.MqttCallback;
+import org.eclipse.paho.mqttv5.client.MqttConnectionOptions;
+import org.eclipse.paho.mqttv5.client.persist.MemoryPersistence;
+import org.eclipse.paho.mqttv5.common.MqttException;
+import org.eclipse.paho.mqttv5.common.MqttMessage;
+import org.thingsboard.server.common.data.StringUtils;
+
+import java.util.concurrent.TimeUnit;
+
+public class MqttV5TestClient { // We should copy part of MqttV3TestClient, due to different package names in import
+
+ private static final String MQTT_URL = "tcp://localhost:1883";
+ private static final int TIMEOUT = 30; // seconds
+ private static final long TIMEOUT_MS = TimeUnit.SECONDS.toMillis(TIMEOUT);
+
+ private final MqttAsyncClient client;
+
+ public void setCallback(MqttCallback callback) {
+ client.setCallback(callback);
+ }
+
+ public MqttV5TestClient() throws MqttException {
+ this.client = createClient();
+ }
+
+ public MqttV5TestClient(String clientId) throws MqttException {
+ this.client = createClient(clientId);
+ }
+
+ public MqttV5TestClient(boolean generateClientId) throws MqttException {
+ this.client = createClient(generateClientId);
+ }
+
+ public IMqttToken connectAndWait(String userName, String password) throws MqttException {
+ IMqttToken connect = connect(userName, password);
+ connect.waitForCompletion(TIMEOUT_MS);
+ return connect;
+ }
+
+ public IMqttToken connectAndWait(String userName) throws MqttException {
+ return connectAndWait(userName, null);
+ }
+
+ public IMqttToken connectAndWait() throws MqttException {
+ return connectAndWait(null, null);
+ }
+
+ public IMqttToken connectAndWait(MqttConnectionOptions options) throws MqttException {
+ IMqttToken iMqttToken = connect(options);
+ iMqttToken.waitForCompletion(TIMEOUT_MS);
+ return iMqttToken;
+ }
+
+ private IMqttToken connect(String userName, String password) throws MqttException {
+ if (client == null) {
+ throw new RuntimeException("Failed to connect! MqttAsyncClient is not initialized!");
+ }
+ MqttConnectionOptions options = new MqttConnectionOptions();
+ if (StringUtils.isNotEmpty(userName)) {
+ options.setUserName(userName);
+ }
+ if (StringUtils.isNotEmpty(password)) {
+ options.setPassword(password.getBytes());
+ }
+ return client.connect(options);
+ }
+
+ public IMqttToken connect(MqttConnectionOptions options) throws MqttException {
+ if (client == null) {
+ throw new RuntimeException("Failed to connect! MqttAsyncClient is not initialized!");
+ }
+ return client.connect(options);
+ }
+
+ public void disconnectAndWait() throws MqttException {
+ disconnect().waitForCompletion(TIMEOUT_MS);
+ }
+
+ public IMqttToken disconnect() throws MqttException {
+ return client.disconnect();
+ }
+
+ public void disconnectForcibly() throws MqttException {
+ client.disconnectForcibly(TIMEOUT_MS);
+ }
+
+ public IMqttToken publishAndWait(String topic, byte[] payload) throws MqttException {
+ IMqttToken iMqttToken = publish(topic, payload);
+ iMqttToken.waitForCompletion(TIMEOUT_MS);
+ return iMqttToken;
+ }
+
+ public IMqttToken publish(String topic, byte[] payload) throws MqttException {
+ MqttMessage message = new MqttMessage();
+ message.setPayload(payload);
+ return publish(topic, message);
+ }
+
+ public IMqttToken publish(String topic, MqttMessage message) throws MqttException {
+ return publish(topic, message.getPayload(), message.getQos(), message.isRetained());
+ }
+
+ public IMqttToken publish(String topic, byte[] payload, int qos, boolean retain) throws MqttException {
+ return client.publish(topic, payload, qos, retain);
+ }
+
+ public IMqttToken subscribeAndWait(String topic, MqttQoS qoS) throws MqttException {
+ IMqttToken iMqttToken = subscribe(topic, qoS);
+ iMqttToken.waitForCompletion(TIMEOUT_MS);
+ return iMqttToken;
+ }
+
+ public IMqttToken subscribe(String topic, MqttQoS qoS) throws MqttException {
+ return client.subscribe(topic, qoS.value());
+ }
+
+ public IMqttToken unsubscribeAndWait(String topic) throws MqttException {
+ IMqttToken iMqttToken = unsubscribe(topic);
+ iMqttToken.waitForCompletion(TIMEOUT_MS);
+ return iMqttToken;
+ }
+
+ public IMqttToken unsubscribe(String topic) throws MqttException {
+ return client.unsubscribe(topic);
+ }
+
+ public boolean isConnected() {
+ return client.isConnected();
+ }
+
+ public void enableManualAcks() {
+ client.setManualAcks(true);
+ }
+
+ public void messageArrivedComplete(MqttMessage mqttMessage) throws MqttException {
+ client.messageArrivedComplete(mqttMessage.getId(), mqttMessage.getQos());
+ }
+
+ private MqttAsyncClient createClient() throws MqttException {
+ return createClient(true);
+ }
+
+ private MqttAsyncClient createClient(boolean generateClientId) throws MqttException {
+ String clientId = null;
+ if (generateClientId) {
+ clientId = "test" + System.nanoTime();
+ }
+ return createClient(clientId);
+ }
+
+ private MqttAsyncClient createClient(String clientId) throws MqttException {
+ return new MqttAsyncClient(MQTT_URL, clientId, new MemoryPersistence());
+ }
+
+}
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/attributes/AbstractAttributesMqttV5Test.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/attributes/AbstractAttributesMqttV5Test.java
new file mode 100644
index 0000000000..204c66e597
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/attributes/AbstractAttributesMqttV5Test.java
@@ -0,0 +1,161 @@
+/**
+ * Copyright © 2016-2022 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.transport.mqtt.mqttv5.attributes;
+
+import com.fasterxml.jackson.core.type.TypeReference;
+import io.netty.handler.codec.mqtt.MqttQoS;
+import org.junit.Before;
+import org.thingsboard.common.util.JacksonUtil;
+import org.thingsboard.server.common.data.device.profile.MqttTopics;
+import org.thingsboard.server.common.data.id.DeviceId;
+import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties;
+import org.thingsboard.server.transport.mqtt.mqttv5.AbstractMqttV5Test;
+import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestCallback;
+import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient;
+
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertTrue;
+import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
+
+public abstract class AbstractAttributesMqttV5Test extends AbstractMqttV5Test {
+
+ private static final String SHARED_ATTRIBUTES_PAYLOAD = "{\"sharedStr\":\"value1\",\"sharedBool\":true,\"sharedDbl\":42.0,\"sharedLong\":73," +
+ "\"sharedJson\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}";
+
+ private static final String SHARED_ATTRIBUTES_DELETED_RESPONSE = "{\"deleted\":[\"sharedJson\"]}";
+
+ protected static final String PAYLOAD_VALUES_STR = "{\"key1\":\"value1\", \"key2\":true, \"key3\": 3.0, \"key4\": 4," +
+ " \"key5\": {\"someNumber\": 42, \"someArray\": [1,2,3], \"someNestedObject\": {\"key\": \"value\"}}}";
+
+ @Before
+ public void beforeTest() throws Exception {
+ MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder()
+ .deviceName("Test Post Attributes device")
+ .build();
+ processBeforeTest(configProperties);
+ }
+
+ protected void processAttributesPublishTest() throws Exception {
+ List expectedKeys = Arrays.asList("key1", "key2", "key3", "key4", "key5");
+
+ MqttV5TestClient client = new MqttV5TestClient();
+ client.connectAndWait(accessToken);
+
+ client.publishAndWait(MqttTopics.DEVICE_ATTRIBUTES_TOPIC, PAYLOAD_VALUES_STR.getBytes());
+ client.disconnectAndWait();
+
+ DeviceId deviceId = savedDevice.getId();
+
+ long start = System.currentTimeMillis();
+ long end = System.currentTimeMillis() + 5000;
+
+ List actualKeys = null;
+ while (start <= end) {
+ actualKeys = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + deviceId + "/keys/attributes/CLIENT_SCOPE", new TypeReference<>() {
+ });
+ if (actualKeys.size() == expectedKeys.size()) {
+ break;
+ }
+ Thread.sleep(100);
+ start += 100;
+ }
+ assertNotNull(actualKeys);
+
+ Set actualKeySet = new HashSet<>(actualKeys);
+
+ Set expectedKeySet = new HashSet<>(expectedKeys);
+
+ assertEquals(expectedKeySet, actualKeySet);
+
+ String getAttributesValuesUrl = getAttributesValuesUrl(deviceId, actualKeySet);
+ List
+
+ org.eclipse.paho
+ org.eclipse.paho.mqttv5.client
+ ${paho.mqttv5.client.version}
+
org.apache.curator
curator-x-discovery
diff --git a/pull_request_template.md b/pull_request_template.md
index 21362429dc..4fd83ec4f0 100644
--- a/pull_request_template.md
+++ b/pull_request_template.md
@@ -12,7 +12,7 @@ Put your PR description here instead of this sentence.
- [ ] Description contains brief notes about what needs to be added to the documentation.
- [ ] No merge conflicts, commented blocks of code, code formatting issues.
- [ ] Changes are backward compatible or upgrade script is provided.
-- [ ] Similar PR is opened for PE version to simplify merge. Required for internal contributors only.
+- [ ] Similar PR is opened for PE version to simplify merge. Crosslinks between PRs added. Required for internal contributors only.
## Front-End feature checklist