Browse Source

sparkplug - add group to name device, refactoring tests

pull/15041/head
nickAS21 7 months ago
parent
commit
14b6495983
  1. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java
  2. 14
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java
  3. 5
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java

2
application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java

@ -88,7 +88,7 @@ public abstract class AbstractMqttIntegrationTest extends AbstractTransportInteg
assertNotNull(accessToken); assertNotNull(accessToken);
} }
if (config.getGatewayName() != null) { if (config.getGatewayName() != null) {
savedGateway = createDevice(config.getGatewayName(), deviceProfile.getName(), true); savedGateway = createDevice(config.getGatewayName(), deviceProfile.getName(), !config.isSparkplug);
DeviceCredentials gatewayCredentials = DeviceCredentials gatewayCredentials =
doGet("/api/device/" + savedGateway.getId().getId().toString() + "/credentials", DeviceCredentials.class); doGet("/api/device/" + savedGateway.getId().getId().toString() + "/credentials", DeviceCredentials.class);
assertNotNull(gatewayCredentials); assertNotNull(gatewayCredentials);

14
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java

@ -179,14 +179,14 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder()
.setTimestamp(ts) .setTimestamp(ts)
.setSeq(getSeqNum()); .setSeq(getSeqNum());
String deviceName = deviceId + "_" + i; String deviceIdName = deviceId + "_" + i;
String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName;
payloadBirthDevice.addMetrics(metric); payloadBirthDevice.addMetrics(metric);
if (client.isConnected()) { if (client.isConnected()) {
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceName, client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceIdName,
payloadBirthDevice.build().toByteArray(), 0, false); payloadBirthDevice.build().toByteArray(), 0, false);
AtomicReference<Device> device = new AtomicReference<>(); AtomicReference<Device> device = new AtomicReference<>();
await(alias + "find device [" + deviceName + "] after created") await(alias + "find device [" + deviceIdName + "] after created")
.atMost(200, TimeUnit.SECONDS) .atMost(200, TimeUnit.SECONDS)
.until(() -> { .until(() -> {
device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class));
@ -194,7 +194,6 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
}); });
devices.add(device.get()); devices.add(device.get());
} }
} }
Assert.assertEquals(cntDevices, devices.size()); Assert.assertEquals(cntDevices, devices.size());
@ -224,11 +223,12 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder()
.setTimestamp(ts) .setTimestamp(ts)
.setSeq(getSeqNum()); .setSeq(getSeqNum());
String deviceName = deviceId + "_" + 1; String deviceIdName = deviceId + "_" + 1;
String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName;
payloadBirthDevice.addMetrics(metric); payloadBirthDevice.addMetrics(metric);
if (client.isConnected()) { if (client.isConnected()) {
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceName, client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceIdName,
payloadBirthDevice.build().toByteArray(), 0, false); payloadBirthDevice.build().toByteArray(), 0, false);
AtomicReference<Device> device = new AtomicReference<>(); AtomicReference<Device> device = new AtomicReference<>();
await(alias + "find device [" + deviceName + "] after created") await(alias + "find device [" + deviceName + "] after created")

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

@ -108,10 +108,9 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<S
public void onAttributesTelemetryProto(int msgId, SparkplugBProto.Payload sparkplugBProto, SparkplugTopic topic) throws AdaptorException, ThingsboardException { public void onAttributesTelemetryProto(int msgId, SparkplugBProto.Payload sparkplugBProto, SparkplugTopic topic) throws AdaptorException, ThingsboardException {
String deviceName = checkDeviceName(this.deviceSessionCtx.getDeviceInfo().getDeviceName()); String deviceName = checkDeviceName(this.deviceSessionCtx.getDeviceInfo().getDeviceName());
ListenableFuture<MqttDeviceAwareSessionContext> contextListenableFuture; ListenableFuture<MqttDeviceAwareSessionContext> contextListenableFuture;
TransportProtos.SessionInfoProto sessionInfo = this.deviceSessionCtx.getSessionInfo();
if (topic.isNode()) { if (topic.isNode()) {
if (topic.isType(NBIRTH)) { if (topic.isType(NBIRTH)) {
sendSparkplugStateOnTelemetry(sessionInfo, deviceName, ONLINE, sendSparkplugStateOnTelemetry(this.deviceSessionCtx.getSessionInfo(), deviceName, ONLINE,
sparkplugBProto.getTimestamp()); sparkplugBProto.getTimestamp());
setNodeBirthMetrics(sparkplugBProto.getMetricsList()); setNodeBirthMetrics(sparkplugBProto.getMetricsList());
} }
@ -123,7 +122,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<S
String finalDeviceName = deviceName; String finalDeviceName = deviceName;
contextListenableFuture = Futures.transform(deviceCtx, ctx -> { contextListenableFuture = Futures.transform(deviceCtx, ctx -> {
if (topic.isType(DBIRTH)) { if (topic.isType(DBIRTH)) {
sendSparkplugStateOnTelemetry(sessionInfo, finalDeviceName, ONLINE, sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), finalDeviceName, ONLINE,
sparkplugBProto.getTimestamp()); sparkplugBProto.getTimestamp());
ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList());
} }

Loading…
Cancel
Save