Browse Source

Merge pull request #7931 from thingsboard/sparkplag_Telemetry

sparkplug: Telemetry
pull/8032/head
Andrew Shvayka 4 years ago
committed by GitHub
parent
commit
87b59cace5
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 4
      application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java
  2. 1
      application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java
  3. 1
      application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java
  4. 4
      application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv5/MqttV5TestClient.java
  5. 256
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java
  6. 140
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java
  7. 62
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java
  8. 505
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java
  9. 61
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java
  10. 93
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  11. 24
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java
  12. 175
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java
  13. 199
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java
  14. 296
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java
  15. 147
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java
  16. 8
      common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java

4
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) {

1
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)) {

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

4
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 {

256
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++;
}
}

140
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<ListenableFuture<Optional<TsKvEntry>>> 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<String> 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> 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();
}
}

62
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);
}
}

505
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<String> listKeys = new ArrayList<>();
SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder()
.setTimestamp(calendar.getTimeInMillis());
long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS;
long valueBdSec = getBdSeqNum();
MetricDataType metricDataType = Int64;
TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, valueBdSec));
payloadBirthNode.addMetrics(createMetric(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<ListenableFuture<List<TsKvEntry>>> finalFuture = new AtomicReference<>();
await(alias + SparkplugMessageType.NBIRTH.name())
.atMost(40, TimeUnit.SECONDS)
.until(() -> {
finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), listKeys));
return !finalFuture.get().get().isEmpty();
});
Assert.assertEquals(listKeys.size(), finalFuture.get().get().size());
}
protected void processClientWithCorrectAccessTokenPublishNCMDReBirth() throws Exception {
processClientWithCorrectNodeAccess();
SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder()
.setTimestamp(calendar.getTimeInMillis());
List<String> listKeys = new ArrayList<>();
long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS;
long valueBdSec = getBdSeqNum();
MetricDataType metricDataType = Int64;
TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, valueBdSec));
payloadBirthNode.addMetrics(createMetric(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<ListenableFuture<List<TsKvEntry>>> finalFuture = new AtomicReference<>();
await(alias + SparkplugMessageType.NCMD.name())
.atMost(40, TimeUnit.SECONDS)
.until(() -> {
finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), listKeys));
return !finalFuture.get().get().isEmpty();
});
Assert.assertEquals(listKeys.size(), finalFuture.get().get().size());
}
protected void processClientWithCorrectAccessTokenPushNodeMetricBuildPrimitiveSimple() throws Exception {
processClientWithCorrectNodeAccess();
String messageTypeName = SparkplugMessageType.NDATA.name();
List<String> listKeys = new ArrayList<>();
List<TsKvEntry> 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<ListenableFuture<List<TsKvEntry>>> 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<String> listKeys = new ArrayList<>();
List<TsKvEntry> 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<ListenableFuture<List<TsKvEntry>>> 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<TsKvEntry> listTsKvEntry, List<String> 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<TsKvEntry> listTsKvEntry, List<String> 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<TsKvEntry> listTsKvEntry,
List<String> 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<TsKvEntry> 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();
}
}

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

93
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);

24
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<String, Lock> deviceCreationLockMap;
private final ConcurrentMap<String, MqttDeviceAwareSessionContext> devices;
private final ConcurrentMap<String, ListenableFuture<MqttDeviceAwareSessionContext>> deviceFutures;
private final ConcurrentMap<MqttTopicMatcher, Integer> 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<MqttDeviceAwareSessionContext>() {
@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<MqttDeviceAwareSessionContext> checkDeviceConnected(String deviceName) {
protected ListenableFuture<MqttDeviceAwareSessionContext> 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));

175
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<TransportProtos.PostTelemetryMsg> msgs = convertToPostTelemetry(sparkplugBProto, topicTypeName);
int finalMsgId = msgId;
ListenableFuture<MqttDeviceAwareSessionContext> 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<Integer> 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<TransportProtos.PostTelemetryMsg> convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto, String topicTypeName) throws AdaptorException {
try {
List<TransportProtos.PostTelemetryMsg> 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<TransportProtos.KeyValueProto> keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(keys, protoMetric);
if (keyValueProtoOpt.isPresent()) {
List<TransportProtos.KeyValueProto> 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<TransportProtos.KeyValueProto> 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);
}
}
}

199
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;
}
}

296
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<TransportProtos.KeyValueProto> 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();
}
}
}

147
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<String, String[]> SPLIT_TOPIC_CACHE = new HashMap<String, String[]>();
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<String, String[]> SPLIT_TOPIC_CACHE = new HashMap<String, String[]>();
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);
}
}

8
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> T clone(T value) {
@SuppressWarnings("unchecked")
Class<T> valueClass = (Class<T>) value.getClass();

Loading…
Cancel
Save