Browse Source

sparkplug: Test Telemetry

pull/7931/head
nickAS21 4 years ago
parent
commit
fa4c00c437
  1. 24
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java
  2. 7
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java
  3. 181
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java
  4. 15
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java
  5. 1
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java

24
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();

7
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> 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<ListenableFuture<Optional<TsKvEntry>>> 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 {

181
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<String> 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<ListenableFuture<List<TsKvEntry>>> 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<String> 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<ListenableFuture<List<TsKvEntry>>> 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<String> 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<ListenableFuture<List<TsKvEntry>>> 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();
}
}

15
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();
}
}

1
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:

Loading…
Cancel
Save