Browse Source

Merge pull request #15753 from thingsboard/sparkplug_fix_bug_for_unique-device-names

sparkplug_fix_bug_for_unique-device-names
pull/15769/head
Viacheslav Klimov 4 months ago
committed by GitHub
parent
commit
b60fb96983
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 52
      application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
  2. 95
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java
  3. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java
  4. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java
  5. 3
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java
  6. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java
  7. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java
  8. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java
  9. 8
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java
  10. 35
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java

52
application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java

@ -141,6 +141,7 @@ public class DefaultTransportApiService implements TransportApiService {
private final OtaPackageService otaPackageService; private final OtaPackageService otaPackageService;
private final OtaPackageDataCache otaPackageDataCache; private final OtaPackageDataCache otaPackageDataCache;
private final QueueService queueService; private final QueueService queueService;
public static final String GATEWAY_CREATED_RELATION = "Created";
private final ConcurrentMap<String, ReentrantLock> deviceCreationLocks = new ConcurrentReferenceHashMap<>(16, ConcurrentReferenceHashMap.ReferenceType.WEAK); private final ConcurrentMap<String, ReentrantLock> deviceCreationLocks = new ConcurrentReferenceHashMap<>(16, ConcurrentReferenceHashMap.ReferenceType.WEAK);
@ -395,47 +396,32 @@ public class DefaultTransportApiService implements TransportApiService {
// Security check: verify that the device was created by this gateway // Security check: verify that the device was created by this gateway
boolean isRelated = false; boolean isRelated = false;
try { isRelated = relationService.checkRelation(
// Security check: verify that the device was originally created by this gateway gateway.getTenantId(),
isRelated = relationService.checkRelation( gateway.getId(),
gateway.getTenantId(), existingDevice.getId(),
GATEWAY_CREATED_RELATION,
RelationTypeGroup.COMMON
);
// If the device is found but not related to this gateway
if (!isRelated) {
log.debug("[{}] Device [{}] exists but is not related to gateway. " +
"Skipping rename and allowing creation of Sparkplug device [{}].",
gateway.getId(), gateway.getId(),
existingDevice.getId(), existingDevice.getId(),
"Created", requestMsg.getDeviceName());
RelationTypeGroup.COMMON
);
} catch (Exception e) {
// Log the error from the relation service but return null to allow potential recovery
log.error("[{}] Error checking relation for device {}", gateway.getId(), existingDevice.getId(), e);
return null;
}
// If the device is found but not related to this gateway, it's a security breach return null;
if (!isRelated) {
log.error("[{}] Security breach attempt! Gateway tried to rename device [{}] without 'Created' relation.",
gateway.getId(), existingDevice.getId());
// Throwing exception to halt the entire connection process
throw new RuntimeException("Security breach attempt! Unauthorized device rename.");
} }
// Logic for renaming the device if it's related and no naming conflicts exist // Logic for renaming the device if it's related and no naming conflicts exist
boolean changed = false; boolean changed = false;
String newName = requestMsg.getDeviceName(); String newName = requestMsg.getDeviceName();
if (!newName.equals(existingDevice.getName())) { if (!newName.equals(existingDevice.getName())) {
// Check if the new name is already taken by another device
Device conflictDevice = deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), newName);
if (conflictDevice != null) {
log.warn("[{}] Cannot rename device [{}] to [{}]: name already exists!",
gateway.getId(), existingDevice.getId(), newName);
return existingDevice;
}
existingDevice.setName(newName); existingDevice.setName(newName);
// Update label only if it's empty to avoid overwriting user changes if (StringUtils.isEmpty(existingDevice.getLabel())) {
if (existingDevice.getLabel() == null || existingDevice.getLabel().isEmpty()) {
existingDevice.setLabel(deviceId); existingDevice.setLabel(deviceId);
} }
@ -452,8 +438,8 @@ public class DefaultTransportApiService implements TransportApiService {
Device device = new Device(); Device device = new Device();
device.setTenantId(tenantId); device.setTenantId(tenantId);
device.setName(requestMsg.getDeviceName()); device.setName(requestMsg.getDeviceName());
if (requestMsg.getIsSparkplug()) { if (requestMsg.getIsSparkplug() && topicPath.length == 3) {
if (topicPath.length == 3) device.setLabel(topicPath[2]); device.setLabel(topicPath[2]);
} }
device.setType(requestMsg.getDeviceType()); device.setType(requestMsg.getDeviceType());
device.setCustomerId(gateway.getCustomerId()); device.setCustomerId(gateway.getCustomerId());
@ -466,7 +452,7 @@ public class DefaultTransportApiService implements TransportApiService {
device = deviceService.saveDevice(device); device = deviceService.saveDevice(device);
relationService.saveRelation( relationService.saveRelation(
tenantId, tenantId,
new EntityRelation(gateway.getId(), device.getId(), "Created") new EntityRelation(gateway.getId(), device.getId(), GATEWAY_CREATED_RELATION)
); );
return device; return device;
} }

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

@ -31,7 +31,6 @@ import org.junit.Assert;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
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.asset.AssetInfo;
import org.thingsboard.server.common.data.exception.ThingsboardException; 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.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
@ -54,14 +53,19 @@ import java.util.Calendar;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.Set; import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
import static java.util.concurrent.Executors.newFixedThreadPool;
import static org.awaitility.Awaitility.await; import static org.awaitility.Awaitility.await;
import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.common.util.JacksonUtil.newArrayNode; import static org.thingsboard.common.util.JacksonUtil.newArrayNode;
import static org.thingsboard.server.service.transport.DefaultTransportApiService.GATEWAY_CREATED_RELATION;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Bytes; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Bytes;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32;
@ -108,7 +112,8 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
protected static final String metricBirthName_Int32 = "Device Metric int32"; protected static final String metricBirthName_Int32 = "Device Metric int32";
protected Set<String> sparkplugAttributesMetricNames; protected Set<String> sparkplugAttributesMetricNames;
public void beforeSparkplugTest(boolean isCreateDevices) throws Exception { public void beforeSparkplugTest() throws Exception {
MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder()
.gatewayName(edgeNodeDeviceName) .gatewayName(edgeNodeDeviceName)
.isSparkplug(true) .isSparkplug(true)
@ -116,24 +121,25 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.transportPayloadType(TransportPayloadType.PROTOBUF) .transportPayloadType(TransportPayloadType.PROTOBUF)
.build(); .build();
processBeforeTest(configProperties); processBeforeTest(configProperties);
if (isCreateDevices) { }
// 1. Create the first device with a short name (legacy style)
String deviceName1 = deviceId + "_1"; public void seedLegacyAndFullPathDevices() throws Exception {
Device device1 = createDevice(deviceName1, deviceProfile.getName(), false); // 1. Create the first device with a short name (legacy style)
String deviceName1 = deviceId + "_1";
// 2. Establish 'Created' relation so the transport identifies this gateway as the owner Device device1 = createDevice(deviceName1, deviceProfile.getName(), false);
String relationType = "Created";
EntityRelation relation1 = createFromRelation(savedGateway, device1, relationType); // 2. Establish 'Created' relation so the transport identifies this gateway as the owner
doPost("/api/relation", relation1).andExpect(status().isOk()); EntityRelation relation1 = createFromRelation(savedGateway, device1, GATEWAY_CREATED_RELATION);
doPost("/api/relation", relation1).andExpect(status().isOk());
// 3. Create the second device with a full-path name
String deviceName2 = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceId + "_2"; // 3. Create the second device with a full-path name
Device device2 = createDevice(deviceName2, deviceProfile.getName(), false); String deviceName2 = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceId + "_2";
Device device2 = createDevice(deviceName2, deviceProfile.getName(), false);
// 4. Establish 'Created' relation for the second device as well
EntityRelation relation2 = createFromRelation(savedGateway, device2, relationType); // 4. Establish 'Created' relation for the second device as well
doPost("/api/relation", relation2).andExpect(status().isOk()); EntityRelation relation2 = createFromRelation(savedGateway, device2, GATEWAY_CREATED_RELATION);
} doPost("/api/relation", relation2).andExpect(status().isOk());
} }
public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception { public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception {
@ -150,9 +156,9 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
public void clientWithCorrectNodeAccessTokenWithNDEATH(long ts, long value) throws Exception { public void clientWithCorrectNodeAccessTokenWithNDEATH(long ts, long value) throws Exception {
IMqttToken connectionResult = clientMqttV5ConnectWithNDEATH(ts, value, -1L); IMqttToken connectionResult = clientMqttV5ConnectWithNDEATH(ts, value, -1L);
MqttWireMessage response = connectionResult.getResponse(); MqttWireMessage response = connectionResult.getResponse();
Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType()); assertEquals(MESSAGE_TYPE_CONNACK, response.getType());
MqttConnAck connAckMsg = (MqttConnAck) response; MqttConnAck connAckMsg = (MqttConnAck) response;
Assert.assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode()); assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode());
} }
public IMqttToken clientMqttV5ConnectWithNDEATH(long ts, long value, Long alias, String... nameSpaceBad) throws Exception { public IMqttToken clientMqttV5ConnectWithNDEATH(long ts, long value, Long alias, String... nameSpaceBad) throws Exception {
@ -208,7 +214,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts) .setTimestamp(ts)
.setSeq(getSeqNum()); .setSeq(getSeqNum());
String deviceIdName = deviceId + "_" + i; String deviceIdName = deviceId + "_" + i;
String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; String deviceName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceIdName;
payloadBirthDevice.addMetrics(metric); payloadBirthDevice.addMetrics(metric);
if (client.isConnected()) { if (client.isConnected()) {
client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode + TOPIC_SPLIT_SEPARATOR + deviceIdName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode + TOPIC_SPLIT_SEPARATOR + deviceIdName,
@ -225,7 +231,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
} }
} }
Assert.assertEquals(cntDevices, devices.size()); assertEquals(cntDevices, devices.size());
return devices; return devices;
} }
@ -290,10 +296,10 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
return device2.get() != null; return device2.get() != null;
}); });
devices.add(device2.get()); devices.add(device2.get());
Assert.assertEquals(cntDevices, devices.size()); assertEquals(cntDevices, devices.size());
state_ONLINE_ALL (devices, calendar.getTimeInMillis()); state_ONLINE_ALL (devices, calendar.getTimeInMillis());
// Without full topic: as it was in the old version. When deviceId is updated to full theme, Label is also updated to old deviceId // Without full topic: as it was in the old version. When deviceId is updated to full theme, Label is also updated to old deviceId
Assert.assertEquals(deviceIdNameLabel1, device1.get().getLabel()); assertEquals(deviceIdNameLabel1, device1.get().getLabel());
// // With a full topic: if new. When creating a device by a client to a full topic, if the Label was not filled in - we do not touch it. // // With a full topic: if new. When creating a device by a client to a full topic, if the Label was not filled in - we do not touch it.
Assert.assertNull(device2.get().getLabel()); Assert.assertNull(device2.get().getLabel());
} }
@ -304,7 +310,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
protected void renameCollisionWhenTargetNameAlreadyExists_Test() throws Exception { protected void renameCollisionWhenTargetNameAlreadyExists_Test() throws Exception {
long ts = calendar.getTimeInMillis(); long ts = calendar.getTimeInMillis();
String shortName = deviceId + "_1"; // Created in beforeTest String shortName = deviceId + "_1"; // Created in beforeTest
String fullPathName = groupId + ":" + edgeNode + ":" + shortName; String fullPathName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + shortName;
// Manually create a device that already has the "new" full-path name to trigger a collision // Manually create a device that already has the "new" full-path name to trigger a collision
createDevice(fullPathName, deviceProfile.getName(), false); createDevice(fullPathName, deviceProfile.getName(), false);
@ -318,7 +324,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
// Gateway sends DBIRTH for the short name. // Gateway sends DBIRTH for the short name.
// Transport will try to rename it but should find a conflict and handle it gracefully. // Transport will try to rename it but should find a conflict and handle it gracefully.
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + shortName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + shortName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
await("Checking stability after collision") await("Checking stability after collision")
@ -352,10 +358,10 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts).setSeq(getSeqNum()); .setTimestamp(ts).setSeq(getSeqNum());
// 2. Unauthorized gateway attempts to rename this device via Sparkplug topic path // 2. Unauthorized gateway attempts to rename this device via Sparkplug topic path
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + strangerName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + strangerName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
String expectedFullPath = groupId + ":" + edgeNode + ":" + strangerName; String expectedFullPath = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + strangerName;
// 3. Verify security: the original device must still be linked to its short name with the same ID // 3. Verify security: the original device must still be linked to its short name with the same ID
await("Verify original device was not hijacked") await("Verify original device was not hijacked")
@ -364,8 +370,8 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.untilAsserted(() -> { .untilAsserted(() -> {
// Check if the original device still exists with its original ID // Check if the original device still exists with its original ID
Device currentStranger = doGet("/api/tenant/devices?deviceName=" + strangerName, Device.class); Device currentStranger = doGet("/api/tenant/devices?deviceName=" + strangerName, Device.class);
Assert.assertNotNull("Original device disappeared!", currentStranger); assertNotNull("Original device disappeared!", currentStranger);
Assert.assertEquals("Security breach: Original device ID changed!", originalStrangerId, currentStranger.getId()); assertEquals("Security breach: Original device ID changed!", originalStrangerId, currentStranger.getId());
// Even if the gateway created a NEW device with a full path, it must have a different ID // Even if the gateway created a NEW device with a full path, it must have a different ID
Device newDevice = doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class); Device newDevice = doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class);
@ -387,6 +393,8 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
stranger.setName(strangerName); stranger.setName(strangerName);
stranger.setType("default"); stranger.setType("default");
doPost("/api/device", stranger); doPost("/api/device", stranger);
Device originalStrangerDevice =
doGet("/api/tenant/devices?deviceName=" + strangerName, Device.class);
clientWithCorrectNodeAccessTokenWithNDEATH(); clientWithCorrectNodeAccessTokenWithNDEATH();
@ -394,13 +402,18 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts).setSeq(getSeqNum()); .setTimestamp(ts).setSeq(getSeqNum());
// Unauthorized gateway attempts to rename the device via Sparkplug topic // Unauthorized gateway attempts to rename the device via Sparkplug topic
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + strangerName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + strangerName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
String expectedFullPath = groupId + ":" + edgeNode + ":" + strangerName; String expectedFullPath = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + strangerName;
await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> await().atMost(30, TimeUnit.SECONDS).untilAsserted(() ->
doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class, status().isNotFound()) doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class, status().isOk())
); );
Device strangerDevice =
doGet("/api/tenant/devices?deviceName=" + strangerName, Device.class);
assertNotNull(strangerDevice);
assertEquals(originalStrangerDevice.getId(), strangerDevice.getId());
} }
protected void state_ONLINE_ALL (List<Device> devices, long ts) { protected void state_ONLINE_ALL (List<Device> devices, long ts) {
@ -438,7 +451,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
String concurrentDeviceName = "concurrent_device"; String concurrentDeviceName = "concurrent_device";
clientWithCorrectNodeAccessTokenWithNDEATH(); clientWithCorrectNodeAccessTokenWithNDEATH();
java.util.concurrent.ExecutorService executor = java.util.concurrent.Executors.newFixedThreadPool(threadCount); ExecutorService executor = newFixedThreadPool(threadCount);
long ts = calendar.getTimeInMillis(); long ts = calendar.getTimeInMillis();
for (int i = 0; i < threadCount; i++) { for (int i = 0; i < threadCount; i++) {
@ -446,7 +459,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
try { try {
SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder() SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder()
.setTimestamp(ts).setSeq(0); .setTimestamp(ts).setSeq(0);
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + concurrentDeviceName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + concurrentDeviceName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
} catch (Exception e) { } catch (Exception e) {
log.error("Concurrent publish failed", e); log.error("Concurrent publish failed", e);
@ -454,9 +467,9 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
}); });
} }
String expectedName = groupId + ":" + edgeNode + ":" + concurrentDeviceName; String expectedName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + concurrentDeviceName;
await("Wait for concurrent registration result") await("Wait for concurrent registration result")
.atMost(40, TimeUnit.SECONDS) // Restored to 40s as requested .atMost(40, TimeUnit.SECONDS)
.until(() -> doGet("/api/tenant/devices?deviceName=" + expectedName, Device.class) != null); .until(() -> doGet("/api/tenant/devices?deviceName=" + expectedName, Device.class) != null);
executor.shutdown(); executor.shutdown();
@ -506,7 +519,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts) .setTimestamp(ts)
.setSeq(getSeqNum()); .setSeq(getSeqNum());
String deviceIdName = deviceId + "_1"; String deviceIdName = deviceId + "_1";
String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; String deviceName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceIdName;
payloadBirthDevice.addMetrics(metric); payloadBirthDevice.addMetrics(metric);
if (client.isConnected()) { if (client.isConnected()) {
@ -523,7 +536,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
devices.add(device.get()); devices.add(device.get());
} }
Assert.assertEquals(1, devices.size()); assertEquals(1, devices.size());
return devices; return devices;
} }

2
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java

@ -33,7 +33,7 @@ public class MqttV5ClientSparkplugBAttributesInProfileTest extends AbstractMqttV
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
sparkplugAttributesMetricNames = new HashSet<>(); sparkplugAttributesMetricNames = new HashSet<>();
sparkplugAttributesMetricNames.add(metricBirthName_Int32); sparkplugAttributesMetricNames.add(metricBirthName_Int32);
beforeSparkplugTest(false); beforeSparkplugTest();
} }
@After @After

2
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java

@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBAttributesTest extends AbstractMqttV5ClientSp
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
beforeSparkplugTest(false); beforeSparkplugTest();
} }
@After @After

3
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java

@ -34,7 +34,8 @@ public class MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest extends A
*/ */
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
beforeSparkplugTest(true); beforeSparkplugTest();
seedLegacyAndFullPathDevices();
} }
@After @After

2
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java

@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBConnectionTest extends AbstractMqttV5ClientSp
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
beforeSparkplugTest(false); beforeSparkplugTest();
} }
@After @After

2
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java

@ -28,7 +28,7 @@ public class MqttV5RpcSparkplugTest extends AbstractMqttV5RpcSparkplugTest {
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
beforeSparkplugTest(false); beforeSparkplugTest();
} }
@After @After

2
application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java

