Browse Source

sparkplug: Telemetry DBIRTH, NBIRH...

pull/7931/head
nickAS21 4 years ago
parent
commit
278935cdc9
  1. 21
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  2. 29
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java

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

29
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<TransportProtos.PostTelemetryMsg> msgs = convertToPostTelemetry(sparkplugBProto);
List<TransportProtos.PostTelemetryMsg> msgs = convertToPostTelemetry(sparkplugBProto, topicTypeName);
int finalMsgId = msgId;
ListenableFuture<MqttDeviceAwareSessionContext> 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<TransportProtos.PostTelemetryMsg> convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto) throws AdaptorException {
private List<TransportProtos.PostTelemetryMsg> convertToPostTelemetry(SparkplugBProto.Payload sparkplugBProto, String topicTypeName) throws AdaptorException {
try {
List<TransportProtos.PostTelemetryMsg> msgs = new ArrayList<>();
for (SparkplugBProto.Payload.Metric protoMetric : sparkplugBProto.getMetricsList()) {
long ts = protoMetric.getTimestamp();
Optional<TransportProtos.KeyValueProto> keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(protoMetric.getName(), protoMetric);
String keys = "bdSeq".equals(protoMetric.getName()) ?
topicTypeName + " " + protoMetric.getName() : protoMetric.getName();
Optional<TransportProtos.KeyValueProto> keyValueProtoOpt = getFromSparkplugBMetricToKeyValueProto(keys, protoMetric);
if (keyValueProtoOpt.isPresent()) {
List<TransportProtos.KeyValueProto> result = new ArrayList<>();
result.add(keyValueProtoOpt.get());
@ -173,6 +174,20 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler {
msgs.add(request.build());
}
}
if (DBIRTH.name().equals(topicTypeName)) {
List<TransportProtos.KeyValueProto> result = new ArrayList<>();
TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder();
keyValueProtoBuilder.setKey(topicTypeName + " " + "seq");
keyValueProtoBuilder.setType(TransportProtos.KeyValueType.LONG_V);
keyValueProtoBuilder.setLongV(sparkplugBProto.getSeq());
result.add(keyValueProtoBuilder.build());
TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder();
TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder();
builder.setTs(sparkplugBProto.getTimestamp());
builder.addAllKv(result);
request.addTsKvList(builder.build());
msgs.add(request.build());
}
return msgs;
} catch (IllegalStateException | JsonSyntaxException | ThingsboardException e) {
log.error("Failed to decode post telemetry request", e);

Loading…
Cancel
Save