Browse Source

sparkplug: Telemetry

pull/7931/head
nickAS21 4 years ago
parent
commit
587e161398
  1. 86
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  2. 24
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java
  3. 174
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java
  4. 199
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java
  5. 244
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMetricUtil.java
  6. 147
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java

86
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.StringUtils;
import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.common.data.device.profile.MqttTopics; 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.DeviceId;
import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.id.OtaPackageId;
import org.thingsboard.server.common.data.ota.OtaPackageType; 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;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509CertRequestMsg; 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.queue.scheduler.SchedulerComponent;
import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor;
import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx;
import org.thingsboard.server.transport.mqtt.session.GatewaySessionHandler; import org.thingsboard.server.transport.mqtt.session.GatewaySessionHandler;
import org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher; import org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher;
import org.thingsboard.server.transport.mqtt.session.SparkplugNodeSessionHandler; 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.ReturnCode;
import org.thingsboard.server.transport.mqtt.util.ReturnCodeResolver; import org.thingsboard.server.transport.mqtt.util.ReturnCodeResolver;
import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic;
import javax.net.ssl.SSLPeerUnverifiedException; import javax.net.ssl.SSLPeerUnverifiedException;
import java.io.IOException; import java.io.IOException;
@ -97,16 +97,15 @@ import java.util.regex.Matcher;
import java.util.regex.Pattern; import java.util.regex.Pattern;
import static com.amazonaws.util.StringUtils.UTF8; 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.CONNECT;
import static io.netty.handler.codec.mqtt.MqttMessageType.PINGRESP; 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.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_LEAST_ONCE;
import static io.netty.handler.codec.mqtt.MqttQoS.AT_MOST_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_CLOSED;
import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_OPEN; 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 * @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 static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE;
private final UUID sessionId; private final UUID sessionId;
private final MqttTransportContext context; protected final MqttTransportContext context;
private final TransportService transportService; private final TransportService transportService;
private final SchedulerComponent scheduler; private final SchedulerComponent scheduler;
private final SslHandler sslHandler; private final SslHandler sslHandler;
@ -324,15 +323,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
String topicName = mqttMsg.variableHeader().topicName(); String topicName = mqttMsg.variableHeader().topicName();
int msgId = mqttMsg.variableHeader().packetId(); int msgId = mqttMsg.variableHeader().packetId();
log.trace("[{}][{}] Processing publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId); log.trace("[{}][{}] Processing publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId);
if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) {
if (sparkplugSessionHandler != null) {
handleSparkplugPublishMsg(ctx, topicName, msgId, mqttMsg);
transportService.reportActivity(deviceSessionCtx.getSessionInfo());
} else if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) {
if (gatewaySessionHandler != null) { if (gatewaySessionHandler != null) {
handleGatewayPublishMsg(ctx, topicName, msgId, mqttMsg); handleGatewayPublishMsg(ctx, topicName, msgId, mqttMsg);
transportService.reportActivity(deviceSessionCtx.getSessionInfo()); transportService.reportActivity(deviceSessionCtx.getSessionInfo());
} }
} else if (sparkplugSessionHandler != null) {
handleSparkplugPublishMsg(ctx, topicName, mqttMsg);
} else { } else {
processDevicePublish(ctx, mqttMsg, topicName, msgId); 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 { 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) { } 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(); ctx.close();
} catch (Exception e) { } catch (AdaptorException | ThingsboardException e) {
log.debug("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); log.error("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e);
sendAckOrCloseSession(ctx, topicName, msgId); sendAckOrCloseSession(ctx, topicName, msgId);
} }
} }
@ -648,7 +691,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
MqttQoS reqQoS = subscription.qualityOfService(); MqttQoS reqQoS = subscription.qualityOfService();
try { try {
if (sparkplugSessionHandler != null) { 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); sparkplugSessionHandler.handleSparkplugSubscribeMsg(grantedQoSList, sparkplugTopic, reqQoS);
} else { } else {
switch (topic) { switch (topic) {
@ -1012,14 +1055,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private void checkSparkplugSession(MqttConnectMessage connectMessage) { private void checkSparkplugSession(MqttConnectMessage connectMessage) {
try { try {
SparkplugTopic sparkplugTopic = parseTopic(connectMessage.payload().willTopic()); SparkplugTopic sparkplugTopic = parseTopicPublish(connectMessage.payload().willTopic());
// Test proto
SparkplugBProto.Payload payloadBProto = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes());
//
if (sparkplugSessionHandler == null) { if (sparkplugSessionHandler == null) {
sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId, sparkplugTopic.toString()); sparkplugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId);
} else {
log.warn("SparkPlugNodeReConnected [{}] [{}]", sparkplugTopic.getDeviceId(), sparkplugTopic.getType());
} }
} catch (Exception e) { } catch (Exception e) {
log.trace("[{}][{}] Failed to fetch sparkplugDevice additional info or sparkplugTopicName", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName(), e); log.trace("[{}][{}] Failed to fetch sparkplugDevice additional info or sparkplugTopicName", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName(), e);

24
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java

@ -79,20 +79,20 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe
@Slf4j @Slf4j
public abstract class AbstractGatewaySessionHandler { 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 CAN_T_PARSE_VALUE = "Can't parse value: ";
private static final String DEVICE_PROPERTY = "device"; private static final String DEVICE_PROPERTY = "device";
private final MqttTransportContext context; protected final MqttTransportContext context;
private final TransportService transportService; private final TransportService transportService;
private final TransportDeviceInfo gateway; protected final TransportDeviceInfo gateway;
private final UUID sessionId; protected final UUID sessionId;
private final ConcurrentMap<String, Lock> deviceCreationLockMap; private final ConcurrentMap<String, Lock> deviceCreationLockMap;
private final ConcurrentMap<String, MqttDeviceAwareSessionContext> devices; private final ConcurrentMap<String, MqttDeviceAwareSessionContext> devices;
private final ConcurrentMap<String, ListenableFuture<MqttDeviceAwareSessionContext>> deviceFutures; private final ConcurrentMap<String, ListenableFuture<MqttDeviceAwareSessionContext>> deviceFutures;
private final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap; private final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap;
private final ChannelHandlerContext channel; protected final ChannelHandlerContext channel;
private final DeviceSessionCtx deviceSessionCtx; protected final DeviceSessionCtx deviceSessionCtx;
public AbstractGatewaySessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId) { public AbstractGatewaySessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId) {
this.context = deviceSessionCtx.getContext(); this.context = deviceSessionCtx.getContext();
@ -198,7 +198,7 @@ public abstract class AbstractGatewaySessionHandler {
return deviceSessionCtx.isJsonPayloadType(); 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); log.trace("[{}] onDeviceConnect: {}", sessionId, deviceName);
Futures.addCallback(onDeviceConnect(deviceName, deviceType), new FutureCallback<MqttDeviceAwareSessionContext>() { Futures.addCallback(onDeviceConnect(deviceName, deviceType), new FutureCallback<MqttDeviceAwareSessionContext>() {
@Override @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)); transportService.process(deviceCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(channel, deviceName, msgId, postTelemetryMsg));
} }
@ -666,7 +666,7 @@ public abstract class AbstractGatewaySessionHandler {
return result.build(); return result.build();
} }
private ListenableFuture<MqttDeviceAwareSessionContext> checkDeviceConnected(String deviceName) { protected ListenableFuture<MqttDeviceAwareSessionContext> checkDeviceConnected(String deviceName) {
MqttDeviceAwareSessionContext ctx = devices.get(deviceName); MqttDeviceAwareSessionContext ctx = devices.get(deviceName);
if (ctx == null) { if (ctx == null) {
log.debug("[{}] Missing device [{}] for the gateway session", sessionId, deviceName); 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)) { if (StringUtils.isEmpty(deviceName)) {
throw new RuntimeException("Device name is empty!"); throw new RuntimeException("Device name is empty!");
} else { } else {
@ -697,11 +697,11 @@ public abstract class AbstractGatewaySessionHandler {
return JsonMqttAdaptor.validateJsonPayload(sessionId, mqttMsg.payload()); return JsonMqttAdaptor.validateJsonPayload(sessionId, mqttMsg.payload());
} }
private byte[] getBytes(ByteBuf payload) { protected byte[] getBytes(ByteBuf payload) {
return ProtoMqttAdaptor.toBytes(payload); return ProtoMqttAdaptor.toBytes(payload);
} }
private void ack(MqttPublishMessage msg, ReturnCode returnCode) { protected void ack(MqttPublishMessage msg, ReturnCode returnCode) {
int msgId = getMsgId(msg); int msgId = getMsgId(msg);
if (msgId > 0) { if (msgId > 0) {
writeAndFlush(MqttTransportHandler.createMqttPubAckMsg(deviceSessionCtx, msgId, returnCode)); writeAndFlush(MqttTransportHandler.createMqttPubAckMsg(deviceSessionCtx, msgId, returnCode));

174
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; 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.MqttPublishMessage;
import io.netty.handler.codec.mqtt.MqttPublishVariableHeader;
import io.netty.handler.codec.mqtt.MqttQoS; import io.netty.handler.codec.mqtt.MqttQoS;
import lombok.extern.slf4j.Slf4j; 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 org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic;
import javax.annotation.Nullable;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Optional;
import java.util.UUID; 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 * Created by nickAS21 on 12.12.22
*/ */
@Slf4j @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) { public TransportProtos.PostTelemetryMsg convertToPostTelemetry(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
super(deviceSessionCtx, sessionId); DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx;
this.nodeTopic = nodeTopic; 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 { public void onDeviceTelemetryProto(int msgId, ByteBuf payload, String deviceName, boolean isNode) throws AdaptorException {
SparkplugTopic sparkplugTopic = parseTopic(topicName); try {
log.warn("SparkplugPublishMsg [{}] [{}]", sparkplugTopic.isNode() ? "node" : "device: " + sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); checkDeviceName(deviceName);
if (sparkplugTopic.isNode()) { SparkplugBProto.Payload sparkplugBProto = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(payload));
// A node topic List<TransportProtos.PostTelemetryMsg> msgs = convertToPostTelemetry(sparkplugBProto);
switch (sparkplugTopic.getType()) { int finalMsgId = msgId;
case STATE: ListenableFuture<MqttDeviceAwareSessionContext> contextListenableFuture = isNode ?
// TODO Futures.immediateFuture(this.deviceSessionCtx) : checkDeviceConnected(deviceName);
break; for (TransportProtos.PostTelemetryMsg msg : msgs) {
case NBIRTH: Futures.addCallback(contextListenableFuture,
// TODO new FutureCallback<>() {
break; @Override
case NCMD: public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) {
// TODO try {
break; processPostTelemetryMsg(deviceCtx, msg, deviceName, finalMsgId);
case NDATA: } catch (Throwable e) {
// TODO log.warn("[{}][{}] Failed to convert telemetry: {}", gateway.getDeviceId(), deviceName, msg, e);
break; channel.close();
case NDEATH: }
onGatewayDeviceDisconnectProto(mqttMsg); }
break;
case NRECORD: @Override
// TODO public void onFailure(Throwable t) {
break; log.debug("[{}] Failed to process device telemetry command: {}", sessionId, deviceName, t);
default: }
} }, context.getExecutor());
} 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:
} }
} catch (RuntimeException | InvalidProtocolBufferException e) {
throw new AdaptorException(e);
} }
} }
public void handleSparkplugSubscribeMsg(List<Integer> grantedQoSList, SparkplugTopic sparkplugTopic, MqttQoS reqQoS) { public void handleSparkplugSubscribeMsg(List<Integer> grantedQoSList, SparkplugTopic sparkplugTopic, MqttQoS reqQoS) {
String topicName = sparkplugTopic.toString();
log.warn("SparkplugSubscribeMsg [{}] [{}]", sparkplugTopic.isNode() ? "node" : "device: " + sparkplugTopic.getDeviceId(), sparkplugTopic.getType());
if (sparkplugTopic.getGroupId() == null) { if (sparkplugTopic.getGroupId() == null) {
// TODO SUBSCRIBE NameSpace // TODO SUBSCRIBE NameSpace
} else if (sparkplugTopic.getType() == null) { } else if (sparkplugTopic.getType() == null) {
// TODO SUBSCRIBE GroupId // TODO SUBSCRIBE GroupId
} } else if (sparkplugTopic.isNode()) {
else if (sparkplugTopic.isNode()) {
// A node topic // A node topic
switch (sparkplugTopic.getType()) { switch (sparkplugTopic.getType()) {
case STATE: case STATE:
@ -150,4 +156,52 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler{
} }
} }
private List<TransportProtos.PostTelemetryMsg> convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto) throws AdaptorException {
try {
List<TransportProtos.PostTelemetryMsg> msgs = new ArrayList<>();
for (SparkplugBProto.Payload.Metric protoMetric : sparkplugBProto.getMetricsList()) {
long ts = protoMetric.getTimestamp();
Optional<TransportProtos.KeyValueProto> keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(protoMetric.getName(), protoMetric);
if (keyValueProtoOpt.isPresent()) {
List<TransportProtos.KeyValueProto> result = new ArrayList<>();
result.add(keyValueProtoOpt.get());
TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder();
TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder();
builder.setTs(ts);
builder.addAllKv(result);
request.addTsKvList(builder.build());
msgs.add(request.build());
}
}
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);
}
}
} }

199
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/MetricDataType.java

@ -0,0 +1,199 @@
/**
* Copyright © 2016-2022 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.mqtt.util.sparkplug;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto;
import java.math.BigInteger;
import java.util.Date;
/**
* Created by nickAS21 on 10.01.23
*/
@Slf4j
public enum MetricDataType {
// Basic Types
Int8(1, Byte.class),
Int16(2, Short.class),
Int32(3, Integer.class),
Int64(4, Long.class),
UInt8(5, Short.class),
UInt16(6, Integer.class),
UInt32(7, Long.class),
UInt64(8, BigInteger.class),
Float(9, Float.class),
Double(10, Double.class),
Boolean(11, Boolean.class),
String(12, String.class),
DateTime(13, Date.class),
Text(14, String.class),
// Custom Types for Metrics
UUID(15, String.class),
DataSet(16, SparkplugBProto.Payload.DataSet.class),
Bytes(17, byte[].class),
File(18, SparkplugMetricUtil.File.class),
Template(19, SparkplugBProto.Payload.Template.class),
// PropertyValue Types (20 and 21) are NOT metric datatypes
// Array Types
Int8Array(22, Byte[].class),
Int16Array(23, Short[].class),
Int32Array(24, Integer[].class),
Int64Array(25, Long[].class),
UInt8Array(26, Short[].class),
UInt16Array(27, Integer[].class),
UInt32Array(28, Long[].class),
UInt64Array(29, BigInteger[].class),
FloatArray(30, Float[].class),
DoubleArray(31, Double[].class),
BooleanArray(32, Boolean[].class),
StringArray(33, String[].class),
DateTimeArray(34, Date[].class),
// Unknown
Unknown(0, Object.class);
private Class<?> clazz = null;
private int intValue = 0;
/**
* Constructor
*
* @param intValue the integer value of this {@link MetricDataType}
* @param clazz the {@link Class} type associated with this {@link MetricDataType}
*/
private MetricDataType(int intValue, Class<?> clazz) {
this.intValue = intValue;
this.clazz = clazz;
}
/**
* Checks the type of a specified value against the specified {@link MetricDataType}
*
* @param value the {@link Object} value to check against the {@link MetricDataType}
* @throws AdaptorException if the value is not a valid type for the given {@link MetricDataType}
*/
public void checkType(Object value) throws AdaptorException {
if (value != null && !clazz.isAssignableFrom(value.getClass())) {
String msgError = "Failed type check - " + clazz + " != " + ((value != null) ? value.getClass().toString() : "null");
log.debug(msgError);
throw new AdaptorException(msgError);
}
}
/**
* Returns an integer representation of the data type.
*
* @return an integer representation of the data type.
*/
public int toIntValue() {
return this.intValue;
}
/**
* Converts the integer representation of the data type into a {@link MetricDataType} instance.
*
* @param i the integer representation of the data type.
* @return a {@link MetricDataType} instance.
*/
public static MetricDataType fromInteger(int i) {
switch (i) {
case 1:
return Int8;
case 2:
return Int16;
case 3:
return Int32;
case 4:
return Int64;
case 5:
return UInt8;
case 6:
return UInt16;
case 7:
return UInt32;
case 8:
return UInt64;
case 9:
return Float;
case 10:
return Double;
case 11:
return Boolean;
case 12:
return String;
case 13:
return DateTime;
case 14:
return Text;
case 15:
return UUID;
case 16:
return DataSet;
case 17:
return Bytes;
case 18:
return File;
case 19:
return Template;
case 22:
return Int8Array;
case 23:
return Int16Array;
case 24:
return Int32Array;
case 25:
return Int64Array;
case 26:
return UInt8Array;
case 27:
return UInt16Array;
case 28:
return UInt32Array;
case 29:
return UInt64Array;
case 30:
return FloatArray;
case 31:
return DoubleArray;
case 32:
return BooleanArray;
case 33:
return StringArray;
case 34:
return DateTimeArray;
default:
return Unknown;
}
}
/**
* Returns the class type for this DataType
*
* @return the class type for this DataType
*/
public Class<?> getClazz() {
return clazz;
}
}

244
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<TransportProtos.KeyValueProto> getFromSparkplugBMetricToKeyValueProto(String key, SparkplugBProto.Payload.Metric protoMetric) throws ThingsboardException {
// Check if the null flag has been set indicating that the value is null
if (protoMetric.getIsNull()) {
return 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<String> 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<Long> dateTimeList = new ArrayList<Long>();
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();
}
}
}

147
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicUtil.java

@ -27,88 +27,79 @@ import java.util.Map;
* Provides utility methods for handling Sparkplug MQTT message topics. * Provides utility methods for handling Sparkplug MQTT message topics.
*/ */
public class SparkplugTopicUtil { public class SparkplugTopicUtil {
private static final Map<String, String[]> SPLIT_TOPIC_CACHE = new HashMap<String, String[]>();
public static String[] getSplitTopic(String topic) {
String[] splitTopic = SPLIT_TOPIC_CACHE.get(topic);
if (splitTopic == null) {
splitTopic = topic.split("/");
SPLIT_TOPIC_CACHE.put(topic, splitTopic);
}
return splitTopic;
}
/** private static final Map<String, String[]> SPLIT_TOPIC_CACHE = new HashMap<String, String[]>();
* Serializes a {@link SparkplugTopic} instance in to a JSON string. private static final String TOPIC_INVALID_NUMBER = "Invalid number of topic elements: ";
*
* @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);
}
/** public static String[] getSplitTopic(String topic) {
* Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance. String[] splitTopic = SPLIT_TOPIC_CACHE.get(topic);
* if (splitTopic == null) {
* @param topic a topic string splitTopic = topic.split("/");
* @return a {@link SparkplugTopic} instance SPLIT_TOPIC_CACHE.put(topic, splitTopic);
* @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));
}
/** return 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;
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]; * Parses a Sparkplug MQTT message topic string and returns a {@link SparkplugTopic} instance.
type = SparkplugMessageType.parseMessageType(splitTopic[2]); *
edgeNodeId = splitTopic[3]; * @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);
}
} }

Loading…
Cancel
Save