@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSpa
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
beforeSparkplugTest(false); beforeSparkplugTest();
} }
@After @After

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

@ -261,7 +261,7 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
log.trace("[{}][{}][{}] onDeviceConnect: [{}]", gateway.getTenantId(), gateway.getDeviceId(), sessionId, deviceName); log.trace("[{}][{}][{}] onDeviceConnect: [{}]", gateway.getTenantId(), gateway.getDeviceId(), sessionId, deviceName);
int msgId = getMsgId(msg); int msgId = getMsgId(msg);
AtomicBoolean ackSent = new AtomicBoolean(false); AtomicBoolean ackSent = new AtomicBoolean(false);
process(onDeviceConnect(deviceName, deviceType, false), process(onDeviceConnect(deviceName, deviceType),
result -> { result -> {
ack(msg, MqttReasonCodes.PubAck.SUCCESS); ack(msg, MqttReasonCodes.PubAck.SUCCESS);
log.trace("[{}][{}][{}] onDeviceConnectOk: [{}]", gateway.getTenantId(), gateway.getDeviceId(), sessionId, deviceName); log.trace("[{}][{}][{}] onDeviceConnectOk: [{}]", gateway.getTenantId(), gateway.getDeviceId(), sessionId, deviceName);
@ -277,6 +277,10 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
} }
} }
ListenableFuture<T> onDeviceConnect(String deviceName, String deviceType) {
return onDeviceConnect(deviceName, deviceType, false);
}
ListenableFuture<T> onDeviceConnect(String deviceName, String deviceType, boolean isSparkplug) { ListenableFuture<T> onDeviceConnect(String deviceName, String deviceType, boolean isSparkplug) {
T result = devices.get(deviceName); T result = devices.get(deviceName);
if (result == null) { if (result == null) {
@ -902,7 +906,7 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
} }
protected void process(String deviceName, Consumer<T> onSuccess, Consumer<Throwable> onFailure) { protected void process(String deviceName, Consumer<T> onSuccess, Consumer<Throwable> onFailure) {
ListenableFuture<T> deviceCtxFuture = onDeviceConnect(deviceName, DEFAULT_DEVICE_TYPE, false); ListenableFuture<T> deviceCtxFuture = onDeviceConnect(deviceName, DEFAULT_DEVICE_TYPE);
process(deviceCtxFuture, onSuccess, onFailure); process(deviceCtxFuture, onSuccess, onFailure);
} }

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

@ -115,26 +115,23 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<S
} }
contextListenableFuture = Futures.immediateFuture(this.deviceSessionCtx); contextListenableFuture = Futures.immediateFuture(this.deviceSessionCtx);
} else { } else {
try { deviceName = checkDeviceName(topic.getNodeDeviceNameAllPath());
deviceName = checkDeviceName(topic.getNodeDeviceNameAllPath()); ListenableFuture<SparkplugDeviceSessionContext> deviceCtx = this.onDeviceConnectProto(topic);
ListenableFuture<SparkplugDeviceSessionContext> deviceCtx = this.onDeviceConnectProto(topic); 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(ctx.getSessionInfo(), finalDeviceName, ONLINE, sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), finalDeviceName, ONLINE,
sparkplugBProto.getTimestamp()); sparkplugBProto.getTimestamp());
try { try {
ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList());
} catch (IllegalArgumentException | DuplicateKeyException e) { } catch (IllegalArgumentException | DuplicateKeyException e) {
log.error("[{}] Failed to set birth metrics", finalDeviceName, e); log.error("[{}] Failed to set birth metrics", finalDeviceName, e);
throw new RuntimeException(e); throw new RuntimeException(e);
}
} }
return ctx; }
}, MoreExecutors.directExecutor()); return ctx;
} catch (IllegalArgumentException | DuplicateKeyException e) { }, MoreExecutors.directExecutor());
throw new RuntimeException(e);
}
} }
Set<String> attributesMetricNames = ((MqttDeviceProfileTransportConfiguration) deviceSessionCtx Set<String> attributesMetricNames = ((MqttDeviceProfileTransportConfiguration) deviceSessionCtx
.getDeviceProfile().getProfileData().getTransportConfiguration()).getSparkplugAttributesMetricNames(); .getDeviceProfile().getProfileData().getTransportConfiguration()).getSparkplugAttributesMetricNames();

Loading…
Cancel
Save