From d7050074bc744200656f2e807de58e5a51dc4fad Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Tue, 14 Feb 2023 18:55:28 +0200 Subject: [PATCH] sparkplug: rpc --- ...ctMqttV5ClientSparkplugAttributesTest.java | 5 + ...ctMqttV5ClientSparkplugConnectionTest.java | 132 ++++++++++++++-- .../MqttV5ClientSparkplugBConnectionTest.java | 15 ++ .../rpc/AbstractMqttV5RpcSparkplugTest.java | 59 ++++++++ .../transport/mqtt/MqttTransportHandler.java | 141 +++++++++++++----- .../AbstractGatewaySessionHandler.java | 4 +- .../SparkplugDeviceSessionContext.java | 38 +++++ .../session/SparkplugNodeSessionHandler.java | 21 ++- .../util/sparkplug/SparkplugMetricUtil.java | 66 ++++++-- .../sparkplug/SparkplugRpcRequestHeader.java | 29 ++++ .../sparkplug/SparkplugRpcResponseBody.java | 31 ++++ 11 files changed, 475 insertions(+), 66 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java create mode 100644 common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcRequestHeader.java create mode 100644 common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcResponseBody.java diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java index 57b14471cc..15427099ce 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java @@ -40,6 +40,11 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra protected ThreadLocalRandom random = ThreadLocalRandom.current(); + /** + * "sparkPlugAttributesMetricNames": ["SN node", "SN device", "Firmware version", "Date version", "Last date update"] + * @throws Exception + */ + protected void processClientWithCorrectAccessTokenPublishNCMDReBirth() throws Exception { clientWithCorrectNodeAccessTokenWithNDEATH(); List listKeys = new ArrayList<>(); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java index f8a4373036..ebf2c7248e 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java @@ -22,6 +22,7 @@ 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.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; @@ -29,7 +30,9 @@ import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSpark import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; +import java.util.ArrayList; import java.util.HashSet; +import java.util.List; import java.util.Optional; import java.util.Set; import java.util.concurrent.TimeUnit; @@ -37,16 +40,18 @@ import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate.OFFLINE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate.ONLINE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; /** * Created by nickAS21 on 12.01.23 */ @Slf4j -public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends AbstractMqttV5ClientSparkplugTest { +public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends AbstractMqttV5ClientSparkplugTest { protected void processClientWithCorrectNodeAccessTokenWithNDEATH_Test() throws Exception { - long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; + long ts = calendar.getTimeInMillis() - PUBLISH_TS_DELTA_MS; long value = bdSeq = 0; clientWithCorrectNodeAccessTokenWithNDEATH(ts, value); @@ -69,20 +74,116 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr String expectedMessage = "Server unavailable."; int expectedReasonCode = 136; Assert.assertEquals(expectedMessage, actualException.getMessage()); - Assert.assertEquals( expectedReasonCode, actualException.getReasonCode()); + Assert.assertEquals(expectedReasonCode, actualException.getReasonCode()); } protected void processClientWithCorrectAccessTokenWithNDEATHCreatedDevices(int cntDevices) throws Exception { - clientWithCorrectNodeAccessTokenWithNDEATH(); + Set devices = new HashSet<>(); + long ts = calendar.getTimeInMillis(); + connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(devices, cntDevices, ts); + } + + protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_ALL(int cntDevices) throws Exception { + Set devices = new HashSet<>(); + long ts = calendar.getTimeInMillis(); + connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(devices, cntDevices, ts); + + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(SparkplugMessageType.STATE.name(), ONLINE.name())); + AtomicReference>> finalFuture = new AtomicReference<>(); + await(alias + SparkplugMessageType.STATE.name() + ", device: " + savedGateway.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId())); + return finalFuture.get().get().contains(tsKvEntry); + }); + + for (Device device : devices) { + await(alias + SparkplugMessageType.STATE.name() + ", device: " + device.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findAllLatest(tenantId, device.getId())); + return finalFuture.get().get().contains(tsKvEntry); + }); + } + } + + protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OneDeviceOFFLINE(int cntDevices, int indexDeviceDisconnect) throws Exception { + Set devices = new HashSet<>(); + long ts = calendar.getTimeInMillis(); + connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(devices, cntDevices, ts); + + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(SparkplugMessageType.STATE.name(), OFFLINE.name())); + AtomicReference>> finalFuture = new AtomicReference<>(); + + SparkplugBProto.Payload.Builder payloadDeathDevice = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getSeqNum()); + if (client.isConnected()) { + List devicesList = new ArrayList<>(devices); + Device device = devicesList.get(indexDeviceDisconnect); + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.DDEATH.name() + "/" + edgeNode + "/" + device.getName(), + payloadDeathDevice.build().toByteArray(), 0, false); + await(alias + SparkplugMessageType.STATE.name() + ", device: " + device.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findAllLatest(tenantId, device.getId())); + return findEqualsKeyValueInKvEntrys(finalFuture.get().get(), tsKvEntry); + }); + } + } + + protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OFFLINE_All(int cntDevices) throws Exception { + Set devices = new HashSet<>(); long ts = calendar.getTimeInMillis(); + connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(devices, cntDevices, ts); + + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(SparkplugMessageType.STATE.name(), OFFLINE.name())); + AtomicReference>> finalFuture = new AtomicReference<>(); + + if (client.isConnected()) { + client.disconnect(); + + await(alias + SparkplugMessageType.STATE.name() + ", device: " + savedGateway.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findAllLatest(tenantId, savedGateway.getId())); + return findEqualsKeyValueInKvEntrys(finalFuture.get().get(), tsKvEntry); + }); + + List devicesList = new ArrayList<>(devices); + for (Device device : devicesList) { + await(alias + SparkplugMessageType.STATE.name() + ", device: " + device.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + finalFuture.set(tsService.findAllLatest(tenantId, device.getId())); + return findEqualsKeyValueInKvEntrys(finalFuture.get().get(), tsKvEntry); + }); + } + } + } + + private void connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(Set devices, int cntDevices, long ts) throws Exception { + clientWithCorrectNodeAccessTokenWithNDEATH(); MetricDataType metricDataType = Int32; - Set deviceIds = new HashSet<>(); - String key = "Device Metric int32"; + String key = "Node Metric int32"; int valueDeviceInt32 = 1024; SparkplugBProto.Payload.Metric metric = createMetric(valueDeviceInt32, ts, key, metricDataType); - for (int i=0; i < cntDevices; i++ ) { + SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getBdSeqNum()); + payloadBirthNode.addMetrics(metric); + payloadBirthNode.setTimestamp(ts); + if (client.isConnected()) { + client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + payloadBirthNode.build().toByteArray(), 0, false); + } + metricDataType = Int32; + key = "Device Metric int32"; + valueDeviceInt32 = 4024; + metric = createMetric(valueDeviceInt32, ts, key, metricDataType); + for (int i = 0; i < cntDevices; i++) { SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() - .setTimestamp(calendar.getTimeInMillis()) + .setTimestamp(ts) .setSeq(getSeqNum()); String deviceName = deviceId + "_" + i; @@ -97,10 +198,21 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); return device.get() != null; }); + devices.add(device.get()); + } + + } + + Assert.assertEquals(cntDevices, devices.size()); + } + + private boolean findEqualsKeyValueInKvEntrys(List finalFuture, TsKvEntry tsKvEntry) { + for (TsKvEntry kvEntry : finalFuture) { + if (kvEntry.getKey().equals(tsKvEntry.getKey()) && kvEntry.getValue().equals(tsKvEntry.getValue())) { + return true; } - deviceIds.add(deviceName); } - Assert.assertEquals(cntDevices, deviceIds.size()); + return false; } } 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 index 6c928358ef..b07c3f3d8d 100644 --- 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 @@ -60,4 +60,19 @@ public class MqttV5ClientSparkplugBConnectionTest extends AbstractMqttV5ClientSp processClientWithCorrectAccessTokenWithNDEATHCreatedDevices(2); } + @Test + public void testClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_ALL() throws Exception { + processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_ALL(3); + } + + @Test + public void testConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OneDeviceOFFLINE() throws Exception { + processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OneDeviceOFFLINE(3, 1); + } + + @Test + public void testConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OFFLINE_All() throws Exception { + processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OFFLINE_All(3); + } + } diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java new file mode 100644 index 0000000000..dfdc2d9ea9 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java @@ -0,0 +1,59 @@ +/** + * 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.rpc; + +import com.nimbusds.jose.util.StandardCharset; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.paho.mqttv5.common.MqttException; +import org.eclipse.paho.mqttv5.common.MqttMessage; +import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; +import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestCallback; +import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; + +@Slf4j +public abstract class AbstractMqttV5RpcSparkplugTest extends AbstractMqttIntegrationTest { + + private static final String DEVICE_RESPONSE = "{\"value1\":\"A\",\"value2\":\"B\"}"; + private static final String setSparklpugRpcNodeRequest = "{\"method\": \"NCMD\", \"params\": {\"MyNodeMetric05_String\":\"MyNodeMetric05_String_Value\"}}"; + private static final String setSparklpugRpcDeviceRequest = "{\"method\": \"DCMD\", \"params\": {\"MyDeviceMetric05_String\":{\"MyDeviceMetric05_String_Value\"}}"; + + protected class MqttV5TestRpcCallback extends MqttV5TestCallback { + + private final MqttV5TestClient client; + + public MqttV5TestRpcCallback(MqttV5TestClient client, String awaitSubTopic) { + super(awaitSubTopic); + this.client = client; + } + + @Override + protected void messageArrivedOnAwaitSubTopic(String requestTopic, MqttMessage mqttMessage) { + log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic); + if (awaitSubTopic.equals(requestTopic)) { + qoS = mqttMessage.getQos(); + payloadBytes = mqttMessage.getPayload(); + String responseTopic = requestTopic.replace("request", "response"); + try { + client.publish(responseTopic, DEVICE_RESPONSE.getBytes(StandardCharset.UTF_8)); + } catch (MqttException e) { + log.warn("Failed to publish response on topic: {} due to: ", responseTopic, e); + } + subscribeLatch.countDown(); + } + } + } + +} 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 5a8927e184..910d257640 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 @@ -43,6 +43,8 @@ import io.netty.util.ReferenceCountUtil; import io.netty.util.concurrent.Future; import io.netty.util.concurrent.GenericFutureListener; import lombok.extern.slf4j.Slf4j; +import org.eclipse.leshan.core.ResponseCode; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -82,6 +84,8 @@ import org.thingsboard.server.transport.mqtt.session.SparkplugNodeSessionHandler import org.thingsboard.server.transport.mqtt.util.ReturnCode; import org.thingsboard.server.transport.mqtt.util.ReturnCodeResolver; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugRpcRequestHeader; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugRpcResponseBody; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic; import javax.net.ssl.SSLPeerUnverifiedException; @@ -110,8 +114,11 @@ 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.common.transport.service.DefaultTransportService.SUBSCRIBE_TO_ATTRIBUTE_UPDATES_ASYNC_MSG; +import static org.thingsboard.server.common.transport.service.DefaultTransportService.SUBSCRIBE_TO_RPC_ASYNC_MSG; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.NDEATH; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate.OFFLINE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.getTsKvProto; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopicPublish; /** @@ -794,12 +801,29 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement registerSubQoS(topic, grantedQoSList, reqQoS); } - public void processAttributesSubscribe(List grantedQoSList, String topic, MqttQoS reqQoS, TopicType topicType) { + private void processAttributesSubscribe(List grantedQoSList, String topic, MqttQoS reqQoS, TopicType topicType) { transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null); attrSubTopicType = topicType; registerSubQoS(topic, grantedQoSList, reqQoS); } + public void processAttributesRpcSubscribeSparkplugNode(List grantedQoSList, MqttQoS reqQoS) { + transportService.process(TransportProtos.TransportToDeviceActorMsg.newBuilder() + .setSessionInfo(deviceSessionCtx.getSessionInfo()) + .setSubscribeToAttributes(SUBSCRIBE_TO_ATTRIBUTE_UPDATES_ASYNC_MSG) + .setSubscribeToRPC(SUBSCRIBE_TO_RPC_ASYNC_MSG) + .build(), null); +// transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null); +// attrSubTopicType = TopicType.V1; +// mqttQoSMap.put(new MqttTopicMatcher(MqttTopics.DEVICE_ATTRIBUTES_TOPIC), getMinSupportedQos(reqQoS)); + registerSubQoS(MqttTopics.DEVICE_ATTRIBUTES_TOPIC, grantedQoSList, reqQoS); +// transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToRPCMsg.newBuilder().build(), null); +// rpcSubTopicType = TopicType.V2; +// mqttQoSMap.put(new MqttTopicMatcher(MqttTopics.DEVICE_RPC_REQUESTS_TOPIC), getMinSupportedQos(reqQoS)); +// mqttQoSMap.put(new MqttTopicMatcher(MqttTopics.GATEWAY_RPC_TOPIC), getMinSupportedQos(reqQoS)); +// grantedQoSList.add(getMinSupportedQos(reqQoS)); + } + public void registerSubQoS(String topic, List grantedQoSList, MqttQoS reqQoS) { grantedQoSList.add(getMinSupportedQos(reqQoS)); mqttQoSMap.put(new MqttTopicMatcher(topic), getMinSupportedQos(reqQoS)); @@ -1129,8 +1153,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement gatewaySessionHandler.onDevicesDisconnect(); } if (sparkplugSessionHandler != null) { - // add Msg Telemetry node: key STATE type: String value: OFFLINE ts: sparkplugBProto.getTimestamp() - sparkplugSessionHandler.stateSparkplugtSendOnTelemetry(deviceSessionCtx.getSessionInfo(), + // add Msg Telemetry node: key STATE type: String value: OFFLINE ts: sparkplugBProto.getTimestamp() + sparkplugSessionHandler.sendSparkplugStateOnTelemetry(deviceSessionCtx.getSessionInfo(), deviceSessionCtx.getDeviceInfo().getDeviceName(), OFFLINE, new Date().getTime()); sparkplugSessionHandler.onDevicesDisconnect(); } @@ -1139,7 +1163,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement deviceSessionCtx.release(); } - private void onValidateDeviceResponse(ValidateDeviceCredentialsResponse msg, ChannelHandlerContext ctx, MqttConnectMessage connectMessage) { if (!msg.hasDeviceInfo()) { context.onAuthFailure(address); @@ -1206,8 +1229,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg notification) { log.trace("[{}] Received attributes update notification to device", sessionId); - String topic = attrSubTopicType.getAttributesSubTopic(); - MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(attrSubTopicType); try { if (sparkplugSessionHandler != null) { log.trace("[{}] Received attributes update notification to sparkplug device", sessionId); @@ -1222,6 +1243,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } }); } else { + String topic = attrSubTopicType.getAttributesSubTopic(); + MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(attrSubTopicType); adaptor.convertToPublish(deviceSessionCtx, notification, topic).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); } } catch (Exception e) { @@ -1239,39 +1262,75 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { log.trace("[{}] Received RPC command to device", sessionId); - String baseTopic = rpcSubTopicType.getRpcRequestTopicBase(); - MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(rpcSubTopicType); try { - adaptor.convertToPublish(deviceSessionCtx, rpcRequest, baseTopic).ifPresent(payload -> { - int msgId = ((MqttPublishMessage) payload).variableHeader().packetId(); - if (isAckExpected(payload)) { - rpcAwaitingAck.put(msgId, rpcRequest); - context.getScheduler().schedule(() -> { - TransportProtos.ToDeviceRpcRequestMsg msg = rpcAwaitingAck.remove(msgId); - if (msg != null) { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); - } - }, Math.max(0, Math.min(deviceSessionCtx.getContext().getTimeout(), rpcRequest.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); + if (sparkplugSessionHandler != null) { + /** + * NCMD {"metricName":"MyNodeMetric05_String","value":"MyNodeMetric05_String_Value"} + * NCMD {"metricName":"MyNodeMetric02_LongInt64","value":2814119464032075444} + * NCMD {"metricName":"MyNodeMetric03_Double","value":6336935578763180333} + * NCMD {"metricName":"MyNodeMetric04_Float","value":413.18222} + * NCMD {"metricName":"Node Control/Rebirth","value":false} + * NCMD {"metricName":"MyNodeMetric06_Json_Bytes", "value":[40,47,-49]} + * NCMD {"metricName":"Node Control/Rebirth", "value":false} + * without backspace + */ + SparkplugMessageType messageType = SparkplugMessageType.parseMessageType(rpcRequest.getMethodName()); + if (messageType == null) { + this.sendErrorRpcResponse(deviceSessionCtx.getSessionInfo(), rpcRequest.getRequestId(), + ResponseCode.METHOD_NOT_ALLOWED, "Unsupported SparkplugMessageType: " + rpcRequest.getMethodName() + rpcRequest.getParams()); + return; } - var cf = publish(payload, deviceSessionCtx); - cf.addListener(result -> { - if (result.cause() == null) { - if (!isAckExpected(payload)) { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); - } else if (rpcRequest.getPersisted()) { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); - } - } else { - // TODO: send error - } - }); - }); + SparkplugRpcRequestHeader header = JacksonUtil.fromString(rpcRequest.getParams(), SparkplugRpcRequestHeader.class); + header.setMessageType(messageType.name()); + TransportProtos.TsKvProto tsKvProto = getTsKvProto(header.getMetricName(), header.getValue(), new Date().getTime()); + if (sparkplugSessionHandler.getNodeBirthMetrics().containsKey(tsKvProto.getKv().getKey())) { + SparkplugTopic sparkplugTopic = new SparkplugTopic(sparkplugSessionHandler.getSparkplugTopicNode(), + messageType); + sparkplugSessionHandler.createSparkplugMqttPublishMsg(tsKvProto, + sparkplugTopic.toString(), + sparkplugSessionHandler.getNodeBirthMetrics().get(tsKvProto.getKv().getKey())) + .ifPresent(payload -> sendToDeviceRpcRequest(payload, rpcRequest)); + } + } else { + String baseTopic = rpcSubTopicType.getRpcRequestTopicBase(); + MqttTransportAdaptor adaptor = deviceSessionCtx.getAdaptor(rpcSubTopicType); + adaptor.convertToPublish(deviceSessionCtx, rpcRequest, baseTopic) + .ifPresent(payload -> sendToDeviceRpcRequest(payload, rpcRequest)); + } } catch (Exception e) { - transportService.process(deviceSessionCtx.getSessionInfo(), - TransportProtos.ToDeviceRpcResponseMsg.newBuilder() - .setRequestId(rpcRequest.getRequestId()).setError("Failed to convert device RPC command to MQTT msg").build(), TransportServiceCallback.EMPTY); - log.trace("[{}] Failed to convert device RPC command to MQTT msg", sessionId, e); + log.trace("[{}] Failed to convert device RPC command to Sparkplug MQTT msg", sessionId, e); + this.sendErrorRpcResponse(deviceSessionCtx.getSessionInfo(), rpcRequest.getRequestId(), + ResponseCode.METHOD_NOT_ALLOWED, + "Failed to convert device RPC command to Sparkplug MQTT msg: " + rpcRequest.getMethodName() + rpcRequest.getParams()); + } + } + + public void sendToDeviceRpcRequest (MqttMessage payload, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { + int msgId = ((MqttPublishMessage) payload).variableHeader().packetId(); + if (isAckExpected(payload)) { + rpcAwaitingAck.put(msgId, rpcRequest); + context.getScheduler().schedule(() -> { + TransportProtos.ToDeviceRpcRequestMsg msg = rpcAwaitingAck.remove(msgId); + if (msg != null) { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); + } + }, Math.max(0, Math.min(deviceSessionCtx.getContext().getTimeout(), rpcRequest.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); } + var cf = publish(payload, deviceSessionCtx); + cf.addListener(result -> { + if (result.cause() == null) { + if (!isAckExpected(payload)) { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } else if (rpcRequest.getPersisted()) { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); + } + this.sendSuccessRpcResponse(deviceSessionCtx.getSessionInfo(), rpcRequest.getRequestId(), ResponseCode.CONTENT, "Success: " + rpcRequest.getMethodName()); + } else { + log.trace("[{}] Failed send To Device Rpc Request [{}]", sessionId, rpcRequest.getMethodName()); + this.sendErrorRpcResponse(deviceSessionCtx.getSessionInfo(), rpcRequest.getRequestId(), + ResponseCode.METHOD_NOT_ALLOWED, " Failed send To Device Rpc Request: " + rpcRequest.getMethodName()); + } + }); } @Override @@ -1311,4 +1370,16 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement ctx.close(); } + public void sendErrorRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, ResponseCode result, String errorMsg) { + String payload = JacksonUtil.toString(SparkplugRpcResponseBody.builder().result(result.getName()).error(errorMsg).build()); + TransportProtos.ToDeviceRpcResponseMsg msg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setError(payload).build(); + transportService.process(sessionInfo, msg, null); + } + + public void sendSuccessRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, ResponseCode result, String successMsg) { + String payload = JacksonUtil.toString(SparkplugRpcResponseBody.builder().result(result.getName()).result(successMsg).build()); + TransportProtos.ToDeviceRpcResponseMsg msg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setError(payload).build(); + transportService.process(sessionInfo, msg, null); + } + } 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 ba8036bb9d..13ffb960a5 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 @@ -729,7 +729,7 @@ public abstract class AbstractGatewaySessionHandler { private void deregisterSession(String deviceName, MqttDeviceAwareSessionContext deviceSessionCtx) { if (this.deviceSessionCtx.isSparkplug()){ // add Msg Telemetry: key STATE type: String value: OFFLINE ts: sparkplugBProto.getTimestamp() - stateSparkplugtSendOnTelemetry (deviceSessionCtx.getSessionInfo(), + sendSparkplugStateOnTelemetry(deviceSessionCtx.getSessionInfo(), deviceSessionCtx.getDeviceInfo().getDeviceName(), OFFLINE, new Date().getTime()); } transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); @@ -738,7 +738,7 @@ public abstract class AbstractGatewaySessionHandler { log.debug("[{}] Removed device [{}] from the gateway session", sessionId, deviceName); } - public void stateSparkplugtSendOnTelemetry (TransportProtos.SessionInfoProto sessionInfo, String deviceName, SparkplugMessageTypeSate typeSate, long ts) { + public void sendSparkplugStateOnTelemetry(TransportProtos.SessionInfoProto sessionInfo, String deviceName, SparkplugMessageTypeSate typeSate, long ts) { TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); keyValueProtoBuilder.setKey(STATE.name()); keyValueProtoBuilder.setType(TransportProtos.KeyValueType.STRING_V); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java index 06a398b9c1..e4db1f04f0 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java @@ -17,6 +17,8 @@ package org.thingsboard.server.transport.mqtt.session; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; +import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.gen.transport.TransportProtos; @@ -53,4 +55,40 @@ public class SparkplugDeviceSessionContext extends AbstractGatewayDeviceSessionC }); } + @Override + public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { + log.trace("[{}] Received RPC Request notification to sparkplug device", sessionId); + try { + /** + * NCMD {"metricName":"MyNodeMetric05_String","value":"MyNodeMetric05_String_Value"} + * NCMD {"metricName":"MyNodeMetric02_LongInt64","value":2814119464032075444} + * NCMD {"metricName":"MyNodeMetric03_Double","value":6336935578763180333} + * NCMD {"metricName":"MyNodeMetric04_Float","value":413.18222} + * NCMD {"metricName":"Node Control/Rebirth","value":false} + * NCMD {"metricName":"MyNodeMetric06_Json_Bytes", "value":[40,47,-49]} + * NCMD {"metricName":"Node Control/Rebirth", "value":false} + * without backspace + */ + SparkplugMessageType messageType = SparkplugMessageType.parseMessageType(rpcRequest.getMethodName()); +// if (messageType == null) { +// parent.sendErrorRpcResponse(parent.deviceSessionCtx., rpcRequest.getRequestId(), +// ResponseCode.METHOD_NOT_ALLOWED, "Unsupported SparkplugMessageType: " + rpcRequest.getMethodName() + rpcRequest.getParams()); +// return; +// } +// SparkplugRpcRequestHeader header = JacksonUtil.fromString(rpcRequest.getParams(), SparkplugRpcRequestHeader.class); +// header.setMessageType(messageType.name()); +// TransportProtos.TsKvProto tsKvProto = getTsKvProto(header.getMetricName(), header.getValue(), new Date().getTime()); +// if (sparkplugSessionHandler.getNodeBirthMetrics().containsKey(tsKvProto.getKv().getKey())) { +// SparkplugTopic sparkplugTopic = new SparkplugTopic(sparkplugSessionHandler.getSparkplugTopicNode(), +// messageType); +// sparkplugSessionHandler.createSparkplugMqttPublishMsg(tsKvProto, +// sparkplugTopic.toString(), +// sparkplugSessionHandler.getNodeBirthMetrics().get(tsKvProto.getKv().getKey())) +// .ifPresent(payload -> sendToDeviceRpcRequest(payload, rpcRequest)); +// } + } catch (ThingsboardException e) { + new ThingsboardException(e.getMessage(), ThingsboardErrorCode.INVALID_ARGUMENTS); + } + } + } 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 c52c5e33a8..b55170e70b 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 @@ -21,11 +21,12 @@ 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.MqttMessage; import io.netty.handler.codec.mqtt.MqttPublishMessage; import io.netty.handler.codec.mqtt.MqttQoS; import io.netty.handler.codec.mqtt.MqttTopicSubscription; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.data.device.profile.MqttTopics; +import org.eclipse.leshan.core.ResponseCode; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.transport.adaptor.AdaptorException; @@ -35,7 +36,6 @@ import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGateway import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; import org.thingsboard.server.transport.mqtt.MqttTransportHandler; -import org.thingsboard.server.transport.mqtt.TopicType; import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic; @@ -104,7 +104,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { if (topic.isType(NBIRTH) || topic.isType(DBIRTH)) { try { // add Msg Telemetry: key STATE type: String value: ONLINE ts: sparkplugBProto.getTimestamp() - stateSparkplugtSendOnTelemetry(contextListenableFuture.get().getSessionInfo(), deviceName, ONLINE, + sendSparkplugStateOnTelemetry(contextListenableFuture.get().getSessionInfo(), deviceName, ONLINE, sparkplugBProto.getTimestamp()); contextListenableFuture.get().setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); } catch (InterruptedException | ExecutionException e) { @@ -152,7 +152,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { // TODO SUBSCRIBE GroupId } else if (sparkplugTopic.isNode()) { // SUBSCRIBE Node - parent.processAttributesSubscribe(grantedQoSList, MqttTopics.DEVICE_ATTRIBUTES_TOPIC, reqQoS, TopicType.V1); + parent.processAttributesRpcSubscribeSparkplugNode(grantedQoSList, reqQoS); } else { // SUBSCRIBE Device - DO NOTHING, WE HAVE ALREADY SUBSCRIBED. // TODO: track that node subscribed to # or to particular device. @@ -238,4 +238,17 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { return new SparkplugDeviceSessionContext(this, msg.getDeviceInfo(), msg.getDeviceProfile(), mqttQoSMap, transportService); } + protected void sendToDeviceRpcRequest (MqttMessage payload, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { + parent.sendToDeviceRpcRequest(payload, rpcRequest); + } + + protected void sendErrorRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, ResponseCode result, String errorMsg) { + parent.sendErrorRpcResponse(sessionInfo, requestId, result, errorMsg); + } + + protected void sendSuccessRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, ResponseCode result, String successMsg) { + parent.sendSuccessRpcResponse(sessionInfo, requestId,result, successMsg); + } + + } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java index 19ec3156e6..6137433029 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java @@ -189,6 +189,41 @@ public class SparkplugMetricUtil { return metric; } + public static TransportProtos.TsKvProto getTsKvProto(String key, Object value, long ts) throws ThingsboardException { + try { + TransportProtos.TsKvProto.Builder tsKvProtoBuilder = TransportProtos.TsKvProto.newBuilder(); + TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); + keyValueProtoBuilder.setKey(key); + if (value instanceof String) { + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.STRING_V); + keyValueProtoBuilder.setStringV((String) value); + } else if (value instanceof Integer) { + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.LONG_V); + keyValueProtoBuilder.setLongV((Integer) value); + } else if (value instanceof Long) { + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.LONG_V); + keyValueProtoBuilder.setLongV((Long) value); + } else if (value instanceof Boolean) { + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.BOOLEAN_V); + keyValueProtoBuilder.setBoolV((Boolean) value); + } else if (value instanceof Double) { + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.DOUBLE_V); + keyValueProtoBuilder.setDoubleV((Double) value); + } else if (value instanceof List) { + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.JSON_V); + ArrayNode arrayNodeBytes = JacksonUtil.convertValue(value, ArrayNode.class); + keyValueProtoBuilder.setJsonV(arrayNodeBytes.toString()); + } else { + throw new ThingsboardException("Failed to convert device/node RPC command to TsKvProto for Sparkplug MQT msg: value [" + value + "]", ThingsboardErrorCode.INVALID_ARGUMENTS); + } + tsKvProtoBuilder.setKv(keyValueProtoBuilder.build()); + tsKvProtoBuilder.setTs(ts); + return tsKvProtoBuilder.build(); + } catch (Exception e) { + throw new ThingsboardException("Failed to convert device/node RPC command to TsKvProto for Sparkplug MQT msg: value [" + value + "]", ThingsboardErrorCode.INVALID_ARGUMENTS); + } + } + public static Optional validatedValueByTypeMetric(TransportProtos.KeyValueProto kv, MetricDataType metricDataType) throws ThingsboardException { if (kv.getTypeValue() <= 3) { return validatedValuePrimitiveByTypeMetric(kv, metricDataType); @@ -214,8 +249,8 @@ public class SparkplugMetricUtil { case UInt8: case UInt16: case Int32: - Optional boolInt8 = booleanStringToInt (valueOpt.get()); - if(boolInt8.isPresent()) { + Optional boolInt8 = booleanStringToInt(valueOpt.get()); + if (boolInt8.isPresent()) { return Optional.of(boolInt8.get()); } try { @@ -233,25 +268,25 @@ public class SparkplugMetricUtil { case Int64: case UInt64: case DateTime: - Optional boolInt64 = booleanStringToInt (valueOpt.get()); - if(boolInt64.isPresent()) { + Optional boolInt64 = booleanStringToInt(valueOpt.get()); + if (boolInt64.isPresent()) { return Optional.of(Long.valueOf(boolInt64.get())); } var l = new BigDecimal(valueOpt.get()); return Optional.of(l.longValue()); - // float + // float case Float: - Optional boolFloat = booleanStringToInt (valueOpt.get()); - if(boolFloat.isPresent()) { + Optional boolFloat = booleanStringToInt(valueOpt.get()); + if (boolFloat.isPresent()) { var fb = new BigDecimal(boolFloat.get()); return Optional.of(fb.floatValue()); } var f = new BigDecimal(valueOpt.get()); return Optional.of(f.floatValue()); - // double + // double case Double: - Optional boolDouble = booleanStringToInt (valueOpt.get()); - if(boolDouble.isPresent()) { + Optional boolDouble = booleanStringToInt(valueOpt.get()); + if (boolDouble.isPresent()) { return Optional.of(Double.valueOf(boolDouble.get())); } var dd = new BigDecimal(valueOpt.get()); @@ -272,7 +307,7 @@ public class SparkplugMetricUtil { case String: case Text: case UUID: - return Optional.of(valueOpt.get()); + return Optional.of(valueOpt.get()); } } catch (Exception e) { log.trace("Invalid type value [{}] for MetricDataType [{}] [{}]", kv, metricDataType.name(), e.getMessage()); @@ -282,15 +317,16 @@ public class SparkplugMetricUtil { return Optional.empty(); } - public static Optional validatedValueJsonByTypeMetric(String arrayNodeStr, MetricDataType metricDataType) { + public static Optional validatedValueJsonByTypeMetric(String arrayNodeStr, MetricDataType metricDataType) { try { Optional valueOpt; switch (metricDataType) { // byte[] case Bytes: - List listBytes = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() {}); + List listBytes = JacksonUtil.fromString(arrayNodeStr, new TypeReference<>() { + }); byte[] bytes = new byte[listBytes.size()]; - for(int i = 0; i < listBytes.size(); i++) { + for (int i = 0; i < listBytes.size(); i++) { bytes[i] = listBytes.get(i).byteValue(); } return Optional.of(bytes); @@ -324,7 +360,7 @@ public class SparkplugMetricUtil { } } - private static Optional booleanStringToInt (String booleanStr) { + private static Optional booleanStringToInt(String booleanStr) { if ("true".equals(booleanStr)) { return Optional.of(1); } else if ("false".equals(booleanStr)) { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcRequestHeader.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcRequestHeader.java new file mode 100644 index 0000000000..47b7cb6dd2 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcRequestHeader.java @@ -0,0 +1,29 @@ +/** + * 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 lombok.Data; + +@Data +@JsonIgnoreProperties(ignoreUnknown = true) +public class SparkplugRpcRequestHeader { + + private String messageType; + private String metricName; + private Object value; + +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcResponseBody.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcResponseBody.java new file mode 100644 index 0000000000..d120190588 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugRpcResponseBody.java @@ -0,0 +1,31 @@ +/** + * 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.JsonInclude; +import lombok.Builder; +import lombok.Data; + +@Data +@Builder +@JsonInclude(JsonInclude.Include.NON_NULL) +public class SparkplugRpcResponseBody { + + private String result; + private String value; + private String error; + +}