Browse Source

sparkplug: add converter kvProto to value Metric array

pull/8032/head
nickAS21 4 years ago
parent
commit
36e12e63af
  1. 17
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java
  2. 13
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java
  3. 128
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java
  4. 4
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java
  5. 85
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  6. 38
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java
  7. 282
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java

17
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java

@ -19,7 +19,10 @@ import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.mqttv5.client.IMqttToken; import org.eclipse.paho.mqttv5.client.IMqttToken;
import org.eclipse.paho.mqttv5.client.MqttConnectionOptions; import org.eclipse.paho.mqttv5.client.MqttConnectionOptions;
import org.eclipse.paho.mqttv5.common.MqttMessage; 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.eclipse.paho.mqttv5.common.packet.MqttWireMessage;
import org.junit.Assert;
import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
@ -30,6 +33,7 @@ import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType
import java.util.Calendar; import java.util.Calendar;
import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric;
@ -52,7 +56,6 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
protected int seq = 0; protected int seq = 0;
protected static final long PUBLISH_TS_DELTA_MS = 86400000;// Publish start TS <-> 24h protected static final long PUBLISH_TS_DELTA_MS = 86400000;// Publish start TS <-> 24h
public void beforeSparkplugTest() throws Exception { public void beforeSparkplugTest() throws Exception {
MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder()
.gatewayName("Test Connect Sparkplug client node") .gatewayName("Test Connect Sparkplug client node")
@ -62,13 +65,13 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
processBeforeTest(configProperties); processBeforeTest(configProperties);
} }
public MqttWireMessage clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception { public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception {
long ts = calendar.getTimeInMillis(); long ts = calendar.getTimeInMillis();
long value = bdSeq = 0; long value = bdSeq = 0;
return clientWithCorrectNodeAccessTokenWithNDEATH(ts, value); clientWithCorrectNodeAccessTokenWithNDEATH(ts, value);
} }
public MqttWireMessage clientWithCorrectNodeAccessTokenWithNDEATH(long ts, long value) throws Exception { public void clientWithCorrectNodeAccessTokenWithNDEATH(long ts, long value) throws Exception {
String key = keysBdSeq; String key = keysBdSeq;
MetricDataType metricDataType = Int64; MetricDataType metricDataType = Int64;
SparkplugBProto.Payload.Builder deathPayload = SparkplugBProto.Payload.newBuilder() SparkplugBProto.Payload.Builder deathPayload = SparkplugBProto.Payload.newBuilder()
@ -84,7 +87,11 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
msg.setPayload(deathBytes); msg.setPayload(deathBytes);
options.setWill(topic, msg); options.setWill(topic, msg);
IMqttToken connectionResult = client.connect(options); IMqttToken connectionResult = client.connect(options);
return connectionResult.getResponse();
MqttWireMessage response = connectionResult.getResponse();
Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType());
MqttConnAck connAckMsg = (MqttConnAck) response;
Assert.assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode());
} }
protected long getBdSeqNum() throws Exception { protected long getBdSeqNum() throws Exception {

13
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java

@ -18,9 +18,6 @@ package org.thingsboard.server.transport.mqtt.sparkplug.connection;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.mqttv5.common.MqttException; import org.eclipse.paho.mqttv5.common.MqttException;
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.junit.Assert;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
@ -39,7 +36,6 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
import static org.awaitility.Awaitility.await; 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.Int32;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric;
@ -52,13 +48,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr
protected void processClientWithCorrectNodeAccessTokenWithNDEATH_Test() throws Exception { protected void processClientWithCorrectNodeAccessTokenWithNDEATH_Test() throws Exception {
long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS;
long value = bdSeq = 0; long value = bdSeq = 0;
MqttWireMessage response = clientWithCorrectNodeAccessTokenWithNDEATH(ts, value); clientWithCorrectNodeAccessTokenWithNDEATH(ts, value);
Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType());
MqttConnAck connAckMsg = (MqttConnAck) response;
Assert.assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode());
String keys = SparkplugMessageType.NDEATH.name() + " " + keysBdSeq; String keys = SparkplugMessageType.NDEATH.name() + " " + keysBdSeq;
TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, value)); TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, value));
@ -113,5 +103,4 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr
Assert.assertEquals(cntDevices, deviceIds.size()); Assert.assertEquals(cntDevices, deviceIds.size());
} }
} }

128
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java

