From fa4c00c4374a4174d186ba92ac164efd8c08e257 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Wed, 18 Jan 2023 23:44:36 +0200 Subject: [PATCH] sparkplug: Test Telemetry --- .../AbstractMqttV5ClientSparkplugTest.java | 24 +-- ...ctMqttV5ClientSparkplugConnectionTest.java | 7 +- ...actMqttV5ClientSparkplugTelemetryTest.java | 181 +++++++++++++++++- .../MqttV5ClientSparkplugBTelemetryTest.java | 15 +- .../util/sparkplug/SparkplugMetricUtil.java | 1 + 5 files changed, 196 insertions(+), 32 deletions(-) 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 18d0af0119..13f364ff72 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 @@ -15,13 +15,9 @@ */ package org.thingsboard.server.transport.mqtt.sparkplug; -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; import org.eclipse.paho.mqttv5.common.packet.MqttConnAck; import org.eclipse.paho.mqttv5.common.packet.MqttReturnCode; import org.eclipse.paho.mqttv5.common.packet.MqttWireMessage; @@ -29,15 +25,12 @@ import org.junit.Assert; import org.thingsboard.server.common.data.TransportPayloadType; 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.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.MqttTestConfigProperties; import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; 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; @@ -47,20 +40,14 @@ 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.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; /** * Created by nickAS21 on 12.01.23 */ @Slf4j -public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttIntegrationTest { +public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttIntegrationTest { protected MqttV5TestClient client; protected Calendar calendar = Calendar.getInstance(); @@ -85,7 +72,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInt processBeforeTest(configProperties); } - public void processClientWithCorrectNodeAccess () throws Exception { + public void processClientWithCorrectNodeAccess() throws Exception { this.client = new MqttV5TestClient(); MqttWireMessage response = clientWithCorrectNodeAccessToken(client); Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType()); @@ -106,8 +93,8 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInt case Int32: case UInt8: case UInt16: - case UInt32: return metric.toBuilder().setIntValue(Integer.parseInt(String.valueOf(value))).build(); + case UInt32: case Int64: case UInt64: case DateTime: @@ -121,9 +108,9 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInt case String: case Text: case UUID: - return metric.toBuilder().setDatasetValue((SparkplugBProto.Payload.DataSet) value).build(); - case DataSet: return metric.toBuilder().setStringValue(String.valueOf(value)).build(); + case DataSet: + return metric.toBuilder().setDatasetValue((SparkplugBProto.Payload.DataSet) value).build(); case Bytes: case Int8Array: ByteString byteString = ByteString.copyFrom((byte[]) value); @@ -257,7 +244,6 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInt } - private MqttWireMessage clientWithCorrectNodeAccessToken(MqttV5TestClient client) throws Exception { IMqttToken connectionResult = client.connectAndWait(gatewayAccessToken); return connectionResult.getResponse(); 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 7d5cb73abb..594acc3aec 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 @@ -113,7 +113,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr for (String deviceName: deviceIds) { AtomicReference device = new AtomicReference<>(); - await(alias + SparkplugMessageType.DBIRTH.name()) + await(alias + "find device [" + deviceName + "] after crete") .atMost(40, TimeUnit.SECONDS) .until(() -> { device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); @@ -121,7 +121,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr }); Assert.assertEquals(deviceName, device.get().getName()); AtomicReference>> finalFuture = new AtomicReference<>(); - await(alias + SparkplugMessageType.NDEATH.name()) + await(alias + SparkplugMessageType.DBIRTH.name()) .atMost(40, TimeUnit.SECONDS) .until(() -> { finalFuture.set(tsService.findLatest(tenantId, device.get().getId(), keys)); @@ -130,9 +130,6 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr TsKvEntry actualTsKvEntry = finalFuture.get().get().get(); Assert.assertEquals(expectedTsKvEntryDeviceInt32, actualTsKvEntry); } - - - } private MqttWireMessage clientWithCorrectNodeAccessTokenWithNDEATH(byte[] deathBytes) throws Exception { 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 7fdb25849b..8bc8593f05 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,16 +15,35 @@ */ package org.thingsboard.server.transport.mqtt.sparkplug.timeseries; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; +import org.junit.Assert; 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.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; 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 java.math.BigInteger; +import java.util.ArrayList; +import java.util.List; +import java.util.Random; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import static org.awaitility.Awaitility.await; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int8; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt16; +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; /** * Created by nickAS21 on 12.01.23 @@ -32,40 +51,196 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataTyp @Slf4j public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends AbstractMqttV5ClientSparkplugTest { - protected void processPushTelemetry() throws Exception { + protected void processClientWithCorrectAccessTokenPushDeviceMetricBuildSimple() throws Exception { processClientWithCorrectNodeAccess(); + Random random = new Random(); + String deviceName = deviceId + "_" + 10; + String messageTypeName = SparkplugMessageType.DDATA.name(); + List listKeys = new ArrayList<>(); + + SparkplugBProto.Payload.Builder ddataPayload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()) + .setSeq(getSeqNum()); + long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; + + String keys = "MyInt8"; + MetricDataType metricDataType = Int8; + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, (long)((byte)random.nextInt()))); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyInt16"; + metricDataType = Int16; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, (long)((short)random.nextInt()))); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyInt32"; + metricDataType = Int32; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, (long)(random.nextInt()))); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyInt64"; + metricDataType = UInt64; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, random.nextLong())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyUInt8"; + metricDataType = UInt8; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, (long)((short)random.nextInt()))); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyUInt16"; + metricDataType = UInt16; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, (long)random.nextInt())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyUInt32"; + metricDataType = UInt32; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, random.nextLong())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyUInt64"; + metricDataType = UInt64; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, BigInteger.valueOf(random.nextLong()).longValue())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyFloat"; + metricDataType = MetricDataType.Float; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, (long)random.nextFloat())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyDateTime"; + metricDataType = MetricDataType.DateTime; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, ts)); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyDouble"; + metricDataType = MetricDataType.Double; + tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, (long) random.nextDouble())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyBoolean"; + metricDataType = MetricDataType.Boolean; + tsKvEntry = new BasicTsKvEntry(ts, new BooleanDataEntry(keys, random.nextBoolean())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyString"; + metricDataType = MetricDataType.String; + tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(keys, newUUID())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyText"; + metricDataType = MetricDataType.Text; + tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(keys, newUUID())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + + keys = "MyUUID"; + metricDataType = MetricDataType.Text; + tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(keys, newUUID())); + ddataPayload.addMetrics(createMetric(tsKvEntry, metricDataType)); + listKeys.add(keys); + if (client.isConnected()) { + client.publish(NAMESPACE + "/" + groupId + "/" + messageTypeName + "/" + edgeNode + "/" + deviceName, + ddataPayload.build().toByteArray(), 0, false); + } + + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NCMD.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), listKeys)); + return !finalFuture.get().get().isEmpty(); + }); + Assert.assertEquals(listKeys.size(), finalFuture.get().get().size()); } - protected void processClientWithCorrectAccessTokenCreatedPublishBirthNode() throws Exception { + protected void processClientWithCorrectAccessTokenPublishNBIRTH() throws Exception { processClientWithCorrectNodeAccess(); + List listKeys = new ArrayList<>(); 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)); + listKeys.add(SparkplugMessageType.NBIRTH.name() + " " + keysBdSeq); 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)); + listKeys.add(keys); 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)); + listKeys.add(keys); client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NBIRTH.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), listKeys)); + return !finalFuture.get().get().isEmpty(); + }); + Assert.assertEquals(listKeys.size(), finalFuture.get().get().size()); + } + protected void processClientWithCorrectAccessTokenPublishNCMDReBirth() throws Exception { + processClientWithCorrectNodeAccess(); + SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()); + List listKeys = new ArrayList<>(); + 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)); + listKeys.add(SparkplugMessageType.NCMD.name() + " " + keysBdSeq); + String keys = "Node Control/Rebirth"; + boolean valueRebirth = true; + metricDataType = MetricDataType.Boolean; + TsKvEntry expectedSsKvEntryRebirth = new BasicTsKvEntry(ts, new BooleanDataEntry(keys, valueRebirth)); + payloadBirthNode.addMetrics(createMetric(expectedSsKvEntryRebirth , metricDataType)); + listKeys.add(keys); + + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.NCMD.name() + "/" + edgeNode, + payloadBirthNode.build().toByteArray(), 0, false); + + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NCMD.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), listKeys)); + return !finalFuture.get().get().isEmpty(); + }); + Assert.assertEquals(listKeys.size(), finalFuture.get().get().size()); + } + private String newUUID() { + return java.util.UUID.randomUUID().toString(); } } 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 b9316e1351..678ce69066 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 @@ -39,13 +39,18 @@ public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSpa } @Test - public void testClientWithCorrectAccessTokenCreatedPublishBirthNode() throws Exception { - processClientWithCorrectAccessTokenCreatedPublishBirthNode(); + public void testClientWithCorrectAccessTokenPushDeviceMetricBuildSimple() throws Exception { + processClientWithCorrectAccessTokenPushDeviceMetricBuildSimple(); } @Test - public void testPushTelemetry() throws Exception { - processPushTelemetry(); + public void testClientWithCorrectAccessTokenPublishNBIRTH() throws Exception { + processClientWithCorrectAccessTokenPublishNBIRTH(); } -} + @Test + public void testClientWithCorrectAccessTokenPublishNCMDReBirth() throws Exception { + processClientWithCorrectAccessTokenPublishNCMDReBirth(); + } + +} \ No newline at end of file diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java index 8019c35774..6b575ba306 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java @@ -69,6 +69,7 @@ public class SparkplugMetricUtil { return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.DOUBLE_V) .setDoubleV(protoMetric.getDoubleValue()).build()); case Int8: + case UInt8: case Int16: case Int32: case UInt16: