diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java index 962d85ec7c..2ed3e0d161 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java @@ -88,7 +88,7 @@ public abstract class AbstractMqttIntegrationTest extends AbstractTransportInteg assertNotNull(accessToken); } if (config.getGatewayName() != null) { - savedGateway = createDevice(config.getGatewayName(), deviceProfile.getName(), true); + savedGateway = createDevice(config.getGatewayName(), deviceProfile.getName(), !config.isSparkplug); DeviceCredentials gatewayCredentials = doGet("/api/device/" + savedGateway.getId().getId().toString() + "/credentials", DeviceCredentials.class); assertNotNull(gatewayCredentials); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java index 6fb08f9eea..bbf8a291dc 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java +++ b/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() .setTimestamp(ts) .setSeq(getSeqNum()); - String deviceName = deviceId + "_" + i; - + String deviceIdName = deviceId + "_" + i; + String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; payloadBirthDevice.addMetrics(metric); 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); AtomicReference device = new AtomicReference<>(); - await(alias + "find device [" + deviceName + "] after created") + await(alias + "find device [" + deviceIdName + "] after created") .atMost(200, TimeUnit.SECONDS) .until(() -> { device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); @@ -194,7 +194,6 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte }); devices.add(device.get()); } - } Assert.assertEquals(cntDevices, devices.size()); @@ -224,11 +223,12 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() .setTimestamp(ts) .setSeq(getSeqNum()); - String deviceName = deviceId + "_" + 1; + String deviceIdName = deviceId + "_" + 1; + String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; payloadBirthDevice.addMetrics(metric); 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); AtomicReference device = new AtomicReference<>(); await(alias + "find device [" + deviceName + "] after created") 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 96f6c8a876..65ec849ac6 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 @@ -108,10 +108,9 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler contextListenableFuture; - TransportProtos.SessionInfoProto sessionInfo = this.deviceSessionCtx.getSessionInfo(); if (topic.isNode()) { if (topic.isType(NBIRTH)) { - sendSparkplugStateOnTelemetry(sessionInfo, deviceName, ONLINE, + sendSparkplugStateOnTelemetry(this.deviceSessionCtx.getSessionInfo(), deviceName, ONLINE, sparkplugBProto.getTimestamp()); setNodeBirthMetrics(sparkplugBProto.getMetricsList()); } @@ -123,7 +122,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { if (topic.isType(DBIRTH)) { - sendSparkplugStateOnTelemetry(sessionInfo, finalDeviceName, ONLINE, + sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), finalDeviceName, ONLINE, sparkplugBProto.getTimestamp()); ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); }