diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java index bbf8a291dc..7d8770bbcf 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java @@ -67,6 +67,9 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataTyp import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt32; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt64; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt8; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConnectionState.ONLINE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; @@ -82,6 +85,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte protected ThreadLocalRandom random = ThreadLocalRandom.current(); protected static final String groupId = "SparkplugBGroupId"; + protected static final String edgeNodeDeviceName = "Test Connect Sparkplug client node"; protected static final String edgeNode = "SparkpluBNode"; protected static final String keysBdSeq = "bdSeq"; protected static final String alias = "Failed Telemetry/Attribute proto sparkplug payload. SparkplugMessageType "; @@ -98,14 +102,39 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte protected static final String metricBirthName_Int32 = "Device Metric int32"; protected Set sparkplugAttributesMetricNames; - public void beforeSparkplugTest() throws Exception { - MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() - .gatewayName("Test Connect Sparkplug client node") - .isSparkplug(true) - .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) - .transportPayloadType(TransportPayloadType.PROTOBUF) - .build(); - processBeforeTest(configProperties); + public void beforeSparkplugTest(boolean isCreateDevices) throws Exception { + if (isCreateDevices) { + MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() + .gatewayName(edgeNodeDeviceName) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + configProperties = MqttTestConfigProperties.builder() + .gatewayName(deviceId) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + configProperties = MqttTestConfigProperties.builder() + .gatewayName(groupId + ":" + edgeNode + ":" + deviceId) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + + } else { + MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() + .gatewayName(edgeNodeDeviceName) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + } } public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception { @@ -200,6 +229,97 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte return devices; } + protected void connectClientWithCorrectAccessTokenWithNDEATHDevicesCreatingBefore_Test(int cntDevices) throws Exception { + long ts = calendar.getTimeInMillis(); + List devices = new ArrayList<>(); + clientWithCorrectNodeAccessTokenWithNDEATH(); + MetricDataType metricDataType = Int32; + String key = "Node Metric int32"; + int valueDeviceInt32 = 1024; + SparkplugBProto.Payload.Metric metric = createMetric(valueDeviceInt32, ts, key, metricDataType, -1L); + SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getBdSeqNum()); + payloadBirthNode.addMetrics(metric); + payloadBirthNode.setTimestamp(ts); + if (client.isConnected()) { + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + payloadBirthNode.build().toByteArray(), 0, false); + } + + valueDeviceInt32 = 4024; + metric = createMetric(valueDeviceInt32, ts, metricBirthName_Int32, metricBirthDataType_Int32, -1L); + // as old device name -> deviceId + String deviceName = deviceId; + AtomicReference device1 = new AtomicReference<>(); + String finalDeviceName1 = deviceName; + await(alias + "find device [" + deviceId + "] before connecting") + .atMost(200, TimeUnit.SECONDS) + .until(() -> { + device1.set(doGet("/api/tenant/devices?deviceName=" + finalDeviceName1, Device.class)); + return device1.get() != null; + }); + + if (client.isConnected()) { + SparkplugBProto.Payload.Builder payloadBirthDevice1 = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getSeqNum()); + payloadBirthDevice1.addMetrics(metric); + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceId, + payloadBirthDevice1.build().toByteArray(), 0, false); + devices.add(device1.get()); + } + // as new device name -> groupId + ":" + edgeNode + ":" + deviceId; + deviceName = groupId + ":" + edgeNode + ":" + deviceId; + AtomicReference device2 = new AtomicReference<>(); + String finalDeviceName2 = deviceName; + await(alias + "find device [" + deviceName + "] before connecting") + .atMost(200, TimeUnit.SECONDS) + .until(() -> { + device2.set(doGet("/api/tenant/devices?deviceName=" + finalDeviceName2, Device.class)); + return device2.get() != null; + }); + + if (client.isConnected()) { + SparkplugBProto.Payload.Builder payloadBirthDevice2 = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getSeqNum()); + payloadBirthDevice2.addMetrics(metric); + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceId, + payloadBirthDevice2.build().toByteArray(), 0, false); + devices.add(device2.get()); + } + Assert.assertEquals(cntDevices, devices.size()); + state_ONLINE_ALL (devices, calendar.getTimeInMillis()); + } + + protected void state_ONLINE_ALL (List devices, long ts) { + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(messageName(STATE), ONLINE.name())); + await(alias + messageName(STATE) + ", device: " + savedGateway.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + var foundEntry = tsService.findAllLatest(tenantId, savedGateway.getId()).get().stream() + .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) + .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) + .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) + .findFirst(); + return foundEntry.isPresent(); + }); + + for (Device device : devices) { + await(alias + messageName(STATE) + ", device: " + device.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + var foundEntry = tsService.findAllLatest(tenantId, device.getId()).get().stream() + .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) + .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) + .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) + .findFirst(); + return foundEntry.isPresent(); + }); + } + } + protected List connectClientWithCorrectAccessTokenWithNDEATHWithAliasCreatedDevices(long ts) throws Exception { List devices = new ArrayList<>(); Long alias = 0L; diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java index 37abade6e7..8cd713ad7d 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java @@ -33,7 +33,7 @@ public class MqttV5ClientSparkplugBAttributesInProfileTest extends AbstractMqttV public void beforeTest() throws Exception { sparkplugAttributesMetricNames = new HashSet<>(); sparkplugAttributesMetricNames.add(metricBirthName_Int32); - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java index 3826b59bdc..b080cfbff0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java @@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBAttributesTest extends AbstractMqttV5ClientSp @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java index 8f9943b393..89b17f6b48 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java @@ -95,31 +95,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstra protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_ALL(int cntDevices) throws Exception { long ts = calendar.getTimeInMillis(); List devices = connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(cntDevices, ts); - - TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(messageName(STATE), ONLINE.name())); - await(alias + messageName(STATE) + ", device: " + savedGateway.getName()) - .atMost(40, TimeUnit.SECONDS) - .until(() -> { - var foundEntry = tsService.findAllLatest(tenantId, savedGateway.getId()).get().stream() - .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) - .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) - .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) - .findFirst(); - return foundEntry.isPresent(); - }); - - for (Device device : devices) { - await(alias + messageName(STATE) + ", device: " + device.getName()) - .atMost(40, TimeUnit.SECONDS) - .until(() -> { - var foundEntry = tsService.findAllLatest(tenantId, device.getId()).get().stream() - .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) - .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) - .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) - .findFirst(); - return foundEntry.isPresent(); - }); - } + state_ONLINE_ALL (devices, ts); } protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OneDeviceOFFLINE(int cntDevices, int indexDeviceDisconnect) throws Exception { diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java new file mode 100644 index 0000000000..797b1beea5 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java @@ -0,0 +1,45 @@ +/** + * Copyright © 2016-2026 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.sparkplug.connection; + +import org.eclipse.paho.mqttv5.common.MqttException; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.thingsboard.server.dao.service.DaoSqlTest; + +/** + * Created by nickAS21 on 12.01.23 + */ +@DaoSqlTest +public class MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest extends AbstractMqttV5ClientSparkplugConnectionTest { + + @Before + public void beforeTest() throws Exception { + beforeSparkplugTest(true); + } + + @After + public void afterTest() throws MqttException { + if (client.isConnected()) { + client.disconnect(); + } + } + + @Test + public void testClientWithCorrectAccessTokenWithNDEATHTwoDevicesCreatingBeforeFirstNameDeviceIdSecondNameFull() throws Exception { + connectClientWithCorrectAccessTokenWithNDEATHDevicesCreatingBefore_Test(2); + }} diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java index 01038c21f1..0d5e73e2f0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java @@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBConnectionTest extends AbstractMqttV5ClientSp @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java index e36b4e6260..08ade0274d 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java @@ -28,7 +28,7 @@ public class MqttV5RpcSparkplugTest extends AbstractMqttV5RpcSparkplugTest { @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After @@ -47,6 +47,7 @@ public class MqttV5RpcSparkplugTest extends AbstractMqttV5RpcSparkplugTest { public void testClientDeviceWithCorrectAccessTokenPublish_TwoWayRpc_Success() throws Exception { processClientDeviceWithCorrectAccessTokenPublish_TwoWayRpc_Success(); } + @Test public void testClientDeviceWithCorrectAccessTokenPublishWithAlias_TwoWayRpc_Success() throws Exception { processClientDeviceWithCorrectAccessTokenPublishWithAlias_TwoWayRpc_Success(); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java index bafe2d81d2..5d9f453415 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java @@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSpa @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After