From 587e161398ce9938c3338f4395210f80c1675e64 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Thu, 12 Jan 2023 15:32:18 +0200 Subject: [PATCH] sparkplug: Telemetry --- .../transport/mqtt/MqttTransportHandler.java | 86 ++++-- .../AbstractGatewaySessionHandler.java | 24 +- .../session/SparkplugNodeSessionHandler.java | 174 ++++++++----- .../mqtt/util/sparkplug/MetricDataType.java | 199 ++++++++++++++ .../util/sparkplug/SparkplugMetricUtil.java | 244 ++++++++++++++++++ .../util/sparkplug/SparkplugTopicUtil.java | 147 +++++------ 6 files changed, 700 insertions(+), 174 deletions(-) create mode 100644 common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java create mode 100644 common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 166eed3347..8f8e760123 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 @@ -49,6 +49,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; @@ -68,16 +69,15 @@ 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.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 +97,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 +122,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 +323,13 @@ 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 if (sparkplugSessionHandler != null) { + handleSparkplugPublishMsg(ctx, topicName, mqttMsg); } else { processDevicePublish(ctx, mqttMsg, topicName, msgId); } @@ -375,14 +372,60 @@ 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 mqttMsgOld) { + MqttPublishMessage mqttMsg = sparkplugSessionHandler.reCreateMqttPublishMessageWithPacketId(mqttMsgOld); + 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: + sparkplugSessionHandler.onDeviceTelemetryProto(msgId, mqttMsg.payload(), deviceName, 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: + sparkplugSessionHandler.onDeviceTelemetryProto(msgId, mqttMsg.payload(), deviceName, sparkplugTopic.isNode()); + break; + case DBIRTH: + sparkplugSessionHandler.onDeviceConnectProto(mqttMsg, deviceSessionCtx.getDeviceInfo().getDeviceType()); + 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 e) { + log.error("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); sendAckOrCloseSession(ctx, topicName, msgId); } } @@ -648,7 +691,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 +1055,9 @@ 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()); - // + SparkplugTopic sparkplugTopic = parseTopicPublish(connectMessage.payload().willTopic()); if (sparkplugSessionHandler == null) { - sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId, sparkplugTopic.toString()); - } else { - log.warn("SparkPlugNodeReConnected [{}] [{}]", sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); + sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId); } } catch (Exception e) { log.trace("[{}][{}] Failed to fetch sparkplugDevice additional info or sparkplugTopicName", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName(), e); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index e330e3d5db..9ea5cbb00f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java @@ -79,20 +79,20 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe @Slf4j public abstract class AbstractGatewaySessionHandler { - private static final String DEFAULT_DEVICE_TYPE = "default"; + protected static final String DEFAULT_DEVICE_TYPE = "default"; private static final String CAN_T_PARSE_VALUE = "Can't parse value: "; private static final String DEVICE_PROPERTY = "device"; - private final MqttTransportContext context; + protected final MqttTransportContext context; private final TransportService transportService; - private final TransportDeviceInfo gateway; - private final UUID sessionId; + protected final TransportDeviceInfo gateway; + protected final UUID sessionId; private final ConcurrentMap deviceCreationLockMap; private final ConcurrentMap devices; private final ConcurrentMap> deviceFutures; private final ConcurrentMap mqttQoSMap; - private final ChannelHandlerContext channel; - private final DeviceSessionCtx deviceSessionCtx; + protected final ChannelHandlerContext channel; + protected final DeviceSessionCtx deviceSessionCtx; public AbstractGatewaySessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId) { this.context = deviceSessionCtx.getContext(); @@ -198,7 +198,7 @@ public abstract class AbstractGatewaySessionHandler { return deviceSessionCtx.isJsonPayloadType(); } - private void processOnConnect(MqttPublishMessage msg, String deviceName, String deviceType) { + protected void processOnConnect(MqttPublishMessage msg, String deviceName, String deviceType) { log.trace("[{}] onDeviceConnect: {}", sessionId, deviceName); Futures.addCallback(onDeviceConnect(deviceName, deviceType), new FutureCallback() { @Override @@ -393,7 +393,7 @@ public abstract class AbstractGatewaySessionHandler { } } - private void processPostTelemetryMsg(MqttDeviceAwareSessionContext deviceCtx, TransportProtos.PostTelemetryMsg postTelemetryMsg, String deviceName, int msgId) { + protected void processPostTelemetryMsg(MqttDeviceAwareSessionContext deviceCtx, TransportProtos.PostTelemetryMsg postTelemetryMsg, String deviceName, int msgId) { transportService.process(deviceCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(channel, deviceName, msgId, postTelemetryMsg)); } @@ -666,7 +666,7 @@ public abstract class AbstractGatewaySessionHandler { return result.build(); } - private ListenableFuture checkDeviceConnected(String deviceName) { + protected ListenableFuture checkDeviceConnected(String deviceName) { MqttDeviceAwareSessionContext ctx = devices.get(deviceName); if (ctx == null) { log.debug("[{}] Missing device [{}] for the gateway session", sessionId, deviceName); @@ -676,7 +676,7 @@ public abstract class AbstractGatewaySessionHandler { } } - private String checkDeviceName(String deviceName) { + protected String checkDeviceName(String deviceName) { if (StringUtils.isEmpty(deviceName)) { throw new RuntimeException("Device name is empty!"); } else { @@ -697,11 +697,11 @@ public abstract class AbstractGatewaySessionHandler { return JsonMqttAdaptor.validateJsonPayload(sessionId, mqttMsg.payload()); } - private byte[] getBytes(ByteBuf payload) { + protected byte[] getBytes(ByteBuf payload) { return ProtoMqttAdaptor.toBytes(payload); } - private void ack(MqttPublishMessage msg, ReturnCode returnCode) { + protected void ack(MqttPublishMessage msg, ReturnCode returnCode) { int msgId = getMsgId(msg); if (msgId > 0) { writeAndFlush(MqttTransportHandler.createMqttPubAckMsg(deviceSessionCtx, msgId, returnCode)); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java index 9594b55515..740f8ed292 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java @@ -15,93 +15,99 @@ */ 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 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; 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.adaptors.ProtoMqttAdaptor; 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.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, ByteBuf payload, String deviceName, boolean isNode) throws AdaptorException { + try { + checkDeviceName(deviceName); + SparkplugBProto.Payload sparkplugBProto = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(payload)); + List msgs = convertToPostTelemetry(sparkplugBProto); + int finalMsgId = msgId; + ListenableFuture contextListenableFuture = isNode ? + Futures.immediateFuture(this.deviceSessionCtx) : checkDeviceConnected(deviceName); + for (TransportProtos.PostTelemetryMsg msg : msgs) { + Futures.addCallback(contextListenableFuture, + new FutureCallback<>() { + @Override + public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + try { + processPostTelemetryMsg(deviceCtx, msg, deviceName, finalMsgId); + } catch (Throwable e) { + log.warn("[{}][{}] Failed to convert telemetry: {}", gateway.getDeviceId(), deviceName, msg, e); + channel.close(); + } + } + + @Override + public void onFailure(Throwable t) { + log.debug("[{}] Failed to process device telemetry command: {}", sessionId, deviceName, t); + } + }, context.getExecutor()); } + } catch (RuntimeException | InvalidProtocolBufferException e) { + throw new AdaptorException(e); } } public void handleSparkplugSubscribeMsg(List grantedQoSList, SparkplugTopic sparkplugTopic, MqttQoS reqQoS) { - String topicName = sparkplugTopic.toString(); - log.warn("SparkplugSubscribeMsg [{}] [{}]", sparkplugTopic.isNode() ? "node" : "device: " + sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); - if (sparkplugTopic.getGroupId() == null) { // TODO SUBSCRIBE NameSpace } else if (sparkplugTopic.getType() == null) { // TODO SUBSCRIBE GroupId - } - else if (sparkplugTopic.isNode()) { + } else if (sparkplugTopic.isNode()) { // A node topic switch (sparkplugTopic.getType()) { case STATE: @@ -150,4 +156,52 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler{ } } + private List convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto) throws AdaptorException { + try { + List msgs = new ArrayList<>(); + for (SparkplugBProto.Payload.Metric protoMetric : sparkplugBProto.getMetricsList()) { + long ts = protoMetric.getTimestamp(); + Optional keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(protoMetric.getName(), protoMetric); + if (keyValueProtoOpt.isPresent()) { + List result = new ArrayList<>(); + result.add(keyValueProtoOpt.get()); + TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); + TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); + builder.setTs(ts); + builder.addAllKv(result); + request.addTsKvList(builder.build()); + msgs.add(request.build()); + } + } + return msgs; + } catch (IllegalStateException | JsonSyntaxException | ThingsboardException e) { + log.error("Failed to decode post telemetry request", e); + throw new AdaptorException(e); + } + } + + public MqttPublishMessage reCreateMqttPublishMessageWithPacketId(MqttPublishMessage mqttMsgOld) { + try { + SparkplugBProto.Payload sparkplugBProto = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(mqttMsgOld.payload())); + MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader(mqttMsgOld.variableHeader().topicName(), (int) sparkplugBProto.getSeq()); + return new MqttPublishMessage(mqttMsgOld.fixedHeader(), variableHeader, mqttMsgOld.payload()); + } catch (InvalidProtocolBufferException e) { + log.error("Failed to deserialize SparkplugBProto.Payload", e); + throw new RuntimeException("Failed to deserialize SparkplugBProto.Payload"); + } + } + + public void onDeviceConnectProto(MqttPublishMessage mqttPublishMessage, String nodeDeviceType) throws ThingsboardException { + try { + String topic = mqttPublishMessage.variableHeader().topicName(); + SparkplugTopic sparkplugTopic = parseTopicSubscribe(topic); + String deviceName = checkDeviceName(sparkplugTopic.getDeviceId()); + String deviceType = StringUtils.isEmpty(nodeDeviceType) ? DEFAULT_DEVICE_TYPE : nodeDeviceType; + processOnConnect(mqttPublishMessage, deviceName, deviceType); + } catch (RuntimeException | ThingsboardException e) { + log.error("Failed Sparkplug Device connect proto!", e); + throw new ThingsboardException(e, ThingsboardErrorCode.BAD_REQUEST_PARAMS); + } + } + } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java new file mode 100644 index 0000000000..f4dc74f46b --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java @@ -0,0 +1,199 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.util.sparkplug; + +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.transport.adaptor.AdaptorException; +import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; + +import java.math.BigInteger; +import java.util.Date; + +/** + * Created by nickAS21 on 10.01.23 + */ + +@Slf4j +public enum MetricDataType { + + // Basic Types + Int8(1, Byte.class), + Int16(2, Short.class), + Int32(3, Integer.class), + Int64(4, Long.class), + UInt8(5, Short.class), + UInt16(6, Integer.class), + UInt32(7, Long.class), + UInt64(8, BigInteger.class), + Float(9, Float.class), + Double(10, Double.class), + Boolean(11, Boolean.class), + String(12, String.class), + DateTime(13, Date.class), + Text(14, String.class), + + // Custom Types for Metrics + UUID(15, String.class), + DataSet(16, SparkplugBProto.Payload.DataSet.class), + Bytes(17, byte[].class), + File(18, SparkplugMetricUtil.File.class), + Template(19, SparkplugBProto.Payload.Template.class), + + // PropertyValue Types (20 and 21) are NOT metric datatypes + + // Array Types + Int8Array(22, Byte[].class), + Int16Array(23, Short[].class), + Int32Array(24, Integer[].class), + Int64Array(25, Long[].class), + UInt8Array(26, Short[].class), + UInt16Array(27, Integer[].class), + UInt32Array(28, Long[].class), + UInt64Array(29, BigInteger[].class), + FloatArray(30, Float[].class), + DoubleArray(31, Double[].class), + BooleanArray(32, Boolean[].class), + StringArray(33, String[].class), + DateTimeArray(34, Date[].class), + + // Unknown + Unknown(0, Object.class); + + private Class clazz = null; + private int intValue = 0; + + /** + * Constructor + * + * @param intValue the integer value of this {@link MetricDataType} + * @param clazz the {@link Class} type associated with this {@link MetricDataType} + */ + private MetricDataType(int intValue, Class clazz) { + this.intValue = intValue; + this.clazz = clazz; + } + + /** + * Checks the type of a specified value against the specified {@link MetricDataType} + * + * @param value the {@link Object} value to check against the {@link MetricDataType} + * @throws AdaptorException if the value is not a valid type for the given {@link MetricDataType} + */ + public void checkType(Object value) throws AdaptorException { + if (value != null && !clazz.isAssignableFrom(value.getClass())) { + String msgError = "Failed type check - " + clazz + " != " + ((value != null) ? value.getClass().toString() : "null"); + log.debug(msgError); + throw new AdaptorException(msgError); + } + } + + /** + * Returns an integer representation of the data type. + * + * @return an integer representation of the data type. + */ + public int toIntValue() { + return this.intValue; + } + + /** + * Converts the integer representation of the data type into a {@link MetricDataType} instance. + * + * @param i the integer representation of the data type. + * @return a {@link MetricDataType} instance. + */ + public static MetricDataType fromInteger(int i) { + switch (i) { + case 1: + return Int8; + case 2: + return Int16; + case 3: + return Int32; + case 4: + return Int64; + case 5: + return UInt8; + case 6: + return UInt16; + case 7: + return UInt32; + case 8: + return UInt64; + case 9: + return Float; + case 10: + return Double; + case 11: + return Boolean; + case 12: + return String; + case 13: + return DateTime; + case 14: + return Text; + case 15: + return UUID; + case 16: + return DataSet; + case 17: + return Bytes; + case 18: + return File; + case 19: + return Template; + case 22: + return Int8Array; + case 23: + return Int16Array; + case 24: + return Int32Array; + case 25: + return Int64Array; + case 26: + return UInt8Array; + case 27: + return UInt16Array; + case 28: + return UInt32Array; + case 29: + return UInt64Array; + case 30: + return FloatArray; + case 31: + return DoubleArray; + case 32: + return BooleanArray; + case 33: + return StringArray; + case 34: + return DateTimeArray; + default: + return Unknown; + } + } + + /** + * Returns the class type for this DataType + * + * @return the class type for this DataType + */ + public Class getClazz() { + return clazz; + } + + +} \ No newline at end of file diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java new file mode 100644 index 0000000000..8019c35774 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java @@ -0,0 +1,244 @@ +/** + * 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.ser.std.FileSerializer; +import com.google.gson.Gson; +import com.google.gson.GsonBuilder; +import com.google.gson.JsonArray; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.codec.binary.Hex; +import org.apache.commons.lang3.StringUtils; +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.math.BigInteger; +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.Optional; + +/** + * Provides utility methods for SparkplugB MQTT Payload Metric. + */ +@Slf4j +public class SparkplugMetricUtil { + + public static Optional getFromSparkplugBMetricToKeyValueProto(String key, SparkplugBProto.Payload.Metric protoMetric) throws ThingsboardException { + // Check if the null flag has been set indicating that the value is null + if (protoMetric.getIsNull()) { + return null; + } + // Otherwise convert the value based on the type + int metricType = protoMetric.getDatatype(); + switch (MetricDataType.fromInteger(metricType)) { + case Boolean: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.BOOLEAN_V) + .setBoolV(protoMetric.getBooleanValue()).build()); + case DateTime: + case Int64: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(protoMetric.getLongValue()).build()); + case File: + String filename = protoMetric.getMetadata().getFileName(); + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key + "_" + filename).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(Hex.encodeHexString((protoMetric.getBytesValue().toByteArray()))).build()); + case Float: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV((long) protoMetric.getFloatValue()).build()); + case Double: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.DOUBLE_V) + .setDoubleV(protoMetric.getDoubleValue()).build()); + case Int8: + case Int16: + case Int32: + case UInt16: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(protoMetric.getIntValue()).build()); + case UInt32: + if (protoMetric.hasIntValue()) { + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(protoMetric.getIntValue()).build()); + } else if (protoMetric.hasLongValue()) { + return Optional.of(TransportProtos.KeyValueProto.newBuilder().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 UInt32 datatype " + metricType, ThingsboardErrorCode.INVALID_ARGUMENTS); + } + case UInt64: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.DOUBLE_V) + .setDoubleV((new BigInteger(Long.toUnsignedString(protoMetric.getLongValue()))).longValue()).build()); + case String: + case Text: + case UUID: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(protoMetric.getStringValue()).build()); + case Bytes: + case Int8Array: + case Int16Array: + case Int32Array: + case Int64Array: + case UInt8Array: + case UInt16Array: + case UInt32Array: + case UInt64Array: + case FloatArray: + case DoubleArray: + case BooleanArray: + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(Hex.encodeHexString(protoMetric.getBytesValue().toByteArray())).build()); + case DataSet: + SparkplugBProto.Payload.DataSet protoDataSet = protoMetric.getDatasetValue(); + //TODO + // Build the and create the DataSet + /** + 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(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(protoDataSet.toString()).build()); + case Template: + //TODO + // Build the and create the Template + SparkplugBProto.Payload.Template protoTemplate = protoMetric.getTemplateValue(); + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV( protoTemplate.toString()).build()); + case StringArray: + ByteBuffer stringByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()); + List stringList = new ArrayList<>(); + stringByteBuffer.order(ByteOrder.LITTLE_ENDIAN); + StringBuilder sb = new StringBuilder(); + while (stringByteBuffer.hasRemaining()) { + byte b = stringByteBuffer.get(); + if (b == (byte) 0) { + stringList.add(sb.toString()); + sb = new StringBuilder(); + } else { + + sb.append((char) b); + } + } + String st = StringUtils.join(stringList, "|"); + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(st).build()); + case DateTimeArray: + ByteBuffer dateTimeByteBuffer = ByteBuffer.wrap(protoMetric.getBytesValue().toByteArray()); + List dateTimeList = new ArrayList(); + dateTimeByteBuffer.order(ByteOrder.LITTLE_ENDIAN); + while (dateTimeByteBuffer.hasRemaining()) { + long longValue = dateTimeByteBuffer.getLong(); + dateTimeList.add(longValue); + } + Gson gson = new GsonBuilder().create(); + JsonArray dateTimeArray = gson.toJsonTree(dateTimeList).getAsJsonArray(); + return Optional.of(TransportProtos.KeyValueProto.newBuilder().setKey(key).setType(TransportProtos.KeyValueType.JSON_V) + .setStringV(dateTimeArray.toString()).build()); + case Unknown: + default: + throw new ThingsboardException("Failed to decode: Unknown MetricDataType " + metricType, ThingsboardErrorCode.INVALID_ARGUMENTS); + } + } + + + @JsonIgnoreProperties( + value = {"fileName"}) + @JsonSerialize( + using = FileSerializer.class) + public class File { + + private String fileName; + private byte[] bytes; + + /** + * Default Constructor + */ + public File() { + super(); + } + + /** + * Constructor + * + * @param fileName the full file name path + * @param bytes the array of bytes that represent the contents of the file + */ + public File(String fileName, byte[] bytes) { + super(); + this.fileName = fileName == null + ? null + : fileName.replace("/", System.getProperty("file.separator")).replace("\\", + System.getProperty("file.separator")); + this.bytes = Arrays.copyOf(bytes, bytes.length); + } + + /** + * Gets the full filename path + * + * @return the full filename path + */ + public String getFileName() { + return fileName; + } + + /** + * Sets the full filename path + * + * @param fileName the full filename path + */ + public void setFileName(String fileName) { + this.fileName = fileName; + } + + /** + * Gets the bytes that represent the contents of the file + * + * @return the bytes that represent the contents of the file + */ + public byte[] getBytes() { + return bytes; + } + + /** + * Sets the bytes that represent the contents of the file + * + * @param bytes the bytes that represent the contents of the file + */ + public void setBytes(byte[] bytes) { + this.bytes = bytes; + } + + @Override + public String toString() { + StringBuilder builder = new StringBuilder(); + builder.append("File [fileName="); + builder.append(fileName); + builder.append(", bytes="); + builder.append(Arrays.toString(bytes)); + builder.append("]"); + return builder.toString(); + } + } + +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java index c2e93cb269..93d0245da6 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java @@ -27,88 +27,79 @@ import java.util.Map; * Provides utility methods for handling Sparkplug MQTT message topics. */ public class SparkplugTopicUtil { - - private static final Map SPLIT_TOPIC_CACHE = new HashMap(); - - public static String[] getSplitTopic(String topic) { - String[] splitTopic = SPLIT_TOPIC_CACHE.get(topic); - if (splitTopic == null) { - splitTopic = topic.split("/"); - SPLIT_TOPIC_CACHE.put(topic, splitTopic); - } - - return splitTopic; - } - /** - * Serializes a {@link SparkplugTopic} instance in to a JSON string. - * - * @param topic a {@link SparkplugTopic} instance - * @return a JSON string - * @throws JsonProcessingException - */ - public static String sparkplugTopicToString(SparkplugTopic topic) throws JsonProcessingException { - ObjectMapper mapper = new ObjectMapper(); - return mapper.writeValueAsString(topic); - } + private static final Map SPLIT_TOPIC_CACHE = new HashMap(); + private static final String TOPIC_INVALID_NUMBER = "Invalid number of topic elements: "; - /** - * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. - * - * @param topic a topic string - * @return a {@link SparkplugTopic} instance - * @throws ThingsboardException if an error occurs while parsing - */ - public static SparkplugTopic parseTopic(String topic) throws ThingsboardException { - topic = topic.indexOf("#") > 0 ? topic.substring(0, topic.indexOf("#")) : topic; - return parseTopic(SparkplugTopicUtil.getSplitTopic(topic)); - } + public static String[] getSplitTopic(String topic) { + String[] splitTopic = SPLIT_TOPIC_CACHE.get(topic); + if (splitTopic == null) { + splitTopic = topic.split("/"); + SPLIT_TOPIC_CACHE.put(topic, splitTopic); + } - /** - * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. - * - * @param splitTopic a topic split into tokens - * @return a {@link SparkplugTopic} instance - * @throws Exception if an error occurs while parsing - */ - @SuppressWarnings("incomplete-switch") - public static SparkplugTopic parseTopic(String[] splitTopic) throws ThingsboardException { - SparkplugMessageType type; - String namespace, edgeNodeId, groupId; - int length = splitTopic.length; + return splitTopic; + } - if (length < 4 || length > 5) { - throw new ThingsboardException("Invalid number of topic elements: " + length, ThingsboardErrorCode.INVALID_ARGUMENTS); - } + /** + * Serializes a {@link SparkplugTopic} instance in to a JSON string. + * + * @param topic a {@link SparkplugTopic} instance + * @return a JSON string + * @throws JsonProcessingException + */ + public static String sparkplugTopicToString(SparkplugTopic topic) throws JsonProcessingException { + ObjectMapper mapper = new ObjectMapper(); + return mapper.writeValueAsString(topic); + } - namespace = splitTopic[0]; - groupId = splitTopic[1]; - type = SparkplugMessageType.parseMessageType(splitTopic[2]); - edgeNodeId = splitTopic[3]; + /** + * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. + * + * @param topic a topic string + * @return a {@link SparkplugTopic} instance + * @throws ThingsboardException if an error occurs while parsing + */ + public static SparkplugTopic parseTopicSubscribe(String topic) throws ThingsboardException { + // TODO "+", "$" + topic = topic.indexOf("#") > 0 ? topic.substring(0, topic.indexOf("#")) : topic; + return parseTopic(SparkplugTopicUtil.getSplitTopic(topic)); + } + + public static SparkplugTopic parseTopicPublish(String topic) throws ThingsboardException { + if (topic.contains("#") || topic.contains("$") || topic.contains("+")) { + throw new ThingsboardException("Invalid of topic elements for Publish", ThingsboardErrorCode.INVALID_ARGUMENTS); + } else { + String[] splitTopic = SparkplugTopicUtil.getSplitTopic(topic); + if (splitTopic.length < 4 || splitTopic.length > 5) { + throw new ThingsboardException(TOPIC_INVALID_NUMBER + splitTopic.length, ThingsboardErrorCode.INVALID_ARGUMENTS); + } + return parseTopic(splitTopic); + } + } + + /** + * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. + * + * @param splitTopic a topic split into tokens + * @return a {@link SparkplugTopic} instance + * @throws Exception if an error occurs while parsing + */ + @SuppressWarnings("incomplete-switch") + public static SparkplugTopic parseTopic(String[] splitTopic) throws ThingsboardException { + int length = splitTopic.length; + if (length == 0) { + throw new ThingsboardException(TOPIC_INVALID_NUMBER + length, ThingsboardErrorCode.INVALID_ARGUMENTS); + } else { + SparkplugMessageType type; + String namespace, edgeNodeId, groupId, deviceId; + namespace = splitTopic[0]; + groupId = length > 1 ? splitTopic[1] : null; + type = length > 2 ? SparkplugMessageType.parseMessageType(splitTopic[2]) : null; + edgeNodeId = length > 3 ? splitTopic[3] : null; + deviceId = length > 4 ? splitTopic[4] : null; + return new SparkplugTopic(namespace, groupId, edgeNodeId, deviceId, type); + } + } - if (length == 4) { - // A node topic - switch (type) { - case STATE: - case NBIRTH: - case NCMD: - case NDATA: - case NDEATH: - case NRECORD: - return new SparkplugTopic(namespace, groupId, edgeNodeId, type); - } - } else { - // A device topic - switch (type) { - case STATE: - case DBIRTH: - case DCMD: - case DDATA: - case DDEATH: - case DRECORD: - return new SparkplugTopic(namespace, groupId, edgeNodeId, splitTopic[4], type); - } - } - throw new ThingsboardException("Invalid number of topic elements " + length + " for topic type " + type, ThingsboardErrorCode.INVALID_ARGUMENTS); - } }