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 5a58c0372c..18d0af0119 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 @@ -69,6 +69,8 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInt protected static final String groupId = "SparkplugBGroupId"; protected static final String edgeNode = "SparkpluBNode"; protected static final String keysBdSeq = "bdSeq"; + protected static final String alias = "Failed Post Telemetry node proto payload. SparkplugMessageType "; + protected String deviceId = "Test Sparkplug B Device"; protected int bdSeq = 0; protected int seq = 0; protected static final long PUBLISH_TS_DELTA_MS = 86400000;// Publish start TS <-> 24h @@ -261,4 +263,18 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInt return connectionResult.getResponse(); } + protected long getBdSeqNum() throws Exception { + if (bdSeq == 256) { + bdSeq = 0; + } + return bdSeq++; + } + + protected long getSeqNum() throws Exception { + if (seq == 256) { + seq = 0; + } + return seq++; + } + } 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 09a0a2ef78..7d5cb73abb 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 @@ -16,9 +16,7 @@ package org.thingsboard.server.transport.mqtt.sparkplug.connection; import com.google.common.util.concurrent.ListenableFuture; -import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; -import org.eclipse.paho.client.mqttv3.MqttClient; import org.eclipse.paho.mqttv5.client.IMqttToken; import org.eclipse.paho.mqttv5.client.MqttConnectionOptions; import org.eclipse.paho.mqttv5.common.MqttMessage; @@ -26,33 +24,25 @@ import org.eclipse.paho.mqttv5.common.packet.MqttConnAck; import org.eclipse.paho.mqttv5.common.packet.MqttReturnCode; import org.eclipse.paho.mqttv5.common.packet.MqttWireMessage; import org.junit.Assert; -import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; -import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; -import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSparkplugTest; import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; -import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil; - -import java.io.ByteArrayOutputStream; -import java.io.DataOutputStream; -import java.io.IOException; -import java.io.ObjectOutputStream; -import java.nio.ByteBuffer; -import java.util.Calendar; -import java.util.Date; + import java.util.Optional; -import java.util.concurrent.ExecutorService; +import java.util.Set; +import java.util.HashSet; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64; /** @@ -67,17 +57,15 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr protected void processClientWithCorrectNodeAccessTokenWithNdeathTest() throws Exception { long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; - int value = bdSeq; + long value = bdSeq = 0; MetricDataType metricDataType = Int64; - TsKvEntry expectedTsKvEntryOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, Integer.toUnsignedLong(value))); + TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, value)); SparkplugBProto.Payload.Builder deathPayload = SparkplugBProto.Payload.newBuilder() .setTimestamp(calendar.getTimeInMillis()); - deathPayload.addMetrics(createMetric(expectedTsKvEntryOriginal, metricDataType)); - - byte[] deathBytes = deathPayload.build().toByteArray(); + deathPayload.addMetrics(createMetric(tsKvEntryBdSecOriginal, metricDataType)); - MqttWireMessage response = clientWithCorrectNodeAccessTokenWithNdeath(deathBytes); + MqttWireMessage response = clientWithCorrectNodeAccessTokenWithNDEATH(deathPayload.build().toByteArray()); Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType()); @@ -86,10 +74,9 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr Assert.assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode()); String keys = SparkplugMessageType.NDEATH.name() + " " + keysBdSeq; - TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, Integer.toUnsignedLong(value))); - ListenableFuture> future = tsService.findLatest(tenantId, savedGateway.getId(), keys); - AtomicReference>> finalFuture = new AtomicReference<>(future); - await("Failed Post Telemetry node proto payload. SparkplugMessageType NDEATH") + TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, value)); + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NDEATH.name()) .atMost(40, TimeUnit.SECONDS) .until(() -> { finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), keys)); @@ -97,17 +84,58 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr }); TsKvEntry actualTsKvEntry = finalFuture.get().get().get(); Assert.assertEquals(expectedTsKvEntry, actualTsKvEntry); - client.disconnect(); } protected void processClientWithCorrectAccessTokenCreatedDevices(int cntDevices) throws Exception { processClientWithCorrectNodeAccess(); + long ts = calendar.getTimeInMillis(); + MetricDataType metricDataType = Int32; + Set deviceIds = new HashSet<>(); + String keys = "Device Metric int32"; + int valueDeviceInt32 = 1024; + TsKvEntry expectedTsKvEntryDeviceInt32 = new BasicTsKvEntry(ts, new LongDataEntry(keys, Integer.toUnsignedLong(valueDeviceInt32))); + SparkplugBProto.Payload.Metric metric = createMetric(expectedTsKvEntryDeviceInt32, metricDataType); + for (int i=0; i < cntDevices; i++ ) { + SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()) + .setSeq(getSeqNum()); + String deviceName = deviceId + "_" + i; + + payloadBirthDevice.addMetrics(metric); + if (client.isConnected()) { + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceName, + payloadBirthDevice.build().toByteArray(), 0, false); + deviceIds.add(deviceName); + } + } + + Assert.assertEquals(cntDevices, deviceIds.size()); + + for (String deviceName: deviceIds) { + AtomicReference device = new AtomicReference<>(); + await(alias + SparkplugMessageType.DBIRTH.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); + return device.get() != null; + }); + Assert.assertEquals(deviceName, device.get().getName()); + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NDEATH.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findLatest(tenantId, device.get().getId(), keys)); + return finalFuture.get().get().isPresent(); + }); + TsKvEntry actualTsKvEntry = finalFuture.get().get().get(); + Assert.assertEquals(expectedTsKvEntryDeviceInt32, actualTsKvEntry); + } } - private MqttWireMessage clientWithCorrectNodeAccessTokenWithNdeath(byte[] deathBytes) throws Exception { + private MqttWireMessage clientWithCorrectNodeAccessTokenWithNDEATH(byte[] deathBytes) throws Exception { this.client = new MqttV5TestClient(); MqttConnectionOptions options = new MqttConnectionOptions(); options.setUserName(gatewayAccessToken); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java index 70f3cdf9eb..7fdb25849b 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java @@ -15,40 +15,15 @@ */ package org.thingsboard.server.transport.mqtt.sparkplug.timeseries; -import com.google.common.util.concurrent.ListenableFuture; -import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; -import org.eclipse.paho.mqttv5.client.IMqttToken; -import org.eclipse.paho.mqttv5.client.MqttConnectionOptions; -import org.eclipse.paho.mqttv5.common.MqttMessage; -import org.eclipse.paho.mqttv5.common.packet.MqttConnAck; -import org.eclipse.paho.mqttv5.common.packet.MqttReturnCode; -import org.eclipse.paho.mqttv5.common.packet.MqttWireMessage; -import org.junit.Assert; -import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; -import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; -import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSparkplugTest; import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; -import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil; - -import java.io.ByteArrayOutputStream; -import java.io.DataOutputStream; -import java.io.IOException; -import java.io.ObjectOutputStream; -import java.nio.ByteBuffer; -import java.util.Date; -import java.util.Optional; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicReference; - -import static org.awaitility.Awaitility.await; -import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64; /** @@ -60,6 +35,37 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstra protected void processPushTelemetry() throws Exception { processClientWithCorrectNodeAccess(); + } + + protected void processClientWithCorrectAccessTokenCreatedPublishBirthNode() throws Exception { + processClientWithCorrectNodeAccess(); + SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()); + + long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; + long valueBdSec = getBdSeqNum(); + MetricDataType metricDataType = Int64; + TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, valueBdSec)); + payloadBirthNode.addMetrics(createMetric(tsKvEntryBdSecOriginal, metricDataType)); + + String keys = "Node Control/Rebirth"; + boolean valueRebirth = false; + metricDataType = MetricDataType.Boolean; + TsKvEntry expectedSsKvEntryRebirth = new BasicTsKvEntry(ts, new BooleanDataEntry(keys, valueRebirth)); + payloadBirthNode.addMetrics(createMetric(expectedSsKvEntryRebirth , metricDataType)); + + keys = "Node Metric int32"; + int valueNodeInt32 = 1024; + metricDataType = MetricDataType.Boolean; + TsKvEntry expectedSsKvEntryNodeInt32 = new BasicTsKvEntry(ts, new LongDataEntry(keys, Integer.toUnsignedLong(valueNodeInt32))); + payloadBirthNode.addMetrics(createMetric(expectedSsKvEntryNodeInt32 , metricDataType)); + + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + payloadBirthNode.build().toByteArray(), 0, false); + + + + } } 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 042cbae51e..b9316e1351 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 @@ -20,7 +20,6 @@ import org.junit.After; import org.junit.Before; import org.junit.Test; import org.thingsboard.server.dao.service.DaoSqlTest; -import org.thingsboard.server.transport.mqtt.sparkplug.connection.AbstractMqttV5ClientSparkplugConnectionTest; /** * Created by nickAS21 on 12.01.23 @@ -39,6 +38,11 @@ public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSpa client.disconnect(); } } + @Test + public void testClientWithCorrectAccessTokenCreatedPublishBirthNode() throws Exception { + processClientWithCorrectAccessTokenCreatedPublishBirthNode(); + } + @Test public void testPushTelemetry() throws Exception { processPushTelemetry();