@ -32,7 +32,6 @@ import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSpark
import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; 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.SparkplugMessageType;
import java.math.BigDecimal;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
@ -143,10 +142,11 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
} }
protected void processClientWithCorrectAccessTokenPushNodeMetricBuildPrimitiveSimple() throws Exception { protected void processClientWithCorrectAccessTokenPushNodeMetricBuildPrimitiveSimple() throws Exception {
clientWithCorrectNodeAccessTokenWithNDEATH();; List<String> listKeys = new ArrayList<>();
clientWithCorrectNodeAccessTokenWithNDEATH();
String messageTypeName = SparkplugMessageType.NDATA.name(); String messageTypeName = SparkplugMessageType.NDATA.name();
List<String> listKeys = new ArrayList<>();
List<TsKvEntry> listTsKvEntry = new ArrayList<>(); List<TsKvEntry> listTsKvEntry = new ArrayList<>();
SparkplugBProto.Payload.Builder ndataPayload = SparkplugBProto.Payload.newBuilder() SparkplugBProto.Payload.Builder ndataPayload = SparkplugBProto.Payload.newBuilder()
@ -166,14 +166,13 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
.atMost(40, TimeUnit.SECONDS) .atMost(40, TimeUnit.SECONDS)
.until(() -> { .until(() -> {
finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId())); finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId()));
return finalFuture.get().get().size() == listTsKvEntry.size(); return finalFuture.get().get().size() == (listTsKvEntry.size() + 1);
}); });
Assert.assertTrue("Expected tsKvEntrys is not equal Actual tsKvEntrys", listTsKvEntry.containsAll(finalFuture.get().get())); Assert.assertTrue("Actual tsKvEntrys is not containsAll Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry));
Assert.assertTrue("Actual tsKvEntrys is not equal Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry));
} }
protected void processClientWithCorrectAccessTokenPushNodeMetricBuildArraysSimple() throws Exception { protected void processClientWithCorrectAccessTokenPushNodeMetricBuildArraysPrimitiveSimple() throws Exception {
clientWithCorrectNodeAccessTokenWithNDEATH();; clientWithCorrectNodeAccessTokenWithNDEATH();
String messageTypeName = SparkplugMessageType.NDATA.name(); String messageTypeName = SparkplugMessageType.NDATA.name();
List<String> listKeys = new ArrayList<>(); List<String> listKeys = new ArrayList<>();
@ -184,7 +183,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
.setSeq(getSeqNum()); .setSeq(getSeqNum());
long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS;
createdAddMetricValueArraysTsKv(listTsKvEntry, listKeys, ndataPayload, ts); createdAddMetricValueArraysPrimitiveTsKv(listTsKvEntry, listKeys, ndataPayload, ts);
if (client.isConnected()) { if (client.isConnected()) {
client.publish(NAMESPACE + "/" + groupId + "/" + messageTypeName + "/" + edgeNode, client.publish(NAMESPACE + "/" + groupId + "/" + messageTypeName + "/" + edgeNode,
@ -196,10 +195,9 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
.atMost(40, TimeUnit.SECONDS) .atMost(40, TimeUnit.SECONDS)
.until(() -> { .until(() -> {
finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId())); finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId()));
return finalFuture.get().get().size() == listTsKvEntry.size(); return finalFuture.get().get().size() == (listTsKvEntry.size() + 1);
}); });
Assert.assertTrue("Expected tsKvEntrys is not equal Actual tsKvEntrys", listTsKvEntry.containsAll(finalFuture.get().get())); Assert.assertTrue("Actual tsKvEntrys is not containsAll Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry));
Assert.assertTrue("Actual tsKvEntrys is not equal Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry));
} }
private void createdAddMetricValuePrimitiveTsKv(List<TsKvEntry> listTsKvEntry, List<String> listKeys, private void createdAddMetricValuePrimitiveTsKv(List<TsKvEntry> listTsKvEntry, List<String> listKeys,
@ -229,12 +227,8 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt16(), ts, UInt16)); listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt16(), ts, UInt16));
listKeys.add(keys); listKeys.add(keys);
keys = "MyUInt32I"; keys = "MyUInt32";
listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt32I(), ts, UInt32)); listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt32(), ts, UInt32));
listKeys.add(keys);
keys = "MyUInt32L";
listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt32L(), ts, UInt32));
listKeys.add(keys); listKeys.add(keys);
keys = "MyUInt64"; keys = "MyUInt64";
@ -271,62 +265,62 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
} }
private void createdAddMetricValueArraysTsKv(List<TsKvEntry> listTsKvEntry, List<String> listKeys, private void createdAddMetricValueArraysPrimitiveTsKv(List<TsKvEntry> listTsKvEntry, List<String> listKeys,
SparkplugBProto.Payload.Builder dataPayload, long ts) throws ThingsboardException { SparkplugBProto.Payload.Builder dataPayload, long ts) throws ThingsboardException {
String keys = "MyBytesArray"; String keys = "MyBytesArray";
byte[] bytes = {nextInt8(), nextInt8(), nextInt8()}; byte[] bytes = {nextInt8(), nextInt8(), nextInt8()};
createdAddMetricTsKvJson(dataPayload, keys, bytes, ts, Bytes, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, bytes, ts, Bytes, listTsKvEntry, listKeys);
keys = "MyInt8Array"; keys = "MyInt8Array";
byte[] int8s = {nextInt8(), nextInt8(), nextInt8()}; Byte[] int8s = {nextInt8(), nextInt8(), nextInt8()};
createdAddMetricTsKvJson(dataPayload, keys, int8s, ts, Int8Array, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, int8s, ts, Int8Array, listTsKvEntry, listKeys);
keys = "MyInt16Array"; keys = "MyInt16Array";
short[] int16s = {nextInt16(), nextInt16(), nextInt16()}; Short[] int16s = {nextInt16(), nextInt16(), nextInt16()};
createdAddMetricTsKvJson(dataPayload, keys, int16s, ts, Int16Array, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, int16s, ts, Int16Array, listTsKvEntry, listKeys);
keys = "MyInt32Array";
int[] int32s = {nextInt32(), nextInt32(), nextInt32()};
createdAddMetricTsKvJson(dataPayload, keys, int32s, ts, Int32Array, listTsKvEntry, listKeys);
keys = "MyInt64Array";
long[] int64s = {nextInt64(), nextInt64(), nextInt64()};
createdAddMetricTsKvJson(dataPayload, keys, int64s, ts, Int64Array, listTsKvEntry, listKeys);
keys = "MyUInt8Array"; keys = "MyUInt8Array";
short[] uInt8s = {nextUInt16(), nextUInt16(), nextUInt16()}; Short[] uInt8s = {nextUInt8(), nextUInt8(), nextUInt8()};
createdAddMetricTsKvJson(dataPayload, keys, uInt8s, ts, UInt8Array, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, uInt8s, ts, UInt8Array, listTsKvEntry, listKeys);
keys = "MyInt32Array";
Integer[] int32s = {nextInt32(), nextInt32(), nextInt32()};
createdAddMetricTsKvJson(dataPayload, keys, int32s, ts, Int32Array, listTsKvEntry, listKeys);
keys = "MyUInt16Array"; keys = "MyUInt16Array";
int[] uInt16s = {nextUInt16(), nextUInt16(), nextUInt16()}; Integer[] uInt16s = {nextUInt16(), nextUInt16(), nextUInt16()};
createdAddMetricTsKvJson(dataPayload, keys, uInt16s, ts, UInt16Array, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, uInt16s, ts, UInt16Array, listTsKvEntry, listKeys);
keys = "MyInt64Array";
Long[] int64s = {nextInt64(), nextInt64(), nextInt64()};
createdAddMetricTsKvJson(dataPayload, keys, int64s, ts, Int64Array, listTsKvEntry, listKeys);
keys = "MyUInt32LArray"; keys = "MyUInt32LArray";
long[] uInt32Ls = {nextUInt32L(), nextUInt32L(), nextUInt32L()}; Long[] uInt32Ls = {nextUInt32(), nextUInt32(), nextUInt32()};
createdAddMetricTsKvJson(dataPayload, keys, uInt32Ls, ts, UInt32Array, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, uInt32Ls, ts, UInt32Array, listTsKvEntry, listKeys);
keys = "MyDateTimeArray";
Long[] dateTimes = {nextDateTime(), nextDateTime(), nextDateTime()};
createdAddMetricTsKvJson(dataPayload, keys, dateTimes, ts, DateTimeArray, listTsKvEntry, listKeys);
keys = "MyUInt64Array"; keys = "MyUInt64Array";
long[] uInt64s = {nextUInt64(), nextUInt64(), nextUInt64()}; Long[] uInt64s = {nextUInt64(), nextUInt64(), nextUInt64()};
createdAddMetricTsKvJson(dataPayload, keys, uInt64s, ts, UInt64Array, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, uInt64s, ts, UInt64Array, listTsKvEntry, listKeys);
keys = "MyFloatArray"; keys = "MyFloatArray";
float[] floats = {nextFloat(0,300), nextFloat(0,4000), nextFloat(10,10000)}; Float[] floats = {nextFloat(0, 300), nextFloat(0, 4000), nextFloat(10, 10000)};
createdAddMetricTsKvJson(dataPayload, keys, floats, ts, FloatArray, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, floats, ts, FloatArray, listTsKvEntry, listKeys);
keys = "MyDateTimeArray";
long[] dateTimes = {nextDateTime(), nextDateTime(), nextDateTime()};
createdAddMetricTsKvJson(dataPayload, keys, dateTimes, ts, DateTimeArray, listTsKvEntry, listKeys);
keys = "MyDoubleArray"; keys = "MyDoubleArray";
double [] doubles = {nextDouble(), nextDouble(), nextDouble()}; Double[] doubles = {nextDouble(), nextDouble(), nextDouble()};
createdAddMetricTsKvJson(dataPayload, keys, doubles, ts, DoubleArray, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, doubles, ts, DoubleArray, listTsKvEntry, listKeys);
keys = "MyBooleanArray"; keys = "MyBooleanArray";
boolean [] booleans = {nextBoolean(), nextBoolean(), nextBoolean()}; Boolean[] booleans = {nextBoolean(), nextBoolean(), nextBoolean()};
createdAddMetricTsKvJson(dataPayload, keys, booleans, ts, BooleanArray, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, booleans, ts, BooleanArray, listTsKvEntry, listKeys);
keys = "MyStringArray"; keys = "MyStringArray";
String [] strings = {nexString(), nexString(), nexString()}; String[] strings = {nexString(), nexString(), nexString()};
createdAddMetricTsKvJson(dataPayload, keys, strings, ts, StringArray, listTsKvEntry, listKeys); createdAddMetricTsKvJson(dataPayload, keys, strings, ts, StringArray, listTsKvEntry, listKeys);
} }
@ -337,18 +331,18 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
return tsKvEntry; return tsKvEntry;
} }
private TsKvEntry createdAddMetricTsKvFloat(SparkplugBProto.Payload.Builder dataPayload, String key, Object value, private TsKvEntry createdAddMetricTsKvFloat(SparkplugBProto.Payload.Builder dataPayload, String key, float value,
long ts, MetricDataType metricDataType) throws ThingsboardException { long ts, MetricDataType metricDataType) throws ThingsboardException {
var f = new BigDecimal(String.valueOf(value)); Double dd = Double.parseDouble(Float.toString(value));
TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new DoubleDataEntry(key, f.doubleValue())); TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new DoubleDataEntry(key, dd));
dataPayload.addMetrics(createMetric(value, ts, key, metricDataType)); dataPayload.addMetrics(createMetric(value, ts, key, metricDataType));
return tsKvEntry; return tsKvEntry;
} }
private TsKvEntry createdAddMetricTsKvDouble(SparkplugBProto.Payload.Builder dataPayload, String key, double value, private TsKvEntry createdAddMetricTsKvDouble(SparkplugBProto.Payload.Builder dataPayload, String key, double value,
long ts, MetricDataType metricDataType) throws ThingsboardException { long ts, MetricDataType metricDataType) throws ThingsboardException {
var d = new BigDecimal(String.valueOf(value)); Long l = Double.valueOf(value).longValue();
TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(key, d.longValueExact())); TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(key, l));
dataPayload.addMetrics(createMetric(value, ts, key, metricDataType)); dataPayload.addMetrics(createMetric(value, ts, key, metricDataType));
return tsKvEntry; return tsKvEntry;
} }
@ -368,20 +362,24 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
} }
private void createdAddMetricTsKvJson(SparkplugBProto.Payload.Builder dataPayload, String key, private void createdAddMetricTsKvJson(SparkplugBProto.Payload.Builder dataPayload, String key,
Object values, long ts, MetricDataType metricDataType, Object values, long ts, MetricDataType metricDataType,
List<TsKvEntry> listTsKvEntry, List<TsKvEntry> listTsKvEntry,
List<String> listKeys) throws ThingsboardException { List<String> listKeys) throws ThingsboardException {
ArrayNode nodeArray = newArrayNode(); ArrayNode nodeArray = newArrayNode();
switch (metricDataType) { switch (metricDataType) {
case Bytes: case Bytes:
for (byte b : (byte[]) values) {
nodeArray.add(b);
}
break;
case Int8Array: case Int8Array:
for (byte b : (byte[])values) { for (Byte b : (Byte[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
break; break;
case Int16Array: case Int16Array:
case UInt8Array: case UInt8Array:
for (short b : (short[])values) { for (Short b : (Short[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
break; break;
@ -391,33 +389,33 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
case UInt32Array: case UInt32Array:
case UInt64Array: case UInt64Array:
case DateTimeArray: case DateTimeArray:
if (values instanceof int[]) { if (values instanceof Integer[]) {
for (int b : (int[])values) { for (Integer b : (Integer[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
} else { } else {
for (long b : (long[])values) { for (Long b : (Long[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
} }
break; break;
case DoubleArray: case DoubleArray:
for (double b : (double[])values) { for (Double b : (Double[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
break; break;
case FloatArray: case FloatArray:
for (float b : (float[])values) { for (Float b : (Float[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
break; break;
case BooleanArray: case BooleanArray:
for (boolean b : (boolean[])values) { for (Boolean b : (Boolean[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
break; break;
case StringArray: case StringArray:
for (String b : (String[])values) { for (String b : (String[]) values) {
nodeArray.add(b); nodeArray.add(b);
} }
break; break;
@ -446,19 +444,15 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac
return (short) random.nextInt(Short.MIN_VALUE, Short.MAX_VALUE); return (short) random.nextInt(Short.MIN_VALUE, Short.MAX_VALUE);
} }
private short nextUInt16() { private int nextUInt16() {
return (short) random.nextInt(0, Short.MAX_VALUE * 2 + 1); return random.nextInt(0, Short.MAX_VALUE * 2 + 1);
} }
private int nextInt32() { private int nextInt32() {
return random.nextInt(Integer.MIN_VALUE, Integer.MAX_VALUE); return random.nextInt(Integer.MIN_VALUE, Integer.MAX_VALUE);
} }
private int nextUInt32I() { private long nextUInt32() {
return random.nextInt(0, Integer.MAX_VALUE);
}
private long nextUInt32L() {
long l = Integer.MAX_VALUE; long l = Integer.MAX_VALUE;
return random.nextLong(0, l * 2 + 1); return random.nextLong(0, l * 2 + 1);
} }

4
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java

@ -54,8 +54,8 @@ public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSpa
} }
@Test @Test
public void testClientWithCorrectAccessTokenPushNodeMetricBuildPArraysSimple() throws Exception { public void testClientWithCorrectAccessTokenPushNodeMetricBuildPArraysPrimitiveSimple() throws Exception {
processClientWithCorrectAccessTokenPushNodeMetricBuildArraysSimple(); processClientWithCorrectAccessTokenPushNodeMetricBuildArraysPrimitiveSimple();
} }
} }

85
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -111,7 +111,6 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe
import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_OPEN; import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_OPEN;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.NDEATH; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.NDEATH;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopicPublish; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopicPublish;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopicSubscribe;
/** /**
* @author Andrew Shvayka * @author Andrew Shvayka
@ -129,7 +128,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private final UUID sessionId; private final UUID sessionId;
protected final MqttTransportContext context; protected final MqttTransportContext context;
private final TransportService transportService; public final TransportService transportService;
private final SchedulerComponent scheduler; private final SchedulerComponent scheduler;
private final SslHandler sslHandler; private final SslHandler sslHandler;
private final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap; private final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap;
@ -143,7 +142,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private final ConcurrentHashMap<String, Integer> chunkSizes; private final ConcurrentHashMap<String, Integer> chunkSizes;
private final ConcurrentMap<Integer, TransportProtos.ToDeviceRpcRequestMsg> rpcAwaitingAck; private final ConcurrentMap<Integer, TransportProtos.ToDeviceRpcRequestMsg> rpcAwaitingAck;
private TopicType attrSubTopicType; public TopicType attrSubTopicType;
private TopicType rpcSubTopicType; private TopicType rpcSubTopicType;
private TopicType attrReqTopicType; private TopicType attrReqTopicType;
private TopicType toServerRpcSubTopicType; private TopicType toServerRpcSubTopicType;
@ -451,62 +450,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
} }
private void handleSparkplugSubscribeMsg(List<Integer> grantedQoSList, MqttTopicSubscription subscription, MqttQoS reqQoS) throws ThingsboardException {
SparkplugTopic sparkplugTopic = parseTopicSubscribe(subscription.topicName());
if (sparkplugTopic.getGroupId() == null) {
// TODO SUBSCRIBE NameSpace
} else if (sparkplugTopic.getType() == null) {
// TODO SUBSCRIBE GroupId
} else if (sparkplugTopic.isNode()) {
// A node topic
processAttributesSubscribe(grantedQoSList, MqttTopics.DEVICE_ATTRIBUTES_TOPIC, reqQoS, TopicType.V1);
switch (sparkplugTopic.getType()) {
case STATE:
// TODO
break;
case NBIRTH:
// TODO
break;
case NCMD:
// TODO
break;
case NDATA:
// TODO
break;
case NDEATH:
// TODO
break;
case NRECORD:
// TODO
break;
default:
}
} else {
// A device topic
switch (sparkplugTopic.getType()) {
case STATE:
// TODO
break;
case DBIRTH:
// TODO
break;
case DCMD:
// TODO
break;
case DDATA:
// TODO
break;
case DDEATH:
// TODO
break;
case DRECORD:
// TODO
break;
default:
}
}
}
private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) { private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) {
try { try {
Matcher fwMatcher; Matcher fwMatcher;
@ -768,7 +711,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
MqttQoS reqQoS = subscription.qualityOfService(); MqttQoS reqQoS = subscription.qualityOfService();
try { try {
if (sparkplugSessionHandler != null) { if (sparkplugSessionHandler != null) {
handleSparkplugSubscribeMsg(grantedQoSList, subscription, reqQoS); sparkplugSessionHandler.handleSparkplugSubscribeMsg(grantedQoSList, subscription, reqQoS);
activityReported = true; activityReported = true;
} else { } else {
switch (topic) { switch (topic) {
@ -853,13 +796,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
registerSubQoS(topic, grantedQoSList, reqQoS); registerSubQoS(topic, grantedQoSList, reqQoS);
} }
private void processAttributesSubscribe(List<Integer> grantedQoSList, String topic, MqttQoS reqQoS, TopicType topicType) { public void processAttributesSubscribe(List<Integer> grantedQoSList, String topic, MqttQoS reqQoS, TopicType topicType) {
transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null); transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null);
attrSubTopicType = topicType; attrSubTopicType = topicType;
registerSubQoS(topic, grantedQoSList, reqQoS); registerSubQoS(topic, grantedQoSList, reqQoS);
} }
private void registerSubQoS(String topic, List<Integer> grantedQoSList, MqttQoS reqQoS) { public void registerSubQoS(String topic, List<Integer> grantedQoSList, MqttQoS reqQoS) {
grantedQoSList.add(getMinSupportedQos(reqQoS)); grantedQoSList.add(getMinSupportedQos(reqQoS));
mqttQoSMap.put(new MqttTopicMatcher(topic), getMinSupportedQos(reqQoS)); mqttQoSMap.put(new MqttTopicMatcher(topic), getMinSupportedQos(reqQoS));
} }
@ -1145,10 +1088,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (sparkplugSessionHandler == null) { if (sparkplugSessionHandler == null) {
SparkplugTopic sparkplugTopicNode = validatedSparkplugTopicConnectedNode(connectMessage); SparkplugTopic sparkplugTopicNode = validatedSparkplugTopicConnectedNode(connectMessage);
if (sparkplugTopicNode != null) { if (sparkplugTopicNode != null) {
SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes());
sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId, sparkplugTopicNode); sparkplugSessionHandler = new SparkplugNodeSessionHandler(this, deviceSessionCtx, sessionId, sparkplugTopicNode);
sparkplugSessionHandler.onTelemetryProto(0, sparkplugBProtoNode, sparkplugSessionHandler.onTelemetryProto(0, sparkplugBProtoNode,
deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopicNode); deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopicNode);
} else { } else {
log.trace("[{}][{}] Failed to fetch sparkplugDevice connect: sparkplugTopicName without SparkplugMessageType.NDEATH.", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName()); log.trace("[{}][{}] Failed to fetch sparkplugDevice connect: sparkplugTopicName without SparkplugMessageType.NDEATH.", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName());
throw new ThingsboardException("Invalid request body", ThingsboardErrorCode.BAD_REQUEST_PARAMS); throw new ThingsboardException("Invalid request body", ThingsboardErrorCode.BAD_REQUEST_PARAMS);
@ -1161,13 +1104,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
} }
private SparkplugTopic validatedSparkplugTopicConnectedNode (MqttConnectMessage connectMessage) throws ThingsboardException { private SparkplugTopic validatedSparkplugTopicConnectedNode(MqttConnectMessage connectMessage) throws ThingsboardException {
if(StringUtils.isNotBlank(connectMessage.payload().willTopic()) if (StringUtils.isNotBlank(connectMessage.payload().willTopic())
&& connectMessage.payload().willMessageInBytes() != null && connectMessage.payload().willMessageInBytes() != null
&& connectMessage.payload().willMessageInBytes().length > 0) { && connectMessage.payload().willMessageInBytes().length > 0) {
SparkplugTopic sparkplugTopicNode = parseTopicPublish(connectMessage.payload().willTopic()); SparkplugTopic sparkplugTopicNode = parseTopicPublish(connectMessage.payload().willTopic());
if(NDEATH.equals(sparkplugTopicNode.getType())){ if (NDEATH.equals(sparkplugTopicNode.getType())) {
return sparkplugTopicNode; return sparkplugTopicNode;
} }
} }
return null; return null;

38
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java

@ -22,7 +22,10 @@ import com.google.gson.JsonParser;
import com.google.gson.JsonSyntaxException; import com.google.gson.JsonSyntaxException;
import com.google.protobuf.Descriptors; import com.google.protobuf.Descriptors;
import io.netty.handler.codec.mqtt.MqttPublishMessage; import io.netty.handler.codec.mqtt.MqttPublishMessage;
import io.netty.handler.codec.mqtt.MqttQoS;
import io.netty.handler.codec.mqtt.MqttTopicSubscription;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.device.profile.MqttTopics;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.transport.adaptor.AdaptorException; import org.thingsboard.server.common.transport.adaptor.AdaptorException;
@ -30,6 +33,8 @@ import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.common.transport.adaptor.ProtoConverter; import org.thingsboard.server.common.transport.adaptor.ProtoConverter;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto;
import org.thingsboard.server.transport.mqtt.MqttTransportHandler;
import org.thingsboard.server.transport.mqtt.TopicType;
import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType;
import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic;
@ -48,6 +53,7 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMess
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.fromSparkplugBMetricToKeyValueProto; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.fromSparkplugBMetricToKeyValueProto;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.validatedValueByTypeMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.validatedValueByTypeMetric;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopicSubscribe;
/** /**
* Created by nickAS21 on 12.12.22 * Created by nickAS21 on 12.12.22
@ -57,10 +63,12 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler {
private final SparkplugTopic sparkplugTopicNode; private final SparkplugTopic sparkplugTopicNode;
private final Map<String, SparkplugBProto.Payload.Metric> nodeBirthMetrics; private final Map<String, SparkplugBProto.Payload.Metric> nodeBirthMetrics;
private final MqttTransportHandler parent;
public SparkplugNodeSessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId, public SparkplugNodeSessionHandler(MqttTransportHandler parent, DeviceSessionCtx deviceSessionCtx, UUID sessionId,
SparkplugTopic sparkplugTopicNode) { SparkplugTopic sparkplugTopicNode) {
super(deviceSessionCtx, sessionId); super(deviceSessionCtx, sessionId);
this.parent = parent;
this.sparkplugTopicNode = sparkplugTopicNode; this.sparkplugTopicNode = sparkplugTopicNode;
this.nodeBirthMetrics = new ConcurrentHashMap<>(); this.nodeBirthMetrics = new ConcurrentHashMap<>();
} }
@ -129,6 +137,34 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler {
} }
} }
public void handleSparkplugSubscribeMsg(List<Integer> grantedQoSList, MqttTopicSubscription subscription,
MqttQoS reqQoS) throws ThingsboardException, AdaptorException,
ExecutionException, InterruptedException {
SparkplugTopic sparkplugTopic = parseTopicSubscribe(subscription.topicName());
if (sparkplugTopic.getGroupId() == null) {
// TODO SUBSCRIBE NameSpace
} else if (sparkplugTopic.getType() == null) {
// TODO SUBSCRIBE GroupId
} else if (sparkplugTopic.isNode()) {
// SUBSCRIBE Node
parent.processAttributesSubscribe(grantedQoSList, MqttTopics.DEVICE_ATTRIBUTES_TOPIC, reqQoS, TopicType.V1);
} else {
// SUBSCRIBE Device
onSparkplugDeviceSubscribe(grantedQoSList, MqttTopics.DEVICE_ATTRIBUTES_TOPIC, reqQoS, TopicType.V1, sparkplugTopic.getDeviceId());
}
}
public void onSparkplugDeviceSubscribe(List<Integer> grantedQoSList, String topic,
MqttQoS reqQoS, TopicType topicType, String deviceName)
throws AdaptorException, ThingsboardException, ExecutionException, InterruptedException {
checkDeviceName(deviceName);
ListenableFuture<MqttDeviceAwareSessionContext> contextListenableFuture = onDeviceConnectProto(deviceName);
parent.transportService.process(contextListenableFuture.get().getSessionInfo(),
TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null);
parent.attrSubTopicType = topicType;
parent.registerSubQoS(topic, grantedQoSList, reqQoS);
}
public void onDeviceDisconnect(MqttPublishMessage mqttMsg, String deviceName) throws AdaptorException { public void onDeviceDisconnect(MqttPublishMessage mqttMsg, String deviceName) throws AdaptorException {
try { try {
processOnDisconnect(mqttMsg, deviceName); processOnDisconnect(mqttMsg, deviceName);

282
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java

@ -16,11 +16,14 @@
package org.thingsboard.server.transport.mqtt.util.sparkplug; package org.thingsboard.server.transport.mqtt.util.sparkplug;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.annotation.JsonSerialize; import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.ser.std.FileSerializer; import com.fasterxml.jackson.databind.ser.std.FileSerializer;
import com.google.protobuf.ByteString; import com.google.protobuf.ByteString;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
@ -42,6 +45,7 @@ import java.nio.LongBuffer;
import java.nio.ShortBuffer; import java.nio.ShortBuffer;
import java.text.NumberFormat; import java.text.NumberFormat;
import java.util.Arrays; import java.util.Arrays;
import java.util.List;
import java.util.Optional; import java.util.Optional;
import static org.thingsboard.common.util.JacksonUtil.newArrayNode; import static org.thingsboard.common.util.JacksonUtil.newArrayNode;
@ -228,34 +232,39 @@ public class SparkplugMetricUtil {
.setDatatype(metricDataType.toIntValue()) .setDatatype(metricDataType.toIntValue())
.build(); .build();
switch (metricDataType) { switch (metricDataType) {
case Int8: case Int8: // (byte)
case Int16: return metric.toBuilder().setIntValue(((Byte) value).intValue()).build();
case Int16: // (short)
case UInt8: case UInt8:
case UInt16: return metric.toBuilder().setIntValue(((Short) value).intValue()).build();
case UInt16: // (int)
case Int32: case Int32:
return metric.toBuilder().setIntValue(((Integer) value).intValue()).build(); return metric.toBuilder().setIntValue(((Integer) value).intValue()).build();
case UInt32: case UInt32: // (long)
case Int64: case Int64:
case UInt64: case UInt64:
case DateTime: case DateTime:
return metric.toBuilder().setLongValue(((Long) value).longValue()).build(); return metric.toBuilder().setLongValue(((Long) value).longValue()).build();
case Float: case Float: // (float)
return metric.toBuilder().setFloatValue(((Float) value).floatValue()).build(); return metric.toBuilder().setFloatValue(((Float) value).floatValue()).build();
case Double: case Double: // (double)
return metric.toBuilder().setDoubleValue(((Double) value).doubleValue()).build(); return metric.toBuilder().setDoubleValue(((Double) value).doubleValue()).build();
case Boolean: case Boolean: // (boolean)
return metric.toBuilder().setBooleanValue(((Boolean) value).booleanValue()).build(); return metric.toBuilder().setBooleanValue(((Boolean) value).booleanValue()).build();
case String: case String: // String)
case Text: case Text:
case UUID: case UUID:
return metric.toBuilder().setStringValue((String) value).build(); return metric.toBuilder().setStringValue((String) value).build();
case Bytes: case Bytes:
case Int8Array:
ByteString byteString = ByteString.copyFrom((byte[]) value); ByteString byteString = ByteString.copyFrom((byte[]) value);
return metric.toBuilder().setBytesValue(byteString).build(); return metric.toBuilder().setBytesValue(byteString).build();
case Int8Array:
byte [] int8Array = byteArrayToByteArray((Byte[]) value);
ByteString byteInt8Array = ByteString.copyFrom((int8Array));
return metric.toBuilder().setBytesValue(byteInt8Array).build();
case Int16Array: case Int16Array:
case UInt8Array: case UInt8Array:
byte[] int16Array = shortArrayToByteArray((short[]) value); byte[] int16Array = shortArrayToByteArray((Short[]) value);
ByteString byteInt16Array = ByteString.copyFrom((int16Array)); ByteString byteInt16Array = ByteString.copyFrom((int16Array));
return metric.toBuilder().setBytesValue(byteInt16Array).build(); return metric.toBuilder().setBytesValue(byteInt16Array).build();
case Int32Array: case Int32Array:
@ -264,25 +273,25 @@ public class SparkplugMetricUtil {
case UInt32Array: case UInt32Array:
case UInt64Array: case UInt64Array:
case DateTimeArray: case DateTimeArray:
if (value instanceof int[]) { if (value instanceof Integer[]) {
byte[] int32Array = integerArrayToByteArray((int[]) value); byte[] int32Array = integerArrayToByteArray((Integer[]) value);
ByteString byteInt32Array = ByteString.copyFrom((int32Array)); ByteString byteInt32Array = ByteString.copyFrom((int32Array));
return metric.toBuilder().setBytesValue(byteInt32Array).build(); return metric.toBuilder().setBytesValue(byteInt32Array).build();
} else { } else {
byte[] int64Array = longArrayToByteArray((long[]) value); byte[] int64Array = longArrayToByteArray((Long[]) value);
ByteString byteInt64Array = ByteString.copyFrom((int64Array)); ByteString byteInt64Array = ByteString.copyFrom((int64Array));
return metric.toBuilder().setBytesValue(byteInt64Array).build(); return metric.toBuilder().setBytesValue(byteInt64Array).build();
} }
case DoubleArray: case DoubleArray:
byte[] doubleArray = doublArrayToByteArray((double[]) value); byte[] doubleArray = doublArrayToByteArray((Double[]) value);
ByteString byteDoubleArray = ByteString.copyFrom(doubleArray); ByteString byteDoubleArray = ByteString.copyFrom(doubleArray);
return metric.toBuilder().setBytesValue(byteDoubleArray).build(); return metric.toBuilder().setBytesValue(byteDoubleArray).build();
case FloatArray: case FloatArray:
byte[] floatArray = floatArrayToByteArray((float[]) value); byte[] floatArray = floatArrayToByteArray((Float[]) value);
ByteString byteFloatArray = ByteString.copyFrom(floatArray); ByteString byteFloatArray = ByteString.copyFrom(floatArray);
return metric.toBuilder().setBytesValue(byteFloatArray).build(); return metric.toBuilder().setBytesValue(byteFloatArray).build();
case BooleanArray: case BooleanArray:
byte[] booleanArray = booleanArrayToByteArray((boolean[]) value); byte[] booleanArray = booleanArrayToByteArray((Boolean[]) value);
ByteString byteBooleanArray = ByteString.copyFrom(booleanArray); ByteString byteBooleanArray = ByteString.copyFrom(booleanArray);
return metric.toBuilder().setBytesValue(byteBooleanArray).build(); return metric.toBuilder().setBytesValue(byteBooleanArray).build();
case StringArray: case StringArray:
@ -307,10 +316,14 @@ public class SparkplugMetricUtil {
if (kv.getTypeValue() <= 3) { if (kv.getTypeValue() <= 3) {
return validatedValuePrimitiveByTypeMetric(kv, metricDataType); return validatedValuePrimitiveByTypeMetric(kv, metricDataType);
} else if (kv.getTypeValue() == 4) { } else if (kv.getTypeValue() == 4) {
return validatedValueJsonByTypeMetric(kv, metricDataType); JsonNode arrayNode = JacksonUtil.fromString(kv.getJsonV(), JsonNode.class);
if (arrayNode.isArray()) {
return validatedValueJsonByTypeMetric(kv.getJsonV(), metricDataType);
}
} else { } else {
throw new ThingsboardException("Invalid type KeyValueProto " + kv.toString() + " for MetricDataType " + metricDataType.name(), ThingsboardErrorCode.INVALID_ARGUMENTS); throw new ThingsboardException("Invalid type KeyValueProto " + kv.toString() + " for MetricDataType " + metricDataType.name(), ThingsboardErrorCode.INVALID_ARGUMENTS);
} }
return Optional.empty();
} }
public static Optional<Object> validatedValuePrimitiveByTypeMetric(TransportProtos.KeyValueProto kv, MetricDataType metricDataType) throws ThingsboardException { public static Optional<Object> validatedValuePrimitiveByTypeMetric(TransportProtos.KeyValueProto kv, MetricDataType metricDataType) throws ThingsboardException {
@ -392,188 +405,115 @@ public class SparkplugMetricUtil {
return Optional.empty(); return Optional.empty();
} }
public static Optional<Object> validatedValueJsonByTypeMetric(String arrayNodeStr, MetricDataType metricDataType) {
public static Optional<Object> validatedValueJsonByTypeMetric(TransportProtos.KeyValueProto kv, MetricDataType metricDataType) { try {
// try { Optional<Object> valueOpt;
// Optional<Object> valueOpt; switch (metricDataType) {
// switch (metricDataType) { // byte[]
// // int case Bytes:
// case Int8: List<Byte> listBytes = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// case Int16: byte[] bytes = new byte[listBytes.size()];
// case UInt8: for(int i = 0; i < listBytes.size(); i++) {
// case UInt16: bytes[i] = listBytes.get(i).byteValue();
// valueOpt = getValueKvProtoPrimitive(kv); }
// return valueOpt.isPresent() ? Optional.of(Integer.valueOf(String.valueOf(valueOpt.get()))) : valueOpt; return Optional.of(bytes);
// // int/long // Byte []
// case Int32: case Int8Array:
// case UInt32: List<Byte> listInt8Array = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// valueOpt = getValueKvProtoPrimitive(kv); Byte[] int8Arrays = listInt8Array.toArray(Byte[]::new) ;
// try { return Optional.of(int8Arrays);
// return Optional.of(Integer.valueOf(String.valueOf(valueOpt.get()))); // Short[]
// } catch (NumberFormatException e) { case Int16Array:
// return Optional.of(Long.valueOf(String.valueOf(valueOpt.get()))); case UInt8Array:
// } List<Short> listShorts = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// // long return Optional.of(listShorts.toArray(Short[]::new));
// case Int64: // Integer []
// case UInt64: case UInt16Array:
// case DateTime: case Int32Array:
// valueOpt = getValueKvProtoPrimitive(kv); List<Integer> listIntegers = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// return Optional.of(Long.valueOf(String.valueOf(valueOpt.get()))); return Optional.of(listIntegers.toArray(Integer[]::new));
// // float // Long[]
// case Float: case UInt32Array:
// valueOpt = getValueKvProtoPrimitive(kv); case Int64Array:
// var f = new BigDecimal(String.valueOf(kv.getDoubleV())); case UInt64Array:
// return Optional.of(f.floatValue()); case DateTimeArray:
// } List<Long> listLongs = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// break; return Optional.of(listLongs.toArray(Long[]::new));
// // double // Double []
// case Double: case DoubleArray:
// if (kv.getTypeValue() == 1) { List<Double> listDoubles = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// return Optional.of(kv.getLongV()); return Optional.of(listDoubles.toArray(Double[]::new));
// } // Float[]
// break; case FloatArray:
// case Boolean: List<Float> listFloats = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// if (kv.getTypeValue() == 0) { // ok 0 return Optional.of(listFloats.toArray(Float[]::new));
// return Optional.of(kv.getBoolV()); // Boolean[]
// } case BooleanArray:
// break; List<Boolean> listBooleans = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// case String: return Optional.of(listBooleans.toArray(Boolean[]::new));
// case Text: case StringArray:
// case UUID: List<String> listStrings = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {});
// if (kv.getTypeValue() == 4) { return Optional.of(listStrings.toArray(String[]::new));
// return Optional.of(kv.getStringV()); case DataSet:
// } case File:
// break; case Template:
// // byte[] log.error("Invalid type value [{}] for MetricDataType [{}]", arrayNodeStr, metricDataType.name());
// case Bytes: return Optional.empty();
// case Int8Array: case Unknown:
// if (kv.getTypeValue() == 5) { default:
//// ByteString byteString = ByteString.copyFrom((byte[]) value); log.error("Invalid MetricDataType [{}] type, value [{}]", arrayNodeStr, metricDataType.name());
//// return metric.toBuilder().setBytesValue(byteString).build(); return Optional.empty();
// return Optional.of(kv.getJsonV()); }
// } } catch (Exception e) {
// break; log.error("Invalid type value [{}] for MetricDataType [{}] [{}]", arrayNodeStr, metricDataType.name(), e.getMessage());
// // short[] return Optional.empty();
// case Int16Array: }
// case UInt8Array:
// if (kv.getTypeValue() == 5) {
//// byte[] int16Array = shortArrayToByteArray((short[]) value);
//// ByteString byteInt16Array = ByteString.copyFrom((int16Array));
// return Optional.of(kv.getJsonV());
// }
// break;
// // int []
// case UInt16Array:
// case Int32Array:
//
// // int[] / long[]
// case UInt32Array:
//
// // long[]
// case Int64Array:
// case UInt64Array:
// case DateTimeArray:
// if (kv.getTypeValue() == 5) {
//// if (value instanceof int[]) {
//// byte[] int32Array = integerArrayToByteArray((int[]) value);
//// ByteString byteInt32Array = ByteString.copyFrom((int32Array));
//// return metric.toBuilder().setBytesValue(byteInt32Array).build();
//// } else {
//// byte[] int64Array = longArrayToByteArray((long[]) value);
//// ByteString byteInt64Array = ByteString.copyFrom((int64Array));
//// return metric.toBuilder().setBytesValue(byteInt64Array).build();
//// }
// return Optional.of(kv.getJsonV());
// }
// break;
// // double []
// case DoubleArray:
// if (kv.getTypeValue() == 5) {
//// byte[] doubleArray = doublArrayToByteArray((double[]) value);
//// ByteString byteDoubleArray = ByteString.copyFrom(doubleArray);
//// return metric.toBuilder().setBytesValue(byteDoubleArray).build();
// return Optional.of(kv.getJsonV());
// }
// break;
// // float[]
// case FloatArray:
// if (kv.getTypeValue() == 5) {
//// byte[] floatArray = floatArrayToByteArray((float[]) value);
//// ByteString byteFloatArray = ByteString.copyFrom(floatArray);
//// return metric.toBuilder().setBytesValue(byteFloatArray).build();
// return Optional.of(kv.getJsonV());
// }
// break;
// // boolean[]
// case BooleanArray:
// if (kv.getTypeValue() == 5) {
//// byte[] booleanArray = booleanArrayToByteArray((boolean[]) value);
//// ByteString byteBooleanArray = ByteString.copyFrom(booleanArray);
//// return metric.toBuilder().setBytesValue(byteBooleanArray).build();
// return Optional.of(kv.getJsonV());
// }
// break;
// // String []
// case StringArray:
// if (kv.getTypeValue() == 5) {
//// byte[] stringArray = stringArrayToByteArray((String[]) value);
//// ByteString byteStringArray = ByteString.copyFrom(stringArray);
//// return metric.toBuilder().setBytesValue(byteStringArray).build();
// return Optional.of(kv.getJsonV());
// }
// break;
// case DataSet:
// case File:
// case Template:
// log.error("Invalid type value [{}] for MetricDataType [{}]", kv, metricDataType.name());
// return Optional.empty();
// case Unknown:
// log.error("Invalid MetricDataType [{}] type, value [{}]", kv, metricDataType.name());
// return Optional.empty();
// } catch (Exception e) {
// log.error("Invalid type value [{}] for MetricDataType [{}] [{}]", kv, metricDataType.name(), e.getMessage());
// return Optional.empty();
// }
return Optional.empty();
} }
private static byte[] byteArrayToByteArray(Byte[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length);
for (Byte d : inputs) {
bb.put(d);
}
return bb.array();
}
private static byte[] shortArrayToByteArray(short[] inputs) { private static byte[] shortArrayToByteArray(Short[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 2); ByteBuffer bb = ByteBuffer.allocate(inputs.length * 2);
for (short d : inputs) { for (Short d : inputs) {
bb.putShort(d); bb.putShort(d);
} }
return bb.array(); return bb.array();
} }
private static byte[] integerArrayToByteArray(int[] inputs) { private static byte[] integerArrayToByteArray(Integer[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 4); ByteBuffer bb = ByteBuffer.allocate(inputs.length * 4);
for (int d : inputs) { for (Integer d : inputs) {
bb.putInt(d); bb.putInt(d);
} }
return bb.array(); return bb.array();
} }
private static byte[] longArrayToByteArray(long[] inputs) { private static byte[] longArrayToByteArray(Long[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8);
for (long d : inputs) { for (Long d : inputs) {
bb.putLong(d); bb.putLong(d);
} }
return bb.array(); return bb.array();
} }
private static byte[] doublArrayToByteArray(double[] inputs) { private static byte[] doublArrayToByteArray(Double[] inputs) {
ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8);
for (double d : inputs) { for (Double d : inputs) {
bb.putDouble(d); bb.putDouble(d);
} }
return bb.array(); return bb.array();
} }
private static byte[] floatArrayToByteArray(float[] inputs) throws ThingsboardException { private static byte[] floatArrayToByteArray(Float[] inputs) throws ThingsboardException {
ByteArrayOutputStream bas = new ByteArrayOutputStream(); ByteArrayOutputStream bas = new ByteArrayOutputStream();
DataOutputStream ds = new DataOutputStream(bas); DataOutputStream ds = new DataOutputStream(bas);
for (float f : inputs) { for (Float f : inputs) {
try { try {
ds.writeFloat(f); ds.writeFloat(f);
} catch (IOException e) { } catch (IOException e) {
@ -583,7 +523,7 @@ public class SparkplugMetricUtil {
return bas.toByteArray(); return bas.toByteArray();
} }
private static byte[] booleanArrayToByteArray(boolean[] inputs) { private static byte[] booleanArrayToByteArray(Boolean[] inputs) {
byte[] toReturn = new byte[inputs.length]; byte[] toReturn = new byte[inputs.length];
for (int entry = 0; entry < toReturn.length; entry++) { for (int entry = 0; entry < toReturn.length; entry++) {
toReturn[entry] = (byte) (inputs[entry] ? 1 : 0); toReturn[entry] = (byte) (inputs[entry] ? 1 : 0);

Loading…
Cancel
Save