diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index 9ea269918e..3d7dc62cba 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -83,6 +83,7 @@ import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.config.ThingsboardSecurityConfiguration; import org.thingsboard.server.dao.Dao; import org.thingsboard.server.dao.tenant.TenantProfileService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.service.mail.TestMailService; import org.thingsboard.server.service.security.auth.jwt.RefreshTokenRequest; import org.thingsboard.server.service.security.auth.rest.LoginRequest; @@ -167,6 +168,9 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { @Autowired private TenantProfileService tenantProfileService; + @Autowired + public TimeseriesService tsService; + @Rule public TestRule watcher = new TestWatcher() { protected void starting(Description description) { diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java index 917e41d1b4..28e458d971 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java @@ -103,6 +103,7 @@ public abstract class AbstractMqttIntegrationTest extends AbstractTransportInteg if (StringUtils.hasLength(config.getAttributesTopicFilter())) { mqttDeviceProfileTransportConfiguration.setDeviceAttributesTopic(config.getAttributesTopicFilter()); } + mqttDeviceProfileTransportConfiguration.setSparkPlug(config.isSparkPlug()); mqttDeviceProfileTransportConfiguration.setSendAckOnValidationException(config.isSendAckOnValidationException()); TransportPayloadTypeConfiguration transportPayloadTypeConfiguration; if (TransportPayloadType.JSON.equals(transportPayloadType)) { diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java index bc535b424b..1a2b6eefa4 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java @@ -26,6 +26,7 @@ public class MqttTestConfigProperties { String deviceName; String gatewayName; + boolean isSparkPlug; TransportPayloadType transportPayloadType; diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java index 0671dd6f95..0af79c5f06 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java @@ -89,7 +89,9 @@ public class MqttV5TestClient { // We should copy part of MqttV3TestClient, due if (client == null) { throw new RuntimeException("Failed to connect! MqttAsyncClient is not initialized!"); } - return client.connect(options); + IMqttToken connect = client.connect(options); + connect.waitForCompletion(TIMEOUT_MS); + return connect; } public void disconnectAndWait() throws MqttException { 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 new file mode 100644 index 0000000000..bef5bf4770 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java @@ -0,0 +1,256 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.sparkplug; + +import com.google.protobuf.ByteString; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.paho.mqttv5.client.IMqttToken; +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.common.data.exception.ThingsboardErrorCode; +import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; +import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; +import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties; +import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; +import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil; + +import java.io.ByteArrayOutputStream; +import java.io.DataOutputStream; +import java.io.IOException; +import java.io.ObjectOutputStream; +import java.nio.ByteBuffer; +import java.util.Calendar; + +import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; + +/** + * Created by nickAS21 on 12.01.23 + */ +@Slf4j +public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttIntegrationTest { + + protected MqttV5TestClient client; + protected Calendar calendar = Calendar.getInstance(); + + protected static final String NAMESPACE = "spBv1.0"; + protected static final String groupId = "SparkplugBGroupId"; + protected static final String edgeNode = "SparkpluBNode"; + protected static final String keysBdSeq = "bdSeq"; + protected static final String alias = "Failed Post Telemetry node proto payload. SparkplugMessageType "; + protected String deviceId = "Test Sparkplug B Device"; + protected int bdSeq = 0; + 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") + .isSparkPlug(true) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + } + + public void processClientWithCorrectNodeAccess() throws Exception { + this.client = new MqttV5TestClient(); + MqttWireMessage response = clientWithCorrectNodeAccessToken(client); + Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType()); + MqttConnAck connAckMsg = (MqttConnAck) response; + Assert.assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode()); + } + + protected SparkplugBProto.Payload.Metric createMetric(Object value, TsKvEntry tsKvEntry, MetricDataType metricDataType) throws ThingsboardException { + SparkplugBProto.Payload.Metric metric = SparkplugBProto.Payload.Metric.newBuilder() + .setTimestamp(tsKvEntry.getTs()) + .setName(tsKvEntry.getKey()) + .setDatatype(metricDataType.toIntValue()) + .build(); + switch (metricDataType) { + case Int8: + case Int16: + case UInt8: + case UInt16: + int valueMetric = Integer.valueOf(String.valueOf(value)); + return metric.toBuilder().setIntValue(valueMetric).build(); + case Int32: + case UInt32: + if (value instanceof Long) { + return metric.toBuilder().setLongValue((long) value).build(); + } else { + return metric.toBuilder().setIntValue((int)value).build(); + } + case Int64: + case UInt64: + case DateTime: + return metric.toBuilder().setLongValue((long) value).build(); + case Float: + return metric.toBuilder().setFloatValue((float) value).build(); + case Double: + return metric.toBuilder().setDoubleValue((double) value).build(); + case Boolean: + return metric.toBuilder().setBooleanValue((boolean) value).build(); + case String: + case Text: + case UUID: + return metric.toBuilder().setStringValue((String) value).build(); + case DataSet: + return metric.toBuilder().setDatasetValue((SparkplugBProto.Payload.DataSet) value).build(); + case Bytes: + case Int8Array: + ByteString byteString = ByteString.copyFrom((byte[]) value); + return metric.toBuilder().setBytesValue(byteString).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 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(); + } + 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 File: + SparkplugMetricUtil.File file = (SparkplugMetricUtil.File) value; + ByteString byteFileString = ByteString.copyFrom(file.getBytes()); + return metric.toBuilder().setBytesValue(byteFileString).build(); + case Template: + return metric.toBuilder().setTemplateValue((SparkplugBProto.Payload.Template) value).build(); + case Unknown: + throw new ThingsboardException("Invalid value for MetricDataType " + metricDataType.name(), ThingsboardErrorCode.INVALID_ARGUMENTS); + } + return metric; + } + + private byte[] shortArrayToByteArray(short[] inputs) { + ByteBuffer bb = ByteBuffer.allocate(inputs.length * 2); + for (short d : inputs) { + bb.putShort(d); + } + return bb.array(); + } + + private byte[] integerArrayToByteArray(int[] inputs) { + ByteBuffer bb = ByteBuffer.allocate(inputs.length * 4); + for (int d : inputs) { + bb.putInt(d); + } + return bb.array(); + } + + private byte[] longArrayToByteArray(long[] inputs) { + ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); + for (long d : inputs) { + bb.putLong(d); + } + return bb.array(); + } + + private byte[] doublArrayToByteArray(double[] inputs) { + ByteBuffer bb = ByteBuffer.allocate(inputs.length * 8); + for (double d : inputs) { + bb.putDouble(d); + } + return bb.array(); + } + + private 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 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 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 MqttWireMessage clientWithCorrectNodeAccessToken(MqttV5TestClient client) throws Exception { + IMqttToken connectionResult = client.connectAndWait(gatewayAccessToken); + return connectionResult.getResponse(); + } + + protected long getBdSeqNum() throws Exception { + if (bdSeq == 256) { + bdSeq = 0; + } + return bdSeq++; + } + + protected long getSeqNum() throws Exception { + if (seq == 256) { + seq = 0; + } + return seq++; + } + +} 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 new file mode 100644 index 0000000000..c0f9bf6797 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java @@ -0,0 +1,140 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +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.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.Device; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; +import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; +import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSparkplugTest; +import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; + +import java.util.Optional; +import java.util.Set; +import java.util.HashSet; +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.MetricDataType.Int64; + +/** + * Created by nickAS21 on 12.01.23 + */ +@Slf4j +public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends AbstractMqttV5ClientSparkplugTest { + + protected void processClientWithCorrectNodeAccessTokenTest() throws Exception { + processClientWithCorrectNodeAccess(); + } + + protected void processClientWithCorrectNodeAccessTokenWithNdeathTest() throws Exception { + long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; + long value = bdSeq = 0; + MetricDataType metricDataType = Int64; + TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, value)); + + SparkplugBProto.Payload.Builder deathPayload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()); + deathPayload.addMetrics(createMetric(value, tsKvEntryBdSecOriginal, metricDataType)); + + MqttWireMessage response = clientWithCorrectNodeAccessTokenWithNDEATH(deathPayload.build().toByteArray()); + + 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; + TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, value)); + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NDEATH.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), keys)); + return finalFuture.get().get().isPresent(); + }); + TsKvEntry actualTsKvEntry = finalFuture.get().get().get(); + Assert.assertEquals(expectedTsKvEntry, actualTsKvEntry); + } + + protected void processClientWithCorrectAccessTokenCreatedDevices(int cntDevices) throws Exception { + processClientWithCorrectNodeAccess(); + long ts = calendar.getTimeInMillis(); + MetricDataType metricDataType = Int32; + Set deviceIds = new HashSet<>(); + String keys = "Device Metric int32"; + int valueDeviceInt32 = 1024; + TsKvEntry expectedTsKvEntryDeviceInt32 = new BasicTsKvEntry(ts, new LongDataEntry(keys, Integer.toUnsignedLong(valueDeviceInt32))); + SparkplugBProto.Payload.Metric metric = createMetric(valueDeviceInt32, expectedTsKvEntryDeviceInt32, metricDataType); + for (int i=0; i < cntDevices; i++ ) { + SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()) + .setSeq(getSeqNum()); + String deviceName = deviceId + "_" + i; + + payloadBirthDevice.addMetrics(metric); + if (client.isConnected()) { + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceName, + payloadBirthDevice.build().toByteArray(), 0, false); + deviceIds.add(deviceName); + } + } + + Assert.assertEquals(cntDevices, deviceIds.size()); + + for (String deviceName: deviceIds) { + AtomicReference device = new AtomicReference<>(); + await(alias + "find device [" + deviceName + "] after crete") + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); + return device.get() != null; + }); + } + } + + private MqttWireMessage clientWithCorrectNodeAccessTokenWithNDEATH(byte[] deathBytes) throws Exception { + this.client = new MqttV5TestClient(); + MqttConnectionOptions options = new MqttConnectionOptions(); + options.setUserName(gatewayAccessToken); + if (deathBytes != null) { + String topic = NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.NDEATH.name() + "/" + edgeNode; + MqttMessage msg = new MqttMessage(); + msg.setId(0); + msg.setPayload(deathBytes); + options.setWill(topic, msg); + } + IMqttToken connectionResult = client.connect(options); + return connectionResult.getResponse(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java new file mode 100644 index 0000000000..85f708d0f8 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java @@ -0,0 +1,62 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.sparkplug.connection; + +import org.eclipse.paho.mqttv5.common.MqttException; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.thingsboard.server.dao.service.DaoSqlTest; + +/** + * Created by nickAS21 on 12.01.23 + */ +@DaoSqlTest +public class MqttV5ClientSparkplugBConnectionTest extends AbstractMqttV5ClientSparkplugConnectionTest { + + @Before + public void beforeTest() throws Exception { + beforeSparkplugTest(); + } + + @After + public void afterTest() throws MqttException { + if (client.isConnected()) { + client.disconnect(); + } + } + + @Test + public void testClientWithCorrectAccessToken() throws Exception { + processClientWithCorrectNodeAccessTokenTest(); + } + + @Test + public void testClientWithCorrectAccessTokenWithNDEATH() throws Exception { + processClientWithCorrectNodeAccessTokenWithNdeathTest(); + } + + @Test + public void testClientWithCorrectAccessTokenCreatedOneDevice() throws Exception { + processClientWithCorrectAccessTokenCreatedDevices(1); + } + + @Test + public void testClientWithCorrectAccessTokenCreatedTwoDevice() throws Exception { + processClientWithCorrectAccessTokenCreatedDevices(2); + } + +} 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 new file mode 100644 index 0000000000..e2892aaa3a --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java @@ -0,0 +1,505 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.sparkplug.timeseries; + +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; +import org.junit.Assert; +import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; +import org.thingsboard.server.common.data.kv.JsonDataEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; +import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSparkplugTest; +import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; + +import java.math.BigDecimal; +import java.util.ArrayList; +import java.util.List; +import java.util.Optional; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.TimeUnit; +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; + +/** + * Created by nickAS21 on 12.01.23 + */ +@Slf4j +public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends AbstractMqttV5ClientSparkplugTest { + + protected ThreadLocalRandom random = ThreadLocalRandom.current(); + + protected void processClientWithCorrectAccessTokenPublishNBIRTH() throws Exception { + processClientWithCorrectNodeAccess(); + List listKeys = new ArrayList<>(); + SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()); + long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; + long valueBdSec = getBdSeqNum(); + MetricDataType metricDataType = Int64; + TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, valueBdSec)); + payloadBirthNode.addMetrics(createMetric(valueBdSec, tsKvEntryBdSecOriginal, metricDataType)); + listKeys.add(SparkplugMessageType.NBIRTH.name() + " " + keysBdSeq); + + String keys = "Node Control/Rebirth"; + boolean valueRebirth = false; + metricDataType = MetricDataType.Boolean; + TsKvEntry expectedSsKvEntryRebirth = new BasicTsKvEntry(ts, new BooleanDataEntry(keys, valueRebirth)); + payloadBirthNode.addMetrics(createMetric(valueRebirth, expectedSsKvEntryRebirth, metricDataType)); + listKeys.add(keys); + + keys = "Node Metric int32"; + int valueNodeInt32 = 1024; + metricDataType = Int32; + TsKvEntry expectedSsKvEntryNodeInt32 = new BasicTsKvEntry(ts, new LongDataEntry(keys, Integer.toUnsignedLong(valueNodeInt32))); + payloadBirthNode.addMetrics(createMetric(valueNodeInt32, expectedSsKvEntryNodeInt32, metricDataType)); + listKeys.add(keys); + + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + payloadBirthNode.build().toByteArray(), 0, false); + + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NBIRTH.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), listKeys)); + return !finalFuture.get().get().isEmpty(); + }); + Assert.assertEquals(listKeys.size(), finalFuture.get().get().size()); + } + + protected void processClientWithCorrectAccessTokenPublishNCMDReBirth() throws Exception { + processClientWithCorrectNodeAccess(); + SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()); + List listKeys = new ArrayList<>(); + long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; + long valueBdSec = getBdSeqNum(); + MetricDataType metricDataType = Int64; + TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, valueBdSec)); + payloadBirthNode.addMetrics(createMetric(valueBdSec, tsKvEntryBdSecOriginal, metricDataType)); + listKeys.add(SparkplugMessageType.NCMD.name() + " " + keysBdSeq); + + String keys = "Node Control/Rebirth"; + boolean valueRebirth = true; + metricDataType = MetricDataType.Boolean; + TsKvEntry expectedSsKvEntryRebirth = new BasicTsKvEntry(ts, new BooleanDataEntry(keys, valueRebirth)); + payloadBirthNode.addMetrics(createMetric(valueRebirth, expectedSsKvEntryRebirth, metricDataType)); + listKeys.add(keys); + + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.NCMD.name() + "/" + edgeNode, + payloadBirthNode.build().toByteArray(), 0, false); + + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NCMD.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), listKeys)); + return !finalFuture.get().get().isEmpty(); + }); + Assert.assertEquals(listKeys.size(), finalFuture.get().get().size()); + } + + protected void processClientWithCorrectAccessTokenPushNodeMetricBuildPrimitiveSimple() throws Exception { + processClientWithCorrectNodeAccess(); + + String messageTypeName = SparkplugMessageType.NDATA.name(); + List listKeys = new ArrayList<>(); + List listTsKvEntry = new ArrayList<>(); + + SparkplugBProto.Payload.Builder ndataPayload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()) + .setSeq(getSeqNum()); + long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; + + createdAddMetricValuePrimitiveTsKv(listTsKvEntry, listKeys, ndataPayload, ts); + + if (client.isConnected()) { + client.publish(NAMESPACE + "/" + groupId + "/" + messageTypeName + "/" + edgeNode, + ndataPayload.build().toByteArray(), 0, false); + } + + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NDATA.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId())); + return finalFuture.get().get().size() == listTsKvEntry.size(); + }); + 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)); + } + + protected void processClientWithCorrectAccessTokenPushNodeMetricBuildArraysSimple() throws Exception { + processClientWithCorrectNodeAccess(); + + String messageTypeName = SparkplugMessageType.NDATA.name(); + List listKeys = new ArrayList<>(); + List listTsKvEntry = new ArrayList<>(); + + SparkplugBProto.Payload.Builder ndataPayload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(calendar.getTimeInMillis()) + .setSeq(getSeqNum()); + long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; + + createdAddMetricValueArraysTsKv(listTsKvEntry, listKeys, ndataPayload, ts); + + if (client.isConnected()) { + client.publish(NAMESPACE + "/" + groupId + "/" + messageTypeName + "/" + edgeNode, + ndataPayload.build().toByteArray(), 0, false); + } + + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.NDATA.name()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId())); + return finalFuture.get().get().size() == listTsKvEntry.size(); + }); + 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)); + } + + private void createdAddMetricValuePrimitiveTsKv(List listTsKvEntry, List listKeys, + SparkplugBProto.Payload.Builder dataPayload, long ts) throws ThingsboardException { + + String keys = "MyInt8"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextInt8(), ts, Int8)); + listKeys.add(keys); + + keys = "MyInt16"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextInt16(), ts, Int16)); + listKeys.add(keys); + + keys = "MyInt32"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextInt32(), ts, Int32)); + listKeys.add(keys); + + keys = "MyInt64"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextInt64(), ts, Int64)); + listKeys.add(keys); + + keys = "MyUInt8"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt8(), ts, UInt8)); + listKeys.add(keys); + + keys = "MyUInt16"; + 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)); + listKeys.add(keys); + + keys = "MyUInt64"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextUInt64(), ts, UInt64)); + listKeys.add(keys); + + keys = "MyFloat"; + listTsKvEntry.add(createdAddMetricTsKvFloat(dataPayload, keys, nextFloat(0, 100), ts, MetricDataType.Float)); + listKeys.add(keys); + + keys = "MyDateTime"; + listTsKvEntry.add(createdAddMetricTsKvLong(dataPayload, keys, nextDateTime(), ts, MetricDataType.DateTime)); + listKeys.add(keys); + + keys = "MyDouble"; + listTsKvEntry.add(createdAddMetricTsKvDouble(dataPayload, keys, nextDouble(), ts, MetricDataType.Double)); + listKeys.add(keys); + + keys = "MyBoolean"; + listTsKvEntry.add(createdAddMetricTsKvBoolean(dataPayload, keys, nextBoolean(), ts, MetricDataType.Boolean)); + listKeys.add(keys); + + keys = "MyString"; + listTsKvEntry.add(createdAddMetricTsKvString(dataPayload, keys, nexString(), ts, MetricDataType.String)); + listKeys.add(keys); + + keys = "MyText"; + listTsKvEntry.add(createdAddMetricTsKvString(dataPayload, keys, nexString(), ts, MetricDataType.Text)); + listKeys.add(keys); + + keys = "MyUUID"; + listTsKvEntry.add(createdAddMetricTsKvString(dataPayload, keys, nexString(), ts, MetricDataType.UUID)); + listKeys.add(keys); + + } + + private void createdAddMetricValueArraysTsKv(List listTsKvEntry, List listKeys, + SparkplugBProto.Payload.Builder dataPayload, long ts) throws ThingsboardException { + 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 = "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()}; + createdAddMetricTsKvJson(dataPayload, keys, uInt8s, ts, UInt8Array, listTsKvEntry, listKeys); + + keys = "MyUInt16Array"; + int[] uInt16s = {nextUInt16(), nextUInt16(), nextUInt16()}; + createdAddMetricTsKvJson(dataPayload, keys, uInt16s, ts, UInt16Array, listTsKvEntry, listKeys); + + keys = "MyUInt32LArray"; + long[] uInt32Ls = {nextUInt32L(), nextUInt32L(), nextUInt32L()}; + createdAddMetricTsKvJson(dataPayload, keys, uInt32Ls, ts, UInt32Array, 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 = "MyDateTimeArray"; + long[] dateTimes = {nextDateTime(), nextDateTime(), nextDateTime()}; + createdAddMetricTsKvJson(dataPayload, keys, dateTimes, ts, DateTimeArray, 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 keys, Object value, + long ts, MetricDataType metricDataType) throws ThingsboardException { + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, Long.valueOf(String.valueOf(value)))); + + dataPayload.addMetrics(createMetric(value, tsKvEntry, metricDataType)); + return tsKvEntry; + } + + private TsKvEntry createdAddMetricTsKvFloat(SparkplugBProto.Payload.Builder dataPayload, String keys, Object value, + long ts, MetricDataType metricDataType) throws ThingsboardException { + var f = new BigDecimal(String.valueOf(value)); + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new DoubleDataEntry(keys, f.doubleValue())); + dataPayload.addMetrics(createMetric(value, tsKvEntry, metricDataType)); + return tsKvEntry; + } + + private TsKvEntry createdAddMetricTsKvDouble(SparkplugBProto.Payload.Builder dataPayload, String keys, double value, + long ts, MetricDataType metricDataType) throws ThingsboardException { + var d = new BigDecimal(String.valueOf(value)); + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, d.longValueExact())); + dataPayload.addMetrics(createMetric(value, tsKvEntry, metricDataType)); + return tsKvEntry; + } + + private TsKvEntry createdAddMetricTsKvBoolean(SparkplugBProto.Payload.Builder dataPayload, String keys, boolean value, + long ts, MetricDataType metricDataType) throws ThingsboardException { + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new BooleanDataEntry(keys, value)); + dataPayload.addMetrics(createMetric(value, tsKvEntry, metricDataType)); + return tsKvEntry; + } + + private TsKvEntry createdAddMetricTsKvString(SparkplugBProto.Payload.Builder dataPayload, String keys, String value, + long ts, MetricDataType metricDataType) throws ThingsboardException { + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(keys, value)); + dataPayload.addMetrics(createMetric(value, tsKvEntry, metricDataType)); + return tsKvEntry; + } + + private void createdAddMetricTsKvJson(SparkplugBProto.Payload.Builder dataPayload, String keys, + Object values, long ts, MetricDataType metricDataType, + List listTsKvEntry, + List listKeys) throws ThingsboardException { + ArrayNode nodeArray = newArrayNode(); + switch (metricDataType) { + case Bytes: + 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 int[]) { + for (int b : (int[])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: + throw new IllegalStateException("Unexpected value: " + metricDataType); + } + if (nodeArray.size() > 0) { + Optional tsKvEntryOptional = Optional.of(new BasicTsKvEntry(ts, new JsonDataEntry(keys, nodeArray.toString()))); + if (tsKvEntryOptional.isPresent()) { + dataPayload.addMetrics(createMetric(values, tsKvEntryOptional.get(), metricDataType)); + listTsKvEntry.add(tsKvEntryOptional.get()); + listKeys.add(keys); + } + } + } + + private byte nextInt8() { + return (byte) random.nextInt(Byte.MIN_VALUE, Byte.MAX_VALUE); + } + + private short nextUInt8() { + return (short) random.nextInt(0, Byte.MAX_VALUE * 2 + 1); + } + + private short nextInt16() { + 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 nextInt32() { + return random.nextInt(Integer.MIN_VALUE, Integer.MAX_VALUE); + } + + private int nextUInt32I() { + return random.nextInt(0, Integer.MAX_VALUE); + } + + private long nextUInt32L() { + long l = Integer.MAX_VALUE; + return random.nextLong(0, l * 2 + 1); + } + + private long nextInt64() { + return random.nextLong(Long.MIN_VALUE, Long.MAX_VALUE); + } + + private long nextUInt64() { + double d = Long.MAX_VALUE; + return random.nextLong(0, (long) (d * 2 + 1)); + } + + private double nextDouble() { + return random.nextDouble(Long.MIN_VALUE, Long.MAX_VALUE); + } + + private long nextDateTime() { + long min = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; + long max = calendar.getTimeInMillis(); + return random.nextLong(min, max); + } + + private float nextFloat(float min, float max) { + if (min >= max) + throw new IllegalArgumentException("max must be greater than min"); + float result = ThreadLocalRandom.current().nextFloat() * (max - min) + min; + if (result >= max) // correct for rounding + result = Float.intBitsToFloat(Float.floatToIntBits(max) - 1); + return result; + } + + private boolean nextBoolean() { + return random.nextBoolean(); + } + + private String nexString() { + return java.util.UUID.randomUUID().toString(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java new file mode 100644 index 0000000000..917d1a9f65 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java @@ -0,0 +1,61 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.sparkplug.timeseries; + +import org.eclipse.paho.mqttv5.common.MqttException; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.thingsboard.server.dao.service.DaoSqlTest; + +/** + * Created by nickAS21 on 12.01.23 + */ +@DaoSqlTest +public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSparkplugTelemetryTest { + + @Before + public void beforeTest() throws Exception { + beforeSparkplugTest(); + } + + @After + public void afterTest () throws MqttException { + if (client.isConnected()) { + client.disconnect(); } + } + + @Test + public void testClientWithCorrectAccessTokenPublishNBIRTH() throws Exception { + processClientWithCorrectAccessTokenPublishNBIRTH(); + } + + @Test + public void testClientWithCorrectAccessTokenPublishNCMDReBirth() throws Exception { + processClientWithCorrectAccessTokenPublishNCMDReBirth(); + } + + @Test + public void testClientWithCorrectAccessTokenPushNodeMetricBuildPrimitiveSimple() throws Exception { + processClientWithCorrectAccessTokenPushNodeMetricBuildPrimitiveSimple(); + } + + @Test + public void testClientWithCorrectAccessTokenPushNodeMetricBuildPArraysSimple() throws Exception { + processClientWithCorrectAccessTokenPushNodeMetricBuildArraysSimple(); + } + +} \ 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 166eed3347..2a4dcbae16 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 @@ -17,6 +17,7 @@ package org.thingsboard.server.transport.mqtt; import com.fasterxml.jackson.databind.JsonNode; import com.google.gson.JsonParseException; +import com.google.protobuf.InvalidProtocolBufferException; import io.netty.channel.ChannelFuture; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; @@ -49,6 +50,7 @@ import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.MqttTopics; +import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.ota.OtaPackageType; @@ -71,13 +73,14 @@ import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509Ce import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; import org.thingsboard.server.queue.scheduler.SchedulerComponent; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; +import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor; import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; import org.thingsboard.server.transport.mqtt.session.GatewaySessionHandler; import org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher; import org.thingsboard.server.transport.mqtt.session.SparkplugNodeSessionHandler; -import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic; import org.thingsboard.server.transport.mqtt.util.ReturnCode; import org.thingsboard.server.transport.mqtt.util.ReturnCodeResolver; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic; import javax.net.ssl.SSLPeerUnverifiedException; import java.io.IOException; @@ -97,16 +100,15 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; import static com.amazonaws.util.StringUtils.UTF8; -import static io.netty.handler.codec.mqtt.MqttMessageType.CONNACK; import static io.netty.handler.codec.mqtt.MqttMessageType.CONNECT; import static io.netty.handler.codec.mqtt.MqttMessageType.PINGRESP; import static io.netty.handler.codec.mqtt.MqttMessageType.SUBACK; -import static io.netty.handler.codec.mqtt.MqttMessageType.UNSUBACK; import static io.netty.handler.codec.mqtt.MqttQoS.AT_LEAST_ONCE; import static io.netty.handler.codec.mqtt.MqttQoS.AT_MOST_ONCE; import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_CLOSED; import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_OPEN; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopic; +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 @@ -123,7 +125,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; private final UUID sessionId; - private final MqttTransportContext context; + protected final MqttTransportContext context; private final TransportService transportService; private final SchedulerComponent scheduler; private final SslHandler sslHandler; @@ -324,15 +326,15 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement String topicName = mqttMsg.variableHeader().topicName(); int msgId = mqttMsg.variableHeader().packetId(); log.trace("[{}][{}] Processing publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId); - - if (sparkplugSessionHandler != null) { - handleSparkplugPublishMsg(ctx, topicName, msgId, mqttMsg); - transportService.reportActivity(deviceSessionCtx.getSessionInfo()); - } else if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) { + if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) { if (gatewaySessionHandler != null) { handleGatewayPublishMsg(ctx, topicName, msgId, mqttMsg); transportService.reportActivity(deviceSessionCtx.getSessionInfo()); + } else { + log.error("[gatewaySessionHandler] is null, [{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId); } + } else if (sparkplugSessionHandler != null) { + handleSparkplugPublishMsg(ctx, topicName, mqttMsg); } else { processDevicePublish(ctx, mqttMsg, topicName, msgId); } @@ -375,14 +377,58 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private void handleSparkplugPublishMsg(ChannelHandlerContext ctx, String topicName, int msgId, MqttPublishMessage mqttMsg) { + private void handleSparkplugPublishMsg(ChannelHandlerContext ctx, String topicName, MqttPublishMessage mqttMsg) { + int msgId = mqttMsg.variableHeader().packetId(); try { - sparkplugSessionHandler.onPublishMsg(ctx, topicName, msgId, mqttMsg); + SparkplugTopic sparkplugTopic = parseTopicPublish(topicName); + String deviceName = sparkplugTopic.isNode() ? deviceSessionCtx.getDeviceInfo().getDeviceName() : sparkplugTopic.getDeviceId(); + if (sparkplugTopic.isNode()) { + // A node topic + switch (sparkplugTopic.getType()) { + case STATE: + // TODO + break; + case NBIRTH: + case NCMD: + case NDATA: + SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(mqttMsg.payload())); + sparkplugSessionHandler.onDeviceTelemetryProto(msgId, sparkplugBProtoNode, deviceName, sparkplugTopic.getType().name(), sparkplugTopic.isNode()); + break; + case NDEATH: + sparkplugSessionHandler.onDeviceDisconnect(mqttMsg); + break; + case NRECORD: + // TODO + break; + default: + } + } else { + // A device topic + switch (sparkplugTopic.getType()) { + case STATE: + // TODO + break; + case DCMD: + case DDATA: + case DBIRTH: + SparkplugBProto.Payload sparkplugBProtoDevice = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(mqttMsg.payload())); + sparkplugSessionHandler.onDeviceTelemetryProto(msgId, sparkplugBProtoDevice, deviceName, sparkplugTopic.getType().name(), sparkplugTopic.isNode()); + break; + case DDEATH: + sparkplugSessionHandler.onDeviceDisconnect(mqttMsg); + break; + case DRECORD: + // TODO + break; + default: + } + } } catch (RuntimeException e) { - log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); + log.error("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); + ack(ctx, msgId, ReturnCode.IMPLEMENTATION_SPECIFIC); ctx.close(); - } catch (Exception e) { - log.debug("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); + } catch (AdaptorException | ThingsboardException | InvalidProtocolBufferException e) { + log.error("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); sendAckOrCloseSession(ctx, topicName, msgId); } } @@ -648,7 +694,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement MqttQoS reqQoS = subscription.qualityOfService(); try { if (sparkplugSessionHandler != null) { - SparkplugTopic sparkplugTopic = parseTopic(mqttMsg.payload().topicSubscriptions().get(0).topicName()); + SparkplugTopic sparkplugTopic = parseTopicSubscribe(mqttMsg.payload().topicSubscriptions().get(0).topicName()); sparkplugSessionHandler.handleSparkplugSubscribeMsg(grantedQoSList, sparkplugTopic, reqQoS); } else { switch (topic) { @@ -1012,14 +1058,15 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void checkSparkplugSession(MqttConnectMessage connectMessage) { try { - SparkplugTopic sparkplugTopic = parseTopic(connectMessage.payload().willTopic()); - // Test proto - SparkplugBProto.Payload payloadBProto = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); - // if (sparkplugSessionHandler == null) { - sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId, sparkplugTopic.toString()); - } else { - log.warn("SparkPlugNodeReConnected [{}] [{}]", sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); + sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId); + if (StringUtils.isNotBlank(connectMessage.payload().willTopic()) + && connectMessage.payload().willMessageInBytes() != null && connectMessage.payload().willMessageInBytes().length > 0) { + SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); + SparkplugTopic sparkplugTopic = parseTopicPublish(connectMessage.payload().willTopic()); + sparkplugSessionHandler.onDeviceTelemetryProto(0, sparkplugBProtoNode, + deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopic.getType().name(), true); + } } } catch (Exception e) { log.trace("[{}][{}] Failed to fetch sparkplugDevice additional info or sparkplugTopicName", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName(), e); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index e330e3d5db..9ea5cbb00f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java @@ -79,20 +79,20 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe @Slf4j public abstract class AbstractGatewaySessionHandler { - private static final String DEFAULT_DEVICE_TYPE = "default"; + protected static final String DEFAULT_DEVICE_TYPE = "default"; private static final String CAN_T_PARSE_VALUE = "Can't parse value: "; private static final String DEVICE_PROPERTY = "device"; - private final MqttTransportContext context; + protected final MqttTransportContext context; private final TransportService transportService; - private final TransportDeviceInfo gateway; - private final UUID sessionId; + protected final TransportDeviceInfo gateway; + protected final UUID sessionId; private final ConcurrentMap deviceCreationLockMap; private final ConcurrentMap devices; private final ConcurrentMap> deviceFutures; private final ConcurrentMap mqttQoSMap; - private final ChannelHandlerContext channel; - private final DeviceSessionCtx deviceSessionCtx; + protected final ChannelHandlerContext channel; + protected final DeviceSessionCtx deviceSessionCtx; public AbstractGatewaySessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId) { this.context = deviceSessionCtx.getContext(); @@ -198,7 +198,7 @@ public abstract class AbstractGatewaySessionHandler { return deviceSessionCtx.isJsonPayloadType(); } - private void processOnConnect(MqttPublishMessage msg, String deviceName, String deviceType) { + protected void processOnConnect(MqttPublishMessage msg, String deviceName, String deviceType) { log.trace("[{}] onDeviceConnect: {}", sessionId, deviceName); Futures.addCallback(onDeviceConnect(deviceName, deviceType), new FutureCallback() { @Override @@ -393,7 +393,7 @@ public abstract class AbstractGatewaySessionHandler { } } - private void processPostTelemetryMsg(MqttDeviceAwareSessionContext deviceCtx, TransportProtos.PostTelemetryMsg postTelemetryMsg, String deviceName, int msgId) { + protected void processPostTelemetryMsg(MqttDeviceAwareSessionContext deviceCtx, TransportProtos.PostTelemetryMsg postTelemetryMsg, String deviceName, int msgId) { transportService.process(deviceCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(channel, deviceName, msgId, postTelemetryMsg)); } @@ -666,7 +666,7 @@ public abstract class AbstractGatewaySessionHandler { return result.build(); } - private ListenableFuture checkDeviceConnected(String deviceName) { + protected ListenableFuture checkDeviceConnected(String deviceName) { MqttDeviceAwareSessionContext ctx = devices.get(deviceName); if (ctx == null) { log.debug("[{}] Missing device [{}] for the gateway session", sessionId, deviceName); @@ -676,7 +676,7 @@ public abstract class AbstractGatewaySessionHandler { } } - private String checkDeviceName(String deviceName) { + protected String checkDeviceName(String deviceName) { if (StringUtils.isEmpty(deviceName)) { throw new RuntimeException("Device name is empty!"); } else { @@ -697,11 +697,11 @@ public abstract class AbstractGatewaySessionHandler { return JsonMqttAdaptor.validateJsonPayload(sessionId, mqttMsg.payload()); } - private byte[] getBytes(ByteBuf payload) { + protected byte[] getBytes(ByteBuf payload) { return ProtoMqttAdaptor.toBytes(payload); } - private void ack(MqttPublishMessage msg, ReturnCode returnCode) { + protected void ack(MqttPublishMessage msg, ReturnCode returnCode) { int msgId = getMsgId(msg); if (msgId > 0) { writeAndFlush(MqttTransportHandler.createMqttPubAckMsg(deviceSessionCtx, msgId, returnCode)); 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 9594b55515..9c0351b189 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 @@ -15,93 +15,95 @@ */ package org.thingsboard.server.transport.mqtt.session; -import io.netty.channel.ChannelHandlerContext; +import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +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 lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; +import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.transport.adaptor.AdaptorException; +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.util.sparkplug.SparkplugTopic; +import javax.annotation.Nullable; +import java.util.ArrayList; import java.util.List; +import java.util.Optional; import java.util.UUID; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopic; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.DBIRTH; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.getFromSparkplugBMetricToKeyValueProto; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopicSubscribe; /** * Created by nickAS21 on 12.12.22 */ @Slf4j -public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler{ +public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { + public SparkplugNodeSessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId) { + super(deviceSessionCtx, sessionId); + } - private String nodeTopic; - public SparkplugNodeSessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId, String nodeTopic) { - super(deviceSessionCtx, sessionId); - this.nodeTopic = nodeTopic; + public TransportProtos.PostTelemetryMsg convertToPostTelemetry(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { + DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx; + byte[] bytes = getBytes(inbound.payload()); + Descriptors.Descriptor telemetryDynamicMsgDescriptor = ProtoConverter.validateDescriptor(deviceSessionCtx.getTelemetryDynamicMsgDescriptor()); + try { + return JsonConverter.convertToTelemetryProto(new JsonParser().parse(ProtoConverter.dynamicMsgToJson(bytes, telemetryDynamicMsgDescriptor))); + } catch (Exception e) { + log.debug("Failed to decode post telemetry request", e); + throw new AdaptorException(e); + } } - public void onPublishMsg(ChannelHandlerContext ctx, String topicName, int msgId, MqttPublishMessage mqttMsg) throws Exception { - SparkplugTopic sparkplugTopic = parseTopic(topicName); - log.warn("SparkplugPublishMsg [{}] [{}]", sparkplugTopic.isNode() ? "node" : "device: " + sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); - if (sparkplugTopic.isNode()) { - // A node topic - switch (sparkplugTopic.getType()) { - case STATE: - // TODO - break; - case NBIRTH: - // TODO - break; - case NCMD: - // TODO - break; - case NDATA: - // TODO - break; - case NDEATH: - onGatewayDeviceDisconnectProto(mqttMsg); - break; - case NRECORD: - // TODO - break; - default: - } - } else { - // A device topic - switch (sparkplugTopic.getType()) { - case STATE: - // TODO - break; - case DBIRTH: - onDeviceConnectProto(mqttMsg); - break; - case DCMD: - // TODO - break; - case DDATA: - // TODO - break; - case DDEATH: - onGatewayDeviceDisconnectProto(mqttMsg); - break; - case DRECORD: - // TODO - break; - default: + public void onDeviceTelemetryProto(int msgId, SparkplugBProto.Payload sparkplugBProto, String deviceName, String topicTypeName, boolean isNode) throws AdaptorException { + try { + checkDeviceName(deviceName); + List msgs = convertToPostTelemetry(sparkplugBProto, topicTypeName); + int finalMsgId = msgId; + ListenableFuture contextListenableFuture = isNode ? + Futures.immediateFuture(this.deviceSessionCtx) : checkDeviceConnected(deviceName); + for (TransportProtos.PostTelemetryMsg msg : msgs) { + Futures.addCallback(contextListenableFuture, + new FutureCallback<>() { + @Override + public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + try { + processPostTelemetryMsg(deviceCtx, msg, deviceName, finalMsgId); + } catch (Throwable e) { + log.warn("[{}][{}] Failed to convert telemetry: {}", gateway.getDeviceId(), deviceName, msg, e); + channel.close(); + } + } + + @Override + public void onFailure(Throwable t) { + log.debug("[{}] Failed to process device telemetry command: {}", sessionId, deviceName, t); + } + }, context.getExecutor()); } + } catch (RuntimeException e) { + throw new AdaptorException(e); } } public void handleSparkplugSubscribeMsg(List grantedQoSList, SparkplugTopic sparkplugTopic, MqttQoS reqQoS) { - String topicName = sparkplugTopic.toString(); - log.warn("SparkplugSubscribeMsg [{}] [{}]", sparkplugTopic.isNode() ? "node" : "device: " + sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); - if (sparkplugTopic.getGroupId() == null) { // TODO SUBSCRIBE NameSpace } else if (sparkplugTopic.getType() == null) { // TODO SUBSCRIBE GroupId - } - else if (sparkplugTopic.isNode()) { + } else if (sparkplugTopic.isNode()) { // A node topic switch (sparkplugTopic.getType()) { case STATE: @@ -150,4 +152,57 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler{ } } + private List convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto, String topicTypeName) throws AdaptorException { + try { + List msgs = new ArrayList<>(); + for (SparkplugBProto.Payload.Metric protoMetric : sparkplugBProto.getMetricsList()) { + long ts = protoMetric.getTimestamp(); + String keys = "bdSeq".equals(protoMetric.getName()) ? + topicTypeName + " " + protoMetric.getName() : protoMetric.getName(); + Optional keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(keys, protoMetric); + if (keyValueProtoOpt.isPresent()) { + List result = new ArrayList<>(); + result.add(keyValueProtoOpt.get()); + TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); + TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); + builder.setTs(ts); + builder.addAllKv(result); + request.addTsKvList(builder.build()); + msgs.add(request.build()); + } + } + if (DBIRTH.name().equals(topicTypeName)) { + List result = new ArrayList<>(); + TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); + keyValueProtoBuilder.setKey(topicTypeName + " " + "seq"); + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.LONG_V); + keyValueProtoBuilder.setLongV(sparkplugBProto.getSeq()); + result.add(keyValueProtoBuilder.build()); + TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); + TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); + builder.setTs(sparkplugBProto.getTimestamp()); + builder.addAllKv(result); + request.addTsKvList(builder.build()); + msgs.add(request.build()); + } + return msgs; + } catch (IllegalStateException | JsonSyntaxException | ThingsboardException e) { + log.error("Failed to decode post telemetry request", e); + throw new AdaptorException(e); + } + } + + public void onDeviceConnectProto(MqttPublishMessage mqttPublishMessage, String nodeDeviceType) throws ThingsboardException { + try { + String topic = mqttPublishMessage.variableHeader().topicName(); + SparkplugTopic sparkplugTopic = parseTopicSubscribe(topic); + String deviceName = checkDeviceName(sparkplugTopic.getDeviceId()); + String deviceType = StringUtils.isEmpty(nodeDeviceType) ? DEFAULT_DEVICE_TYPE : nodeDeviceType; + processOnConnect(mqttPublishMessage, deviceName, deviceType); + } catch (RuntimeException | ThingsboardException e) { + log.error("Failed Sparkplug Device connect proto!", e); + throw new ThingsboardException(e, ThingsboardErrorCode.BAD_REQUEST_PARAMS); + } + } + } 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 new file mode 100644 index 0000000000..f4dc74f46b --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java @@ -0,0 +1,199 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.util.sparkplug; + +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.transport.adaptor.AdaptorException; +import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; + +import java.math.BigInteger; +import java.util.Date; + +/** + * Created by nickAS21 on 10.01.23 + */ + +@Slf4j +public enum MetricDataType { + + // Basic Types + Int8(1, Byte.class), + Int16(2, Short.class), + Int32(3, Integer.class), + Int64(4, Long.class), + UInt8(5, Short.class), + UInt16(6, Integer.class), + UInt32(7, Long.class), + UInt64(8, BigInteger.class), + Float(9, Float.class), + Double(10, Double.class), + Boolean(11, Boolean.class), + String(12, String.class), + DateTime(13, Date.class), + Text(14, String.class), + + // Custom Types for Metrics + UUID(15, String.class), + DataSet(16, SparkplugBProto.Payload.DataSet.class), + Bytes(17, byte[].class), + File(18, SparkplugMetricUtil.File.class), + Template(19, SparkplugBProto.Payload.Template.class), + + // 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); + + private Class clazz = null; + private int intValue = 0; + + /** + * Constructor + * + * @param intValue the integer value of this {@link MetricDataType} + * @param clazz the {@link Class} type associated with this {@link MetricDataType} + */ + private MetricDataType(int intValue, Class clazz) { + this.intValue = intValue; + this.clazz = clazz; + } + + /** + * Checks the type of a specified value against the specified {@link MetricDataType} + * + * @param value the {@link Object} value to check against the {@link MetricDataType} + * @throws AdaptorException if the value is not a valid type for the given {@link MetricDataType} + */ + public void checkType(Object value) throws AdaptorException { + if (value != null && !clazz.isAssignableFrom(value.getClass())) { + String msgError = "Failed type check - " + clazz + " != " + ((value != null) ? value.getClass().toString() : "null"); + log.debug(msgError); + throw new AdaptorException(msgError); + } + } + + /** + * Returns an integer representation of the data type. + * + * @return an integer representation of the data type. + */ + public int toIntValue() { + return this.intValue; + } + + /** + * Converts the integer representation of the data type into a {@link MetricDataType} instance. + * + * @param i the integer representation of the data type. + * @return a {@link MetricDataType} instance. + */ + public static MetricDataType fromInteger(int i) { + switch (i) { + case 1: + return Int8; + case 2: + return Int16; + case 3: + return Int32; + case 4: + return Int64; + case 5: + return UInt8; + case 6: + return UInt16; + case 7: + return UInt32; + case 8: + return UInt64; + case 9: + return Float; + case 10: + return Double; + case 11: + return Boolean; + case 12: + return String; + case 13: + return DateTime; + case 14: + return Text; + case 15: + return UUID; + case 16: + return DataSet; + case 17: + return Bytes; + case 18: + 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; + } + } + + /** + * Returns the class type for this DataType + * + * @return the class type for this DataType + */ + public Class getClazz() { + return clazz; + } + + +} \ No newline at end of file diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java new file mode 100644 index 0000000000..c24d8a5d36 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java @@ -0,0 +1,296 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.util.sparkplug; + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.fasterxml.jackson.databind.ser.std.FileSerializer; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; +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.ObjectInputStream; +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.util.Arrays; +import java.util.Optional; + +import static org.thingsboard.common.util.JacksonUtil.newArrayNode; + +/** + * Provides utility methods for SparkplugB MQTT Payload Metric. + */ +@Slf4j +public class SparkplugMetricUtil { + + public static Optional getFromSparkplugBMetricToKeyValueProto(String key, SparkplugBProto.Payload.Metric protoMetric) throws ThingsboardException { + // Check if the null flag has been set indicating that the value is null + if (protoMetric.getIsNull()) { + return Optional.empty(); + } + // Otherwise convert the value based on the type + int metricType = protoMetric.getDatatype(); + TransportProtos.KeyValueProto.Builder builderProto = TransportProtos.KeyValueProto.newBuilder(); + ArrayNode nodeArray = newArrayNode(); + MetricDataType metricDataType = MetricDataType.fromInteger(metricType); + try { + switch (metricDataType) { + case Boolean: + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.BOOLEAN_V) + .setBoolV(protoMetric.getBooleanValue()).build()); + case DateTime: + case Int64: + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(protoMetric.getLongValue()).build()); + case Float: + var f = new BigDecimal(String.valueOf(protoMetric.getFloatValue())); + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.DOUBLE_V) + .setDoubleV(f.doubleValue()).build()); + case Double: + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.DOUBLE_V) + .setDoubleV(protoMetric.getDoubleValue()).build()); + case Int8: + case UInt8: + case Int16: + case Int32: + case UInt16: + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(protoMetric.getIntValue()).build()); + case UInt32: + case UInt64: + if (protoMetric.hasIntValue()) { + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(protoMetric.getIntValue()).build()); + } else if (protoMetric.hasLongValue()) { + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(protoMetric.getLongValue()).build()); + } else { + log.error("Invalid value for UInt32 datatype"); + throw new ThingsboardException("Invalid value for " + MetricDataType.fromInteger(metricType).name() + " datatype " + metricType, ThingsboardErrorCode.INVALID_ARGUMENTS); + } + case String: + case Text: + case UUID: + 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: + //TODO + // Build the and create the DataSet + /** + SparkplugBProto.Payload.DataSet protoDataSet = protoMetric.getDatasetValue(); + return new SparkplugBProto.Payload.DataSet.Builder(protoDataSet.getNumOfColumns()).addColumnNames(protoDataSet.getColumnsList()) + .addTypes(convertDataSetDataTypes(protoDataSet.getTypesList())) + .addRows(convertDataSetRows(protoDataSet.getRowsList(), protoDataSet.getTypesList())) + .createDataSet(); + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(protoDataSet.toString()).build()); + **/ + //TODO + // Build the and create the Template + /** + SparkplugBProto.Payload.Template protoTemplate = protoMetric.getTemplateValue(); + return Optional.of(builderProto.setKey(key).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV( protoTemplate.toString()).build()); + **/ + //TODO + // Build the and create the File + /** + String filename = protoMetric.getMetadata().getFileName(); + return Optional.of(builderPrbyteValueoto.setKey(key + "_" + filename).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(Hex.encodeHexString((protoMetric.getBytesValue().toByteArray()))).build()); + **/ + return Optional.empty(); + case Unknown: + default: + throw new ThingsboardException("Failed to decode: Unknown MetricDataType " + metricType, ThingsboardErrorCode.INVALID_ARGUMENTS); + } + } catch (Exception e){ + log.error("", e); + return Optional.empty(); + } + } + + @JsonIgnoreProperties( + value = {"fileName"}) + @JsonSerialize( + using = FileSerializer.class) + public class File { + + private String fileName; + private byte[] bytes; + + /** + * Default Constructor + */ + public File() { + super(); + } + + /** + * Constructor + * + * @param fileName the full file name path + * @param bytes the array of bytes that represent the contents of the file + */ + public File(String fileName, byte[] bytes) { + super(); + this.fileName = fileName == null + ? null + : fileName.replace("/", System.getProperty("file.separator")).replace("\\", + System.getProperty("file.separator")); + this.bytes = Arrays.copyOf(bytes, bytes.length); + } + + /** + * Gets the full filename path + * + * @return the full filename path + */ + public String getFileName() { + return fileName; + } + + /** + * Sets the full filename path + * + * @param fileName the full filename path + */ + public void setFileName(String fileName) { + this.fileName = fileName; + } + + /** + * Gets the bytes that represent the contents of the file + * + * @return the bytes that represent the contents of the file + */ + public byte[] getBytes() { + return bytes; + } + + /** + * Sets the bytes that represent the contents of the file + * + * @param bytes the bytes that represent the contents of the file + */ + public void setBytes(byte[] bytes) { + this.bytes = bytes; + } + + @Override + public String toString() { + StringBuilder builder = new StringBuilder(); + builder.append("File [fileName="); + builder.append(fileName); + builder.append(", bytes="); + builder.append(Arrays.toString(bytes)); + builder.append("]"); + return builder.toString(); + } + } + +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java index c2e93cb269..93d0245da6 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java @@ -27,88 +27,79 @@ import java.util.Map; * Provides utility methods for handling Sparkplug MQTT message topics. */ public class SparkplugTopicUtil { - - private static final Map SPLIT_TOPIC_CACHE = new HashMap(); - - public static String[] getSplitTopic(String topic) { - String[] splitTopic = SPLIT_TOPIC_CACHE.get(topic); - if (splitTopic == null) { - splitTopic = topic.split("/"); - SPLIT_TOPIC_CACHE.put(topic, splitTopic); - } - - return splitTopic; - } - /** - * Serializes a {@link SparkplugTopic} instance in to a JSON string. - * - * @param topic a {@link SparkplugTopic} instance - * @return a JSON string - * @throws JsonProcessingException - */ - public static String sparkplugTopicToString(SparkplugTopic topic) throws JsonProcessingException { - ObjectMapper mapper = new ObjectMapper(); - return mapper.writeValueAsString(topic); - } + private static final Map SPLIT_TOPIC_CACHE = new HashMap(); + private static final String TOPIC_INVALID_NUMBER = "Invalid number of topic elements: "; - /** - * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. - * - * @param topic a topic string - * @return a {@link SparkplugTopic} instance - * @throws ThingsboardException if an error occurs while parsing - */ - public static SparkplugTopic parseTopic(String topic) throws ThingsboardException { - topic = topic.indexOf("#") > 0 ? topic.substring(0, topic.indexOf("#")) : topic; - return parseTopic(SparkplugTopicUtil.getSplitTopic(topic)); - } + public static String[] getSplitTopic(String topic) { + String[] splitTopic = SPLIT_TOPIC_CACHE.get(topic); + if (splitTopic == null) { + splitTopic = topic.split("/"); + SPLIT_TOPIC_CACHE.put(topic, splitTopic); + } - /** - * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. - * - * @param splitTopic a topic split into tokens - * @return a {@link SparkplugTopic} instance - * @throws Exception if an error occurs while parsing - */ - @SuppressWarnings("incomplete-switch") - public static SparkplugTopic parseTopic(String[] splitTopic) throws ThingsboardException { - SparkplugMessageType type; - String namespace, edgeNodeId, groupId; - int length = splitTopic.length; + return splitTopic; + } - if (length < 4 || length > 5) { - throw new ThingsboardException("Invalid number of topic elements: " + length, ThingsboardErrorCode.INVALID_ARGUMENTS); - } + /** + * Serializes a {@link SparkplugTopic} instance in to a JSON string. + * + * @param topic a {@link SparkplugTopic} instance + * @return a JSON string + * @throws JsonProcessingException + */ + public static String sparkplugTopicToString(SparkplugTopic topic) throws JsonProcessingException { + ObjectMapper mapper = new ObjectMapper(); + return mapper.writeValueAsString(topic); + } - namespace = splitTopic[0]; - groupId = splitTopic[1]; - type = SparkplugMessageType.parseMessageType(splitTopic[2]); - edgeNodeId = splitTopic[3]; + /** + * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. + * + * @param topic a topic string + * @return a {@link SparkplugTopic} instance + * @throws ThingsboardException if an error occurs while parsing + */ + public static SparkplugTopic parseTopicSubscribe(String topic) throws ThingsboardException { + // TODO "+", "$" + topic = topic.indexOf("#") > 0 ? topic.substring(0, topic.indexOf("#")) : topic; + return parseTopic(SparkplugTopicUtil.getSplitTopic(topic)); + } + + public static SparkplugTopic parseTopicPublish(String topic) throws ThingsboardException { + if (topic.contains("#") || topic.contains("$") || topic.contains("+")) { + throw new ThingsboardException("Invalid of topic elements for Publish", ThingsboardErrorCode.INVALID_ARGUMENTS); + } else { + String[] splitTopic = SparkplugTopicUtil.getSplitTopic(topic); + if (splitTopic.length < 4 || splitTopic.length > 5) { + throw new ThingsboardException(TOPIC_INVALID_NUMBER + splitTopic.length, ThingsboardErrorCode.INVALID_ARGUMENTS); + } + return parseTopic(splitTopic); + } + } + + /** + * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. + * + * @param splitTopic a topic split into tokens + * @return a {@link SparkplugTopic} instance + * @throws Exception if an error occurs while parsing + */ + @SuppressWarnings("incomplete-switch") + public static SparkplugTopic parseTopic(String[] splitTopic) throws ThingsboardException { + int length = splitTopic.length; + if (length == 0) { + throw new ThingsboardException(TOPIC_INVALID_NUMBER + length, ThingsboardErrorCode.INVALID_ARGUMENTS); + } else { + SparkplugMessageType type; + String namespace, edgeNodeId, groupId, deviceId; + namespace = splitTopic[0]; + groupId = length > 1 ? splitTopic[1] : null; + type = length > 2 ? SparkplugMessageType.parseMessageType(splitTopic[2]) : null; + edgeNodeId = length > 3 ? splitTopic[3] : null; + deviceId = length > 4 ? splitTopic[4] : null; + return new SparkplugTopic(namespace, groupId, edgeNodeId, deviceId, type); + } + } - if (length == 4) { - // A node topic - switch (type) { - case STATE: - case NBIRTH: - case NCMD: - case NDATA: - case NDEATH: - case NRECORD: - return new SparkplugTopic(namespace, groupId, edgeNodeId, type); - } - } else { - // A device topic - switch (type) { - case STATE: - case DBIRTH: - case DCMD: - case DDATA: - case DDEATH: - case DRECORD: - return new SparkplugTopic(namespace, groupId, edgeNodeId, splitTopic[4], type); - } - } - throw new ThingsboardException("Invalid number of topic elements " + length + " for topic type " + type, ThingsboardErrorCode.INVALID_ARGUMENTS); - } } diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index ec464caf7b..55feab05c2 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -155,6 +155,14 @@ public class JacksonUtil { return mapper.createObjectNode(); } + public static ArrayNode newArrayNode() { + return newArrayNode(OBJECT_MAPPER); + } + + public static ArrayNode newArrayNode(ObjectMapper mapper) { + return mapper.createArrayNode(); + } + public static T clone(T value) { @SuppressWarnings("unchecked") Class valueClass = (Class) value.getClass();