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 615a255ac2..b0b936f05f 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 @@ -41,28 +41,15 @@ import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; import static org.thingsboard.common.util.JacksonUtil.newArrayNode; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.BooleanArray; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Bytes; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.DateTimeArray; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.DoubleArray; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.FloatArray; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16Array; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32Array; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64Array; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int8; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int8Array; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.StringArray; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt16; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt16Array; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt32; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt32Array; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt64; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt64Array; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt8; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt8Array; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; /** @@ -270,58 +257,6 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac String keys = "MyBytesArray"; byte[] bytes = {nextInt8(), nextInt8(), nextInt8()}; createdAddMetricTsKvJson(dataPayload, keys, bytes, ts, Bytes, listTsKvEntry, listKeys); - - keys = "MyInt8Array"; - Byte[] int8s = {nextInt8(), nextInt8(), nextInt8()}; - createdAddMetricTsKvJson(dataPayload, keys, int8s, ts, Int8Array, listTsKvEntry, listKeys); - - keys = "MyInt16Array"; - Short[] int16s = {nextInt16(), nextInt16(), nextInt16()}; - createdAddMetricTsKvJson(dataPayload, keys, int16s, ts, Int16Array, listTsKvEntry, listKeys); - - keys = "MyUInt8Array"; - 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"; - 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 = {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()}; - createdAddMetricTsKvJson(dataPayload, keys, uInt64s, ts, UInt64Array, listTsKvEntry, listKeys); - - keys = "MyFloatArray"; - Float[] floats = {nextFloat(0, 300), nextFloat(0, 4000), nextFloat(10, 10000)}; - createdAddMetricTsKvJson(dataPayload, keys, floats, ts, FloatArray, listTsKvEntry, listKeys); - - keys = "MyDoubleArray"; - Double[] doubles = {nextDouble(), nextDouble(), nextDouble()}; - createdAddMetricTsKvJson(dataPayload, keys, doubles, ts, DoubleArray, listTsKvEntry, listKeys); - - keys = "MyBooleanArray"; - Boolean[] booleans = {nextBoolean(), nextBoolean(), nextBoolean()}; - createdAddMetricTsKvJson(dataPayload, keys, booleans, ts, BooleanArray, listTsKvEntry, listKeys); - - keys = "MyStringArray"; - String[] strings = {nexString(), nexString(), nexString()}; - createdAddMetricTsKvJson(dataPayload, keys, strings, ts, StringArray, listTsKvEntry, listKeys); } private TsKvEntry createdAddMetricTsKvLong(SparkplugBProto.Payload.Builder dataPayload, String key, Object value, @@ -372,54 +307,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac nodeArray.add(b); } break; - case Int8Array: - for (Byte b : (Byte[]) values) { - nodeArray.add(b); - } - break; - case Int16Array: - case UInt8Array: - for (Short b : (Short[]) values) { - nodeArray.add(b); - } - break; - case Int32Array: - case UInt16Array: - case Int64Array: - case UInt32Array: - case UInt64Array: - case DateTimeArray: - if (values instanceof Integer[]) { - for (Integer b : (Integer[]) values) { - nodeArray.add(b); - } - } else { - for (Long b : (Long[]) values) { - nodeArray.add(b); - } - } - break; - case DoubleArray: - for (Double b : (Double[]) values) { - nodeArray.add(b); - } - break; - case FloatArray: - for (Float b : (Float[]) values) { - nodeArray.add(b); - } - break; - case BooleanArray: - for (Boolean b : (Boolean[]) values) { - nodeArray.add(b); - } - break; - case StringArray: - for (String b : (String[]) values) { - nodeArray.add(b); - } - break; - default: + default: throw new IllegalStateException("Unexpected value: " + metricDataType); } if (nodeArray.size() > 0) { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java index f4dc74f46b..ceda342acb 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java @@ -54,21 +54,6 @@ public enum MetricDataType { // PropertyValue Types (20 and 21) are NOT metric datatypes - // Array Types - Int8Array(22, Byte[].class), - Int16Array(23, Short[].class), - Int32Array(24, Integer[].class), - Int64Array(25, Long[].class), - UInt8Array(26, Short[].class), - UInt16Array(27, Integer[].class), - UInt32Array(28, Long[].class), - UInt64Array(29, BigInteger[].class), - FloatArray(30, Float[].class), - DoubleArray(31, Double[].class), - BooleanArray(32, Boolean[].class), - StringArray(33, String[].class), - DateTimeArray(34, Date[].class), - // Unknown Unknown(0, Object.class); @@ -155,32 +140,6 @@ public enum MetricDataType { return File; case 19: return Template; - case 22: - return Int8Array; - case 23: - return Int16Array; - case 24: - return Int32Array; - case 25: - return Int64Array; - case 26: - return UInt8Array; - case 27: - return UInt16Array; - case 28: - return UInt32Array; - case 29: - return UInt64Array; - case 30: - return FloatArray; - case 31: - return DoubleArray; - case 32: - return BooleanArray; - case 33: - return StringArray; - case 34: - return DateTimeArray; default: return Unknown; } 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 f1796b2643..19ec3156e6 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 @@ -30,19 +30,8 @@ import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; -import java.io.ByteArrayInputStream; -import java.io.ByteArrayOutputStream; -import java.io.DataOutputStream; -import java.io.IOException; -import java.io.ObjectInputStream; -import java.io.ObjectOutputStream; import java.math.BigDecimal; import java.nio.ByteBuffer; -import java.nio.DoubleBuffer; -import java.nio.FloatBuffer; -import java.nio.IntBuffer; -import java.nio.LongBuffer; -import java.nio.ShortBuffer; import java.text.NumberFormat; import java.util.Arrays; import java.util.List; @@ -107,85 +96,13 @@ public class SparkplugMetricUtil { return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.STRING_V) .setStringV(protoMetric.getStringValue()).build()); // byte[] - case BooleanArray: - ByteBuffer booleanByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()); - while (booleanByteBuffer.hasRemaining()) { - nodeArray.add(booleanByteBuffer.get() == (byte) 0 ? false : true); - } - return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) - .setJsonV(nodeArray.toString()).build()); - // byte[] case Bytes: - case Int8Array: ByteBuffer byteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()); while (byteBuffer.hasRemaining()) { nodeArray.add(byteBuffer.get()); } return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) .setJsonV(nodeArray.toString()).build()); - // short[] - case Int16Array: - case UInt8Array: - ShortBuffer shortByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()) - .asShortBuffer(); - while (shortByteBuffer.hasRemaining()) { - nodeArray.add(shortByteBuffer.get()); - } - return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) - .setJsonV(nodeArray.toString()).build()); - // int[] - case Int32Array: - case UInt16Array: - IntBuffer intByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()) - .asIntBuffer(); - while (intByteBuffer.hasRemaining()) { - nodeArray.add(intByteBuffer.get()); - } - return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) - .setJsonV(nodeArray.toString()).build()); - // float[] - case FloatArray: - FloatBuffer floatByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()) - .asFloatBuffer(); - while (floatByteBuffer.hasRemaining()) { - nodeArray.add(floatByteBuffer.get()); - } - return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) - .setJsonV(nodeArray.toString()).build()); - // double[] - case DoubleArray: - DoubleBuffer doubleByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()) - .asDoubleBuffer(); - while (doubleByteBuffer.hasRemaining()) { - nodeArray.add(doubleByteBuffer.get()); - } - return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) - .setJsonV(nodeArray.toString()).build()); - // long[] - case DateTimeArray: - case Int64Array: - case UInt64Array: - case UInt32Array: - LongBuffer longByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()) - .asLongBuffer(); - while (longByteBuffer.hasRemaining()) { - nodeArray.add(longByteBuffer.get()); - } - return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) - .setJsonV(nodeArray.toString()).build()); - case StringArray: - ByteBuffer stringByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()); - final ByteArrayInputStream byteArrayInputStream = - new ByteArrayInputStream(stringByteBuffer.array()); - final ObjectInputStream objectInputStream = - new ObjectInputStream(byteArrayInputStream); - final String[] stringArray = (String[]) objectInputStream.readObject(); - objectInputStream.close(); - for (String s : stringArray) { - nodeArray.add(s); - } - return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.JSON_V) - .setJsonV(nodeArray.toString()).build()); case DataSet: case Template: case File: @@ -258,46 +175,6 @@ public class SparkplugMetricUtil { case Bytes: 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); - ByteString byteInt16Array = ByteString.copyFrom((int16Array)); - return metric.toBuilder().setBytesValue(byteInt16Array).build(); - case Int32Array: - case UInt16Array: - case Int64Array: - case UInt32Array: - case UInt64Array: - case DateTimeArray: - 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); - ByteString byteInt64Array = ByteString.copyFrom((int64Array)); - return metric.toBuilder().setBytesValue(byteInt64Array).build(); - } - case DoubleArray: - byte[] doubleArray = doublArrayToByteArray((Double[]) value); - ByteString byteDoubleArray = ByteString.copyFrom(doubleArray); - return metric.toBuilder().setBytesValue(byteDoubleArray).build(); - case FloatArray: - byte[] floatArray = floatArrayToByteArray((Float[]) value); - ByteString byteFloatArray = ByteString.copyFrom(floatArray); - return metric.toBuilder().setBytesValue(byteFloatArray).build(); - case BooleanArray: - byte[] booleanArray = booleanArrayToByteArray((Boolean[]) value); - ByteString byteBooleanArray = ByteString.copyFrom(booleanArray); - return metric.toBuilder().setBytesValue(byteBooleanArray).build(); - case StringArray: - byte[] stringArray = stringArrayToByteArray((String[]) value); - ByteString byteStringArray = ByteString.copyFrom(stringArray); - return metric.toBuilder().setBytesValue(byteStringArray).build(); case DataSet: return metric.toBuilder().setDatasetValue((SparkplugBProto.Payload.DataSet) value).build(); case File: @@ -417,43 +294,6 @@ public class SparkplugMetricUtil { 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: @@ -470,80 +310,6 @@ public class SparkplugMetricUtil { } } - 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) { - ByteBuffer bb = ByteBuffer.allocate(inputs.length * 2); - for (Short d : inputs) { - bb.putShort(d); - } - return bb.array(); - } - - private static byte[] integerArrayToByteArray(Integer[] inputs) { - ByteBuffer bb = ByteBuffer.allocate(inputs.length * 4); - for (Integer d : inputs) { - bb.putInt(d); - } - return bb.array(); - } - - private static byte[] longArrayToByteArray(Long[] inputs) { - ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); - for (Long d : inputs) { - bb.putLong(d); - } - return bb.array(); - } - - private static byte[] doublArrayToByteArray(Double[] inputs) { - ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); - for (Double d : inputs) { - bb.putDouble(d); - } - return bb.array(); - } - - private static byte[] floatArrayToByteArray(Float[] inputs) throws ThingsboardException { - ByteArrayOutputStream bas = new ByteArrayOutputStream(); - DataOutputStream ds = new DataOutputStream(bas); - for (Float f : inputs) { - try { - ds.writeFloat(f); - } catch (IOException e) { - throw new ThingsboardException("Invalid value float ", ThingsboardErrorCode.INVALID_ARGUMENTS); - } - } - return bas.toByteArray(); - } - - 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); - } - return toReturn; - } - - private static byte[] stringArrayToByteArray(String[] inputs) throws ThingsboardException { - final ByteArrayOutputStream bas = new ByteArrayOutputStream(); - try { - final ObjectOutputStream os = new ObjectOutputStream(bas); - os.writeObject(inputs); - os.flush(); - os.close(); - } catch (Exception e) { - throw new ThingsboardException("Invalid value float ", ThingsboardErrorCode.INVALID_ARGUMENTS); - } - return bas.toByteArray(); - } - private static Optional getValueKvProtoPrimitive(TransportProtos.KeyValueProto kv) { if (kv.getTypeValue() == 0) { // boolean return Optional.of(String.valueOf(kv.getBoolV()));