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 d57fb2e511..884bd8c721 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 @@ -19,7 +19,10 @@ import lombok.extern.slf4j.Slf4j; import org.eclipse.paho.mqttv5.client.IMqttToken; import org.eclipse.paho.mqttv5.client.MqttConnectionOptions; import org.eclipse.paho.mqttv5.common.MqttMessage; +import org.eclipse.paho.mqttv5.common.packet.MqttConnAck; +import org.eclipse.paho.mqttv5.common.packet.MqttReturnCode; import org.eclipse.paho.mqttv5.common.packet.MqttWireMessage; +import org.junit.Assert; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; 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 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.SparkplugMetricUtil.createMetric; @@ -52,7 +56,6 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte protected int seq = 0; protected static final long PUBLISH_TS_DELTA_MS = 86400000;// Publish start TS <-> 24h - public void beforeSparkplugTest() throws Exception { MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() .gatewayName("Test Connect Sparkplug client node") @@ -62,13 +65,13 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte processBeforeTest(configProperties); } - public MqttWireMessage clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception { + public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception { long ts = calendar.getTimeInMillis(); 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; MetricDataType metricDataType = Int64; SparkplugBProto.Payload.Builder deathPayload = SparkplugBProto.Payload.newBuilder() @@ -84,7 +87,11 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte msg.setPayload(deathBytes); options.setWill(topic, msg); 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 { 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 c02c23e213..f8a4373036 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 @@ -18,9 +18,6 @@ package org.thingsboard.server.transport.mqtt.sparkplug.connection; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; 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.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; @@ -39,7 +36,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; -import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; @@ -52,13 +48,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr protected void processClientWithCorrectNodeAccessTokenWithNDEATH_Test() throws Exception { long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; long value = bdSeq = 0; - MqttWireMessage response = clientWithCorrectNodeAccessTokenWithNDEATH(ts, value); - - Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType()); - - MqttConnAck connAckMsg = (MqttConnAck) response; - - Assert.assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode()); + clientWithCorrectNodeAccessTokenWithNDEATH(ts, value); String keys = SparkplugMessageType.NDEATH.name() + " " + keysBdSeq; TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, value)); @@ -113,5 +103,4 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr Assert.assertEquals(cntDevices, deviceIds.size()); } - } 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 af283c3603..615a255ac2 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 @@ -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.SparkplugMessageType; -import java.math.BigDecimal; import java.util.ArrayList; import java.util.List; import java.util.Optional; @@ -143,10 +142,11 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac } protected void processClientWithCorrectAccessTokenPushNodeMetricBuildPrimitiveSimple() throws Exception { - clientWithCorrectNodeAccessTokenWithNDEATH();; + List listKeys = new ArrayList<>(); + clientWithCorrectNodeAccessTokenWithNDEATH(); String messageTypeName = SparkplugMessageType.NDATA.name(); - List listKeys = new ArrayList<>(); + List listTsKvEntry = new ArrayList<>(); SparkplugBProto.Payload.Builder ndataPayload = SparkplugBProto.Payload.newBuilder() @@ -166,14 +166,13 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac .atMost(40, TimeUnit.SECONDS) .until(() -> { 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 equal Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry)); + Assert.assertTrue("Actual tsKvEntrys is not containsAll Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry)); } - protected void processClientWithCorrectAccessTokenPushNodeMetricBuildArraysSimple() throws Exception { - clientWithCorrectNodeAccessTokenWithNDEATH();; + protected void processClientWithCorrectAccessTokenPushNodeMetricBuildArraysPrimitiveSimple() throws Exception { + clientWithCorrectNodeAccessTokenWithNDEATH(); String messageTypeName = SparkplugMessageType.NDATA.name(); List listKeys = new ArrayList<>(); @@ -184,7 +183,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac .setSeq(getSeqNum()); long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; - createdAddMetricValueArraysTsKv(listTsKvEntry, listKeys, ndataPayload, ts); + createdAddMetricValueArraysPrimitiveTsKv(listTsKvEntry, listKeys, ndataPayload, ts); if (client.isConnected()) { client.publish(NAMESPACE + "/" + groupId + "/" + messageTypeName + "/" + edgeNode, @@ -196,10 +195,9 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac .atMost(40, TimeUnit.SECONDS) .until(() -> { 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 equal Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry)); + Assert.assertTrue("Actual tsKvEntrys is not containsAll Expected tsKvEntrys", finalFuture.get().get().containsAll(listTsKvEntry)); } private void createdAddMetricValuePrimitiveTsKv(List listTsKvEntry, List listKeys, @@ -229,12 +227,8 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt16(), ts, UInt16)); listKeys.add(keys); - keys = "MyUInt32I"; - listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt32I(), ts, UInt32)); - listKeys.add(keys); - - keys = "MyUInt32L"; - listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt32L(), ts, UInt32)); + keys = "MyUInt32"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt32(), ts, UInt32)); listKeys.add(keys); keys = "MyUInt64"; @@ -271,62 +265,62 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac } - private void createdAddMetricValueArraysTsKv(List listTsKvEntry, List listKeys, - SparkplugBProto.Payload.Builder dataPayload, long ts) throws ThingsboardException { + private void createdAddMetricValueArraysPrimitiveTsKv(List listTsKvEntry, List listKeys, + SparkplugBProto.Payload.Builder dataPayload, long ts) throws ThingsboardException { String keys = "MyBytesArray"; - byte[] bytes = {nextInt8(), nextInt8(), nextInt8()}; + byte[] bytes = {nextInt8(), nextInt8(), nextInt8()}; createdAddMetricTsKvJson(dataPayload, keys, bytes, ts, Bytes, listTsKvEntry, listKeys); keys = "MyInt8Array"; - byte[] int8s = {nextInt8(), nextInt8(), nextInt8()}; + Byte[] int8s = {nextInt8(), nextInt8(), nextInt8()}; createdAddMetricTsKvJson(dataPayload, keys, int8s, ts, Int8Array, listTsKvEntry, listKeys); keys = "MyInt16Array"; - short[] int16s = {nextInt16(), nextInt16(), nextInt16()}; + Short[] int16s = {nextInt16(), nextInt16(), nextInt16()}; 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"; - short[] uInt8s = {nextUInt16(), nextUInt16(), nextUInt16()}; + Short[] uInt8s = {nextUInt8(), nextUInt8(), nextUInt8()}; 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"; - int[] uInt16s = {nextUInt16(), nextUInt16(), nextUInt16()}; + Integer[] uInt16s = {nextUInt16(), nextUInt16(), nextUInt16()}; 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"; - long[] uInt32Ls = {nextUInt32L(), nextUInt32L(), nextUInt32L()}; + Long[] uInt32Ls = {nextUInt32(), nextUInt32(), nextUInt32()}; 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"; - long[] uInt64s = {nextUInt64(), nextUInt64(), nextUInt64()}; + Long[] uInt64s = {nextUInt64(), nextUInt64(), nextUInt64()}; createdAddMetricTsKvJson(dataPayload, keys, uInt64s, ts, UInt64Array, listTsKvEntry, listKeys); 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); - keys = "MyDateTimeArray"; - long[] dateTimes = {nextDateTime(), nextDateTime(), nextDateTime()}; - createdAddMetricTsKvJson(dataPayload, keys, dateTimes, ts, DateTimeArray, listTsKvEntry, listKeys); - keys = "MyDoubleArray"; - double [] doubles = {nextDouble(), nextDouble(), nextDouble()}; + Double[] doubles = {nextDouble(), nextDouble(), nextDouble()}; createdAddMetricTsKvJson(dataPayload, keys, doubles, ts, DoubleArray, listTsKvEntry, listKeys); keys = "MyBooleanArray"; - boolean [] booleans = {nextBoolean(), nextBoolean(), nextBoolean()}; + Boolean[] booleans = {nextBoolean(), nextBoolean(), nextBoolean()}; createdAddMetricTsKvJson(dataPayload, keys, booleans, ts, BooleanArray, listTsKvEntry, listKeys); keys = "MyStringArray"; - String [] strings = {nexString(), nexString(), nexString()}; + String[] strings = {nexString(), nexString(), nexString()}; createdAddMetricTsKvJson(dataPayload, keys, strings, ts, StringArray, listTsKvEntry, listKeys); } @@ -337,18 +331,18 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac 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 { - var f = new BigDecimal(String.valueOf(value)); - TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new DoubleDataEntry(key, f.doubleValue())); + Double dd = Double.parseDouble(Float.toString(value)); + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new DoubleDataEntry(key, dd)); dataPayload.addMetrics(createMetric(value, ts, key, metricDataType)); return tsKvEntry; } private TsKvEntry createdAddMetricTsKvDouble(SparkplugBProto.Payload.Builder dataPayload, String key, double value, long ts, MetricDataType metricDataType) throws ThingsboardException { - var d = new BigDecimal(String.valueOf(value)); - TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(key, d.longValueExact())); + Long l = Double.valueOf(value).longValue(); + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(key, l)); dataPayload.addMetrics(createMetric(value, ts, key, metricDataType)); return tsKvEntry; } @@ -368,20 +362,24 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac } private void createdAddMetricTsKvJson(SparkplugBProto.Payload.Builder dataPayload, String key, - Object values, long ts, MetricDataType metricDataType, - List listTsKvEntry, - List listKeys) throws ThingsboardException { + Object values, long ts, MetricDataType metricDataType, + List listTsKvEntry, + List listKeys) throws ThingsboardException { ArrayNode nodeArray = newArrayNode(); switch (metricDataType) { case Bytes: + for (byte b : (byte[]) values) { + nodeArray.add(b); + } + break; case Int8Array: - for (byte b : (byte[])values) { + for (Byte b : (Byte[]) values) { nodeArray.add(b); } break; case Int16Array: case UInt8Array: - for (short b : (short[])values) { + for (Short b : (Short[]) values) { nodeArray.add(b); } break; @@ -391,33 +389,33 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac case UInt32Array: case UInt64Array: case DateTimeArray: - if (values instanceof int[]) { - for (int b : (int[])values) { + if (values instanceof Integer[]) { + for (Integer b : (Integer[]) values) { nodeArray.add(b); } } else { - for (long b : (long[])values) { + for (Long b : (Long[]) values) { nodeArray.add(b); } } break; case DoubleArray: - for (double b : (double[])values) { + for (Double b : (Double[]) values) { nodeArray.add(b); } break; case FloatArray: - for (float b : (float[])values) { + for (Float b : (Float[]) values) { nodeArray.add(b); } break; case BooleanArray: - for (boolean b : (boolean[])values) { + for (Boolean b : (Boolean[]) values) { nodeArray.add(b); } break; case StringArray: - for (String b : (String[])values) { + for (String b : (String[]) values) { nodeArray.add(b); } break; @@ -446,19 +444,15 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac return (short) random.nextInt(Short.MIN_VALUE, Short.MAX_VALUE); } - private short nextUInt16() { - return (short) random.nextInt(0, Short.MAX_VALUE * 2 + 1); + private int nextUInt16() { + return random.nextInt(0, Short.MAX_VALUE * 2 + 1); } private int nextInt32() { return random.nextInt(Integer.MIN_VALUE, Integer.MAX_VALUE); } - private int nextUInt32I() { - return random.nextInt(0, Integer.MAX_VALUE); - } - - private long nextUInt32L() { + private long nextUInt32() { long l = Integer.MAX_VALUE; return random.nextLong(0, l * 2 + 1); } 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 917d1a9f65..4ec4e9627a 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 @@ -54,8 +54,8 @@ public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSpa } @Test - public void testClientWithCorrectAccessTokenPushNodeMetricBuildPArraysSimple() throws Exception { - processClientWithCorrectAccessTokenPushNodeMetricBuildArraysSimple(); + public void testClientWithCorrectAccessTokenPushNodeMetricBuildPArraysPrimitiveSimple() throws Exception { + processClientWithCorrectAccessTokenPushNodeMetricBuildArraysPrimitiveSimple(); } } \ No newline at end of file diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 31037b075d..eb1c688950 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/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.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.parseTopicSubscribe; /** * @author Andrew Shvayka @@ -129,7 +128,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private final UUID sessionId; protected final MqttTransportContext context; - private final TransportService transportService; + public final TransportService transportService; private final SchedulerComponent scheduler; private final SslHandler sslHandler; private final ConcurrentMap mqttQoSMap; @@ -143,7 +142,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private final ConcurrentHashMap chunkSizes; private final ConcurrentMap rpcAwaitingAck; - private TopicType attrSubTopicType; + public TopicType attrSubTopicType; private TopicType rpcSubTopicType; private TopicType attrReqTopicType; private TopicType toServerRpcSubTopicType; @@ -451,62 +450,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private void handleSparkplugSubscribeMsg(List 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) { try { Matcher fwMatcher; @@ -768,7 +711,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement MqttQoS reqQoS = subscription.qualityOfService(); try { if (sparkplugSessionHandler != null) { - handleSparkplugSubscribeMsg(grantedQoSList, subscription, reqQoS); + sparkplugSessionHandler.handleSparkplugSubscribeMsg(grantedQoSList, subscription, reqQoS); activityReported = true; } else { switch (topic) { @@ -853,13 +796,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement registerSubQoS(topic, grantedQoSList, reqQoS); } - private void processAttributesSubscribe(List grantedQoSList, String topic, MqttQoS reqQoS, TopicType topicType) { + public void processAttributesSubscribe(List grantedQoSList, String topic, MqttQoS reqQoS, TopicType topicType) { transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null); attrSubTopicType = topicType; registerSubQoS(topic, grantedQoSList, reqQoS); } - private void registerSubQoS(String topic, List grantedQoSList, MqttQoS reqQoS) { + public void registerSubQoS(String topic, List grantedQoSList, MqttQoS reqQoS) { grantedQoSList.add(getMinSupportedQos(reqQoS)); mqttQoSMap.put(new MqttTopicMatcher(topic), getMinSupportedQos(reqQoS)); } @@ -1145,10 +1088,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement if (sparkplugSessionHandler == null) { SparkplugTopic sparkplugTopicNode = validatedSparkplugTopicConnectedNode(connectMessage); if (sparkplugTopicNode != null) { - SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); - sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId, sparkplugTopicNode); - sparkplugSessionHandler.onTelemetryProto(0, sparkplugBProtoNode, - deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopicNode); + SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); + sparkplugSessionHandler = new SparkplugNodeSessionHandler(this, deviceSessionCtx, sessionId, sparkplugTopicNode); + sparkplugSessionHandler.onTelemetryProto(0, sparkplugBProtoNode, + deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopicNode); } else { 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); @@ -1161,13 +1104,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private SparkplugTopic validatedSparkplugTopicConnectedNode (MqttConnectMessage connectMessage) throws ThingsboardException { - if(StringUtils.isNotBlank(connectMessage.payload().willTopic()) + private SparkplugTopic validatedSparkplugTopicConnectedNode(MqttConnectMessage connectMessage) throws ThingsboardException { + if (StringUtils.isNotBlank(connectMessage.payload().willTopic()) && connectMessage.payload().willMessageInBytes() != null && connectMessage.payload().willMessageInBytes().length > 0) { - SparkplugTopic sparkplugTopicNode = parseTopicPublish(connectMessage.payload().willTopic()); - if(NDEATH.equals(sparkplugTopicNode.getType())){ - return sparkplugTopicNode; + SparkplugTopic sparkplugTopicNode = parseTopicPublish(connectMessage.payload().willTopic()); + if (NDEATH.equals(sparkplugTopicNode.getType())) { + return sparkplugTopicNode; } } return null; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java index 4a379e8a85..4b371acbbc 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java +++ b/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.protobuf.Descriptors; 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 org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; 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.gen.transport.TransportProtos; 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.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.fromSparkplugBMetricToKeyValueProto; 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 @@ -57,10 +63,12 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { private final SparkplugTopic sparkplugTopicNode; private final Map nodeBirthMetrics; + private final MqttTransportHandler parent; - public SparkplugNodeSessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId, + public SparkplugNodeSessionHandler(MqttTransportHandler parent, DeviceSessionCtx deviceSessionCtx, UUID sessionId, SparkplugTopic sparkplugTopicNode) { super(deviceSessionCtx, sessionId); + this.parent = parent; this.sparkplugTopicNode = sparkplugTopicNode; this.nodeBirthMetrics = new ConcurrentHashMap<>(); } @@ -129,6 +137,34 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { } } + public void handleSparkplugSubscribeMsg(List 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 grantedQoSList, String topic, + MqttQoS reqQoS, TopicType topicType, String deviceName) + throws AdaptorException, ThingsboardException, ExecutionException, InterruptedException { + checkDeviceName(deviceName); + ListenableFuture 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 { try { processOnDisconnect(mqttMsg, deviceName); 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 0b36e1fe90..f1796b2643 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 @@ -16,11 +16,14 @@ package org.thingsboard.server.transport.mqtt.util.sparkplug; 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.node.ArrayNode; import com.fasterxml.jackson.databind.ser.std.FileSerializer; import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; @@ -42,6 +45,7 @@ import java.nio.LongBuffer; import java.nio.ShortBuffer; import java.text.NumberFormat; import java.util.Arrays; +import java.util.List; import java.util.Optional; import static org.thingsboard.common.util.JacksonUtil.newArrayNode; @@ -228,34 +232,39 @@ public class SparkplugMetricUtil { .setDatatype(metricDataType.toIntValue()) .build(); switch (metricDataType) { - case Int8: - case Int16: + case Int8: // (byte) + return metric.toBuilder().setIntValue(((Byte) value).intValue()).build(); + case Int16: // (short) case UInt8: - case UInt16: + return metric.toBuilder().setIntValue(((Short) value).intValue()).build(); + case UInt16: // (int) case Int32: return metric.toBuilder().setIntValue(((Integer) value).intValue()).build(); - case UInt32: + case UInt32: // (long) case Int64: case UInt64: case DateTime: return metric.toBuilder().setLongValue(((Long) value).longValue()).build(); - case Float: + case Float: // (float) return metric.toBuilder().setFloatValue(((Float) value).floatValue()).build(); - case Double: + case Double: // (double) return metric.toBuilder().setDoubleValue(((Double) value).doubleValue()).build(); - case Boolean: + case Boolean: // (boolean) return metric.toBuilder().setBooleanValue(((Boolean) value).booleanValue()).build(); - case String: + case String: // String) case Text: case UUID: return metric.toBuilder().setStringValue((String) value).build(); case Bytes: - case Int8Array: ByteString byteString = ByteString.copyFrom((byte[]) value); 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 UInt8Array: - byte[] int16Array = shortArrayToByteArray((short[]) value); + byte[] int16Array = shortArrayToByteArray((Short[]) value); ByteString byteInt16Array = ByteString.copyFrom((int16Array)); return metric.toBuilder().setBytesValue(byteInt16Array).build(); case Int32Array: @@ -264,25 +273,25 @@ public class SparkplugMetricUtil { case UInt32Array: case UInt64Array: case DateTimeArray: - if (value instanceof int[]) { - byte[] int32Array = integerArrayToByteArray((int[]) value); + if (value instanceof Integer[]) { + byte[] int32Array = integerArrayToByteArray((Integer[]) value); ByteString byteInt32Array = ByteString.copyFrom((int32Array)); return metric.toBuilder().setBytesValue(byteInt32Array).build(); } else { - byte[] int64Array = longArrayToByteArray((long[]) value); + byte[] int64Array = longArrayToByteArray((Long[]) value); ByteString byteInt64Array = ByteString.copyFrom((int64Array)); return metric.toBuilder().setBytesValue(byteInt64Array).build(); } case DoubleArray: - byte[] doubleArray = doublArrayToByteArray((double[]) value); + byte[] doubleArray = doublArrayToByteArray((Double[]) value); ByteString byteDoubleArray = ByteString.copyFrom(doubleArray); return metric.toBuilder().setBytesValue(byteDoubleArray).build(); case FloatArray: - byte[] floatArray = floatArrayToByteArray((float[]) value); + byte[] floatArray = floatArrayToByteArray((Float[]) value); ByteString byteFloatArray = ByteString.copyFrom(floatArray); return metric.toBuilder().setBytesValue(byteFloatArray).build(); case BooleanArray: - byte[] booleanArray = booleanArrayToByteArray((boolean[]) value); + byte[] booleanArray = booleanArrayToByteArray((Boolean[]) value); ByteString byteBooleanArray = ByteString.copyFrom(booleanArray); return metric.toBuilder().setBytesValue(byteBooleanArray).build(); case StringArray: @@ -307,10 +316,14 @@ public class SparkplugMetricUtil { if (kv.getTypeValue() <= 3) { return validatedValuePrimitiveByTypeMetric(kv, metricDataType); } 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 { throw new ThingsboardException("Invalid type KeyValueProto " + kv.toString() + " for MetricDataType " + metricDataType.name(), ThingsboardErrorCode.INVALID_ARGUMENTS); } + return Optional.empty(); } public static Optional validatedValuePrimitiveByTypeMetric(TransportProtos.KeyValueProto kv, MetricDataType metricDataType) throws ThingsboardException { @@ -392,188 +405,115 @@ public class SparkplugMetricUtil { return Optional.empty(); } - - public static Optional validatedValueJsonByTypeMetric(TransportProtos.KeyValueProto kv, MetricDataType metricDataType) { -// try { -// Optional valueOpt; -// switch (metricDataType) { -// // int -// case Int8: -// case Int16: -// case UInt8: -// case UInt16: -// valueOpt = getValueKvProtoPrimitive(kv); -// return valueOpt.isPresent() ? Optional.of(Integer.valueOf(String.valueOf(valueOpt.get()))) : valueOpt; -// // int/long -// case Int32: -// case UInt32: -// valueOpt = getValueKvProtoPrimitive(kv); -// try { -// return Optional.of(Integer.valueOf(String.valueOf(valueOpt.get()))); -// } catch (NumberFormatException e) { -// return Optional.of(Long.valueOf(String.valueOf(valueOpt.get()))); -// } -// // long -// case Int64: -// case UInt64: -// case DateTime: -// valueOpt = getValueKvProtoPrimitive(kv); -// return Optional.of(Long.valueOf(String.valueOf(valueOpt.get()))); -// // float -// case Float: -// valueOpt = getValueKvProtoPrimitive(kv); -// var f = new BigDecimal(String.valueOf(kv.getDoubleV())); -// return Optional.of(f.floatValue()); -// } -// break; -// // double -// case Double: -// if (kv.getTypeValue() == 1) { -// return Optional.of(kv.getLongV()); -// } -// break; -// case Boolean: -// if (kv.getTypeValue() == 0) { // ok 0 -// return Optional.of(kv.getBoolV()); -// } -// break; -// case String: -// case Text: -// case UUID: -// if (kv.getTypeValue() == 4) { -// return Optional.of(kv.getStringV()); -// } -// break; -// // byte[] -// case Bytes: -// case Int8Array: -// if (kv.getTypeValue() == 5) { -//// ByteString byteString = ByteString.copyFrom((byte[]) value); -//// return metric.toBuilder().setBytesValue(byteString).build(); -// return Optional.of(kv.getJsonV()); -// } -// break; -// // short[] -// 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(); + public static Optional validatedValueJsonByTypeMetric(String arrayNodeStr, MetricDataType metricDataType) { + try { + Optional valueOpt; + switch (metricDataType) { + // byte[] + case Bytes: + List listBytes = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + byte[] bytes = new byte[listBytes.size()]; + for(int i = 0; i < listBytes.size(); i++) { + bytes[i] = listBytes.get(i).byteValue(); + } + return Optional.of(bytes); + // Byte [] + case Int8Array: + List listInt8Array = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + Byte[] int8Arrays = listInt8Array.toArray(Byte[]::new) ; + return Optional.of(int8Arrays); + // Short[] + case Int16Array: + case UInt8Array: + List listShorts = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + return Optional.of(listShorts.toArray(Short[]::new)); + // Integer [] + case UInt16Array: + case Int32Array: + List listIntegers = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + return Optional.of(listIntegers.toArray(Integer[]::new)); + // Long[] + case UInt32Array: + case Int64Array: + case UInt64Array: + case DateTimeArray: + List listLongs = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + return Optional.of(listLongs.toArray(Long[]::new)); + // Double [] + case DoubleArray: + List listDoubles = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + return Optional.of(listDoubles.toArray(Double[]::new)); + // Float[] + case FloatArray: + List listFloats = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + return Optional.of(listFloats.toArray(Float[]::new)); + // Boolean[] + case BooleanArray: + List listBooleans = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + return Optional.of(listBooleans.toArray(Boolean[]::new)); + case StringArray: + List listStrings = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + return Optional.of(listStrings.toArray(String[]::new)); + case DataSet: + case File: + case Template: + log.error("Invalid type value [{}] for MetricDataType [{}]", arrayNodeStr, metricDataType.name()); + return Optional.empty(); + case Unknown: + default: + log.error("Invalid MetricDataType [{}] type, value [{}]", arrayNodeStr, metricDataType.name()); + return Optional.empty(); + } + } catch (Exception e) { + log.error("Invalid type value [{}] for MetricDataType [{}] [{}]", arrayNodeStr, metricDataType.name(), e.getMessage()); + 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); - for (short d : inputs) { + for (Short d : inputs) { bb.putShort(d); } return bb.array(); } - private static byte[] integerArrayToByteArray(int[] inputs) { + private static byte[] integerArrayToByteArray(Integer[] inputs) { ByteBuffer bb = ByteBuffer.allocate(inputs.length * 4); - for (int d : inputs) { + for (Integer d : inputs) { bb.putInt(d); } return bb.array(); } - private static byte[] longArrayToByteArray(long[] inputs) { + private static byte[] longArrayToByteArray(Long[] inputs) { ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); - for (long d : inputs) { + for (Long d : inputs) { bb.putLong(d); } return bb.array(); } - private static byte[] doublArrayToByteArray(double[] inputs) { + private static byte[] doublArrayToByteArray(Double[] inputs) { ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); - for (double d : inputs) { + for (Double d : inputs) { bb.putDouble(d); } return bb.array(); } - private static byte[] floatArrayToByteArray(float[] inputs) throws ThingsboardException { + private static byte[] floatArrayToByteArray(Float[] inputs) throws ThingsboardException { ByteArrayOutputStream bas = new ByteArrayOutputStream(); DataOutputStream ds = new DataOutputStream(bas); - for (float f : inputs) { + for (Float f : inputs) { try { ds.writeFloat(f); } catch (IOException e) { @@ -583,7 +523,7 @@ public class SparkplugMetricUtil { return bas.toByteArray(); } - private static byte[] booleanArrayToByteArray(boolean[] inputs) { + private static byte[] booleanArrayToByteArray(Boolean[] inputs) { byte[] toReturn = new byte[inputs.length]; for (int entry = 0; entry < toReturn.length; entry++) { toReturn[entry] = (byte) (inputs[entry] ? 1 : 0);