From 278935cdc98eba7d4ee97997ac7f0870f7a3c7a1 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Sun, 15 Jan 2023 19:53:00 +0200 Subject: [PATCH] sparkplug: Telemetry DBIRTH, NBIRH... --- .../transport/mqtt/MqttTransportHandler.java | 21 ++++++++++---- .../session/SparkplugNodeSessionHandler.java | 29 ++++++++++++++----- 2 files changed, 37 insertions(+), 13 deletions(-) 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 8f8e760123..24012e1e7a 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -17,6 +17,7 @@ package org.thingsboard.server.transport.mqtt; import com.fasterxml.jackson.databind.JsonNode; import com.google.gson.JsonParseException; +import com.google.protobuf.InvalidProtocolBufferException; import io.netty.channel.ChannelFuture; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; @@ -69,8 +70,10 @@ import org.thingsboard.server.common.transport.util.SslUtil; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509CertRequestMsg; +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; @@ -388,7 +391,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement case NBIRTH: case NCMD: case NDATA: - sparkplugSessionHandler.onDeviceTelemetryProto(msgId, mqttMsg.payload(), deviceName, sparkplugTopic.isNode()); + 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); @@ -406,10 +410,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement break; case DCMD: case DDATA: - sparkplugSessionHandler.onDeviceTelemetryProto(msgId, mqttMsg.payload(), deviceName, sparkplugTopic.isNode()); - break; case DBIRTH: - sparkplugSessionHandler.onDeviceConnectProto(mqttMsg, deviceSessionCtx.getDeviceInfo().getDeviceType()); + 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); @@ -424,7 +427,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement log.error("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); ack(ctx, msgId, ReturnCode.IMPLEMENTATION_SPECIFIC); ctx.close(); - } catch (AdaptorException | ThingsboardException e) { + } catch (AdaptorException | ThingsboardException | InvalidProtocolBufferException e) { log.error("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); sendAckOrCloseSession(ctx, topicName, msgId); } @@ -1055,9 +1058,15 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void checkSparkplugSession(MqttConnectMessage connectMessage) { try { - SparkplugTopic sparkplugTopic = parseTopicPublish(connectMessage.payload().willTopic()); if (sparkplugSessionHandler == null) { sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId); + if (StringUtils.isNotBlank(connectMessage.payload().willTopic()) + && connectMessage.payload().willMessageInBytes() != null && connectMessage.payload().willMessageInBytes().length > 0) { + SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); + SparkplugTopic sparkplugTopic = parseTopicPublish(connectMessage.payload().willTopic()); + sparkplugSessionHandler.onDeviceTelemetryProto(0, sparkplugBProtoNode, + deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopic.getType().name(), true); + } } } catch (Exception e) { log.trace("[{}][{}] Failed to fetch sparkplugDevice additional info or sparkplugTopicName", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName(), e); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java index 740f8ed292..547139ebfb 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java @@ -22,7 +22,6 @@ import com.google.gson.JsonParser; import com.google.gson.JsonSyntaxException; import com.google.protobuf.Descriptors; import com.google.protobuf.InvalidProtocolBufferException; -import io.netty.buffer.ByteBuf; import io.netty.handler.codec.mqtt.MqttPublishMessage; import io.netty.handler.codec.mqtt.MqttPublishVariableHeader; import io.netty.handler.codec.mqtt.MqttQoS; @@ -44,6 +43,7 @@ import java.util.List; import java.util.Optional; import java.util.UUID; +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; @@ -70,11 +70,10 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { } } - public void onDeviceTelemetryProto(int msgId, ByteBuf payload, String deviceName, boolean isNode) throws AdaptorException { + public void onDeviceTelemetryProto(int msgId, SparkplugBProto.Payload sparkplugBProto, String deviceName, String topicTypeName, boolean isNode) throws AdaptorException { try { checkDeviceName(deviceName); - SparkplugBProto.Payload sparkplugBProto = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(payload)); - List msgs = convertToPostTelemetry(sparkplugBProto); + List msgs = convertToPostTelemetry(sparkplugBProto, topicTypeName); int finalMsgId = msgId; ListenableFuture contextListenableFuture = isNode ? Futures.immediateFuture(this.deviceSessionCtx) : checkDeviceConnected(deviceName); @@ -97,7 +96,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { } }, context.getExecutor()); } - } catch (RuntimeException | InvalidProtocolBufferException e) { + } catch (RuntimeException e) { throw new AdaptorException(e); } } @@ -156,12 +155,14 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { } } - private List convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto) throws AdaptorException { + private List convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto, String topicTypeName) throws AdaptorException { try { List msgs = new ArrayList<>(); for (SparkplugBProto.Payload.Metric protoMetric : sparkplugBProto.getMetricsList()) { long ts = protoMetric.getTimestamp(); - Optional keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(protoMetric.getName(), protoMetric); + String keys = "bdSeq".equals(protoMetric.getName()) ? + topicTypeName + " " + protoMetric.getName() : protoMetric.getName(); + Optional keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(keys, protoMetric); if (keyValueProtoOpt.isPresent()) { List result = new ArrayList<>(); result.add(keyValueProtoOpt.get()); @@ -173,6 +174,20 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { msgs.add(request.build()); } } + if (DBIRTH.name().equals(topicTypeName)) { + List result = new ArrayList<>(); + TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); + keyValueProtoBuilder.setKey(topicTypeName + " " + "seq"); + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.LONG_V); + keyValueProtoBuilder.setLongV(sparkplugBProto.getSeq()); + result.add(keyValueProtoBuilder.build()); + TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); + TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); + builder.setTs(sparkplugBProto.getTimestamp()); + builder.addAllKv(result); + request.addTsKvList(builder.build()); + msgs.add(request.build()); + } return msgs; } catch (IllegalStateException | JsonSyntaxException | ThingsboardException e) { log.error("Failed to decode post telemetry request", e);