Browse Source

Merge pull request #8222 from thingsboard/sparkplug-transport-proto

[fix_bug][3.5]Sparkplug-transport-proto
pull/8245/head
Andrew Shvayka 4 years ago
committed by GitHub
parent
commit
1147cae8a1
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 9
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java
  2. 1
      common/cluster-api/src/main/proto/queue.proto
  3. 7
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  4. 1
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java
  5. 9
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java
  6. 5
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java

9
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java

@ -137,6 +137,12 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstra
} }
} }
/**
* OFFLINE_All
* @param cntDevices
* @throws Exception
*/
protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OFFLINE_All(int cntDevices) throws Exception { protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OFFLINE_All(int cntDevices) throws Exception {
long ts = calendar.getTimeInMillis(); long ts = calendar.getTimeInMillis();
List<Device> devices = connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(cntDevices, ts); List<Device> devices = connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(cntDevices, ts);
@ -166,6 +172,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstra
} }
} }
private boolean findEqualsKeyValueInKvEntrys(List<TsKvEntry> finalFuture, TsKvEntry tsKvEntry) { private boolean findEqualsKeyValueInKvEntrys(List<TsKvEntry> finalFuture, TsKvEntry tsKvEntry) {
for (TsKvEntry kvEntry : finalFuture) { for (TsKvEntry kvEntry : finalFuture) {
if (kvEntry.getKey().equals(tsKvEntry.getKey()) && kvEntry.getValue().equals(tsKvEntry.getValue())) { if (kvEntry.getKey().equals(tsKvEntry.getKey()) && kvEntry.getValue().equals(tsKvEntry.getValue())) {
@ -174,5 +181,5 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstra
} }
return false; return false;
} }
} }

1
common/cluster-api/src/main/proto/queue.proto

@ -186,7 +186,6 @@ message GetOrCreateDeviceFromGatewayRequestMsg {
int64 gatewayIdLSB = 2; int64 gatewayIdLSB = 2;
string deviceName = 3; string deviceName = 3;
string deviceType = 4; string deviceType = 4;
bool sparkplug = 5;
} }
message GetOrCreateDeviceFromGatewayResponseMsg { message GetOrCreateDeviceFromGatewayResponseMsg {

7
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -399,7 +399,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
case NBIRTH: case NBIRTH:
case NCMD: case NCMD:
case NDATA: case NDATA:
sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoNode, deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopic); sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoNode, sparkplugTopic);
break; break;
default: default:
} }
@ -410,7 +410,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
case DBIRTH: case DBIRTH:
case DCMD: case DCMD:
case DDATA: case DDATA:
sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoDevice, sparkplugTopic.getDeviceId(), sparkplugTopic); sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoDevice, sparkplugTopic);
break; break;
case DDEATH: case DDEATH:
sparkplugSessionHandler.onDeviceDisconnect(mqttMsg, sparkplugTopic.getDeviceId()); sparkplugSessionHandler.onDeviceDisconnect(mqttMsg, sparkplugTopic.getDeviceId());
@ -1067,8 +1067,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (sparkplugTopicNode != null) { if (sparkplugTopicNode != null) {
SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes());
sparkplugSessionHandler = new SparkplugNodeSessionHandler(this, deviceSessionCtx, sessionId, sparkplugTopicNode); sparkplugSessionHandler = new SparkplugNodeSessionHandler(this, deviceSessionCtx, sessionId, sparkplugTopicNode);
sparkplugSessionHandler.onAttributesTelemetryProto(0, sparkplugBProtoNode, sparkplugSessionHandler.onAttributesTelemetryProto(0, sparkplugBProtoNode, sparkplugTopicNode);
deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopicNode);
sessionMetaData.setOverwriteActivityTime(true); sessionMetaData.setOverwriteActivityTime(true);
} else { } else {
log.trace("[{}][{}] Failed to fetch sparkplugDevice connect: sparkplugTopicName without SparkplugMessageType.NDEATH.", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName()); log.trace("[{}][{}] Failed to fetch sparkplugDevice connect: sparkplugTopicName without SparkplugMessageType.NDEATH.", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName());

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

@ -253,7 +253,6 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
.setDeviceType(deviceType) .setDeviceType(deviceType)
.setGatewayIdMSB(gateway.getDeviceId().getId().getMostSignificantBits()) .setGatewayIdMSB(gateway.getDeviceId().getId().getMostSignificantBits())
.setGatewayIdLSB(gateway.getDeviceId().getId().getLeastSignificantBits()) .setGatewayIdLSB(gateway.getDeviceId().getId().getLeastSignificantBits())
.setSparkplug(this.deviceSessionCtx.isSparkplug())
.build(), .build(),
new TransportServiceCallback<>() { new TransportServiceCallback<>() {
@Override @Override

9
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java

@ -101,7 +101,8 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<S
} }
} }
public void onAttributesTelemetryProto(int msgId, SparkplugBProto.Payload sparkplugBProto, String deviceName, SparkplugTopic topic) throws AdaptorException, ThingsboardException { public void onAttributesTelemetryProto(int msgId, SparkplugBProto.Payload sparkplugBProto, SparkplugTopic topic) throws AdaptorException, ThingsboardException {
String deviceName = topic.getNodeDeviceName();
checkDeviceName(deviceName); checkDeviceName(deviceName);
ListenableFuture<MqttDeviceAwareSessionContext> contextListenableFuture; ListenableFuture<MqttDeviceAwareSessionContext> contextListenableFuture;
@ -113,7 +114,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<S
} }
contextListenableFuture = Futures.immediateFuture(this.deviceSessionCtx); contextListenableFuture = Futures.immediateFuture(this.deviceSessionCtx);
} else { } else {
ListenableFuture<SparkplugDeviceSessionContext> deviceCtx = onDeviceConnectProto(deviceName); ListenableFuture<SparkplugDeviceSessionContext> deviceCtx = onDeviceConnectProto(topic);
contextListenableFuture = Futures.transform(deviceCtx, ctx -> { contextListenableFuture = Futures.transform(deviceCtx, ctx -> {
if (topic.isType(DBIRTH)) { if (topic.isType(DBIRTH)) {
sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), deviceName, ONLINE, sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), deviceName, ONLINE,
@ -218,10 +219,10 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<S
} }
} }
private ListenableFuture<SparkplugDeviceSessionContext> onDeviceConnectProto(String deviceName) throws ThingsboardException { private ListenableFuture<SparkplugDeviceSessionContext> onDeviceConnectProto(SparkplugTopic topic) throws ThingsboardException {
try { try {
String deviceType = this.gateway.getDeviceType() + "-node"; String deviceType = this.gateway.getDeviceType() + "-node";
return onDeviceConnect(deviceName, deviceType); return onDeviceConnect(topic.getNodeDeviceName(), deviceType);
} catch (RuntimeException e) { } catch (RuntimeException e) {
log.error("Failed Sparkplug Device connect proto!", e); log.error("Failed Sparkplug Device connect proto!", e);
throw new ThingsboardException(e, ThingsboardErrorCode.BAD_REQUEST_PARAMS); throw new ThingsboardException(e, ThingsboardErrorCode.BAD_REQUEST_PARAMS);

5
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java

@ -161,4 +161,9 @@ public class SparkplugTopic {
public boolean isNode() { public boolean isNode() {
return this.deviceId == null; return this.deviceId == null;
} }
public String getNodeDeviceName() {
return isNode() ? edgeNodeId : deviceId;
}
} }

Loading…
Cancel
Save