From ad8c56c5843fe020e3195df6018fe9898eae1d7b Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Fri, 6 Feb 2026 18:39:57 +0200 Subject: [PATCH 1/8] sparkplug - add group to name device --- .../AbstractGatewaySessionHandler.java | 12 +++++- .../session/SparkplugNodeSessionHandler.java | 37 ++++++++++--------- .../mqtt/util/sparkplug/SparkplugTopic.java | 11 ++++++ 3 files changed, 41 insertions(+), 19 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index 8ea6480142..cbb805731e 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java @@ -68,6 +68,7 @@ import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor; import org.thingsboard.server.transport.mqtt.gateway.GatewayMetricsService; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConnectionState; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic; import java.util.ArrayList; import java.util.Collections; @@ -296,6 +297,15 @@ public abstract class AbstractGatewaySessionHandler onDeviceConnectSparkplug(SparkplugTopic topic, String deviceType) { + T result = devices.get(topic.getNodeDeviceName()); + if (result == null) { + return onDeviceConnect(topic.getNodeDeviceNameAllPath(), deviceType); + } else { + return Futures.immediateFuture(result); + } + } + private ListenableFuture getDeviceCreationFuture(String deviceName, String deviceType) { final SettableFuture futureToSet = SettableFuture.create(); ListenableFuture future = deviceFutures.putIfAbsent(deviceName, futureToSet); @@ -844,7 +854,7 @@ public abstract class AbstractGatewaySessionHandler contextListenableFuture; + TransportProtos.SessionInfoProto sessionInfo = this.deviceSessionCtx.getSessionInfo(); if (topic.isNode()) { if (topic.isType(NBIRTH)) { - sendSparkplugStateOnTelemetry(this.deviceSessionCtx.getSessionInfo(), deviceName, ONLINE, + sendSparkplugStateOnTelemetry(sessionInfo, deviceName, ONLINE, sparkplugBProto.getTimestamp()); setNodeBirthMetrics(sparkplugBProto.getMetricsList()); } contextListenableFuture = Futures.immediateFuture(this.deviceSessionCtx); } else { - ListenableFuture deviceCtx = onDeviceConnectProto(topic); - contextListenableFuture = Futures.transform(deviceCtx, ctx -> { - if (topic.isType(DBIRTH)) { - sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), deviceName, ONLINE, - sparkplugBProto.getTimestamp()); - try { - ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); - } catch (IllegalArgumentException | DuplicateKeyException e) { - throw new RuntimeException(e); + try { + ListenableFuture deviceCtx = this.onDeviceConnectProto(topic); + deviceName = checkDeviceName(deviceCtx.get().getDeviceInfo().getDeviceName()); + String finalDeviceName = deviceName; + contextListenableFuture = Futures.transform(deviceCtx, ctx -> { + if (topic.isType(DBIRTH)) { + sendSparkplugStateOnTelemetry(sessionInfo, finalDeviceName, ONLINE, + sparkplugBProto.getTimestamp()); + ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); } - } - return ctx; - }, MoreExecutors.directExecutor()); + return ctx; + }, MoreExecutors.directExecutor()); + } catch (IllegalArgumentException | DuplicateKeyException | ExecutionException | InterruptedException e) { + throw new RuntimeException(e); + } } Set attributesMetricNames = ((MqttDeviceProfileTransportConfiguration) deviceSessionCtx .getDeviceProfile().getProfileData().getTransportConfiguration()).getSparkplugAttributesMetricNames(); @@ -222,7 +223,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler Date: Sat, 7 Feb 2026 19:47:59 +0200 Subject: [PATCH 2/8] sparkplug - add group to name device, refactoring tests --- .../mqtt/AbstractMqttIntegrationTest.java | 2 +- .../AbstractMqttV5ClientSparkplugTest.java | 14 +++++++------- .../mqtt/session/SparkplugNodeSessionHandler.java | 5 ++--- 3 files changed, 10 insertions(+), 11 deletions(-) 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()); } From ee8233c6fad0bc7ae675ea15120a73bdec143b52 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Sat, 14 Feb 2026 19:20:27 +0200 Subject: [PATCH 3/8] sparkplug - add tests with devices when devices are already created --- .../AbstractMqttV5ClientSparkplugTest.java | 136 ++++++++++++++++-- ...ientSparkplugBAttributesInProfileTest.java | 2 +- .../MqttV5ClientSparkplugBAttributesTest.java | 2 +- ...ctMqttV5ClientSparkplugConnectionTest.java | 26 +--- ...gBConnectionDevicesCreatingBeforeTest.java | 45 ++++++ .../MqttV5ClientSparkplugBConnectionTest.java | 2 +- .../sparkplug/rpc/MqttV5RpcSparkplugTest.java | 3 +- .../MqttV5ClientSparkplugBTelemetryTest.java | 2 +- 8 files changed, 180 insertions(+), 38 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java 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 bbf8a291dc..7d8770bbcf 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 @@ -67,6 +67,9 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataTyp import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt32; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt64; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.UInt8; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConnectionState.ONLINE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; @@ -82,6 +85,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte protected ThreadLocalRandom random = ThreadLocalRandom.current(); protected static final String groupId = "SparkplugBGroupId"; + protected static final String edgeNodeDeviceName = "Test Connect Sparkplug client node"; protected static final String edgeNode = "SparkpluBNode"; protected static final String keysBdSeq = "bdSeq"; protected static final String alias = "Failed Telemetry/Attribute proto sparkplug payload. SparkplugMessageType "; @@ -98,14 +102,39 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte protected static final String metricBirthName_Int32 = "Device Metric int32"; protected Set sparkplugAttributesMetricNames; - public void beforeSparkplugTest() throws Exception { - MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() - .gatewayName("Test Connect Sparkplug client node") - .isSparkplug(true) - .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) - .transportPayloadType(TransportPayloadType.PROTOBUF) - .build(); - processBeforeTest(configProperties); + public void beforeSparkplugTest(boolean isCreateDevices) throws Exception { + if (isCreateDevices) { + MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() + .gatewayName(edgeNodeDeviceName) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + configProperties = MqttTestConfigProperties.builder() + .gatewayName(deviceId) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + configProperties = MqttTestConfigProperties.builder() + .gatewayName(groupId + ":" + edgeNode + ":" + deviceId) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + + } else { + MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() + .gatewayName(edgeNodeDeviceName) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); + } } public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception { @@ -200,6 +229,97 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte return devices; } + protected void connectClientWithCorrectAccessTokenWithNDEATHDevicesCreatingBefore_Test(int cntDevices) throws Exception { + long ts = calendar.getTimeInMillis(); + List devices = new ArrayList<>(); + clientWithCorrectNodeAccessTokenWithNDEATH(); + MetricDataType metricDataType = Int32; + String key = "Node Metric int32"; + int valueDeviceInt32 = 1024; + SparkplugBProto.Payload.Metric metric = createMetric(valueDeviceInt32, ts, key, metricDataType, -1L); + SparkplugBProto.Payload.Builder payloadBirthNode = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getBdSeqNum()); + payloadBirthNode.addMetrics(metric); + payloadBirthNode.setTimestamp(ts); + if (client.isConnected()) { + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + payloadBirthNode.build().toByteArray(), 0, false); + } + + valueDeviceInt32 = 4024; + metric = createMetric(valueDeviceInt32, ts, metricBirthName_Int32, metricBirthDataType_Int32, -1L); + // as old device name -> deviceId + String deviceName = deviceId; + AtomicReference device1 = new AtomicReference<>(); + String finalDeviceName1 = deviceName; + await(alias + "find device [" + deviceId + "] before connecting") + .atMost(200, TimeUnit.SECONDS) + .until(() -> { + device1.set(doGet("/api/tenant/devices?deviceName=" + finalDeviceName1, Device.class)); + return device1.get() != null; + }); + + if (client.isConnected()) { + SparkplugBProto.Payload.Builder payloadBirthDevice1 = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getSeqNum()); + payloadBirthDevice1.addMetrics(metric); + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceId, + payloadBirthDevice1.build().toByteArray(), 0, false); + devices.add(device1.get()); + } + // as new device name -> groupId + ":" + edgeNode + ":" + deviceId; + deviceName = groupId + ":" + edgeNode + ":" + deviceId; + AtomicReference device2 = new AtomicReference<>(); + String finalDeviceName2 = deviceName; + await(alias + "find device [" + deviceName + "] before connecting") + .atMost(200, TimeUnit.SECONDS) + .until(() -> { + device2.set(doGet("/api/tenant/devices?deviceName=" + finalDeviceName2, Device.class)); + return device2.get() != null; + }); + + if (client.isConnected()) { + SparkplugBProto.Payload.Builder payloadBirthDevice2 = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getSeqNum()); + payloadBirthDevice2.addMetrics(metric); + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceId, + payloadBirthDevice2.build().toByteArray(), 0, false); + devices.add(device2.get()); + } + Assert.assertEquals(cntDevices, devices.size()); + state_ONLINE_ALL (devices, calendar.getTimeInMillis()); + } + + protected void state_ONLINE_ALL (List devices, long ts) { + TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(messageName(STATE), ONLINE.name())); + await(alias + messageName(STATE) + ", device: " + savedGateway.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + var foundEntry = tsService.findAllLatest(tenantId, savedGateway.getId()).get().stream() + .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) + .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) + .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) + .findFirst(); + return foundEntry.isPresent(); + }); + + for (Device device : devices) { + await(alias + messageName(STATE) + ", device: " + device.getName()) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + var foundEntry = tsService.findAllLatest(tenantId, device.getId()).get().stream() + .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) + .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) + .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) + .findFirst(); + return foundEntry.isPresent(); + }); + } + } + protected List connectClientWithCorrectAccessTokenWithNDEATHWithAliasCreatedDevices(long ts) throws Exception { List devices = new ArrayList<>(); Long alias = 0L; diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java index 37abade6e7..8cd713ad7d 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java +++ b/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 { sparkplugAttributesMetricNames = new HashSet<>(); sparkplugAttributesMetricNames.add(metricBirthName_Int32); - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java index 3826b59bdc..b080cfbff0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java @@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBAttributesTest extends AbstractMqttV5ClientSp @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java index 8f9943b393..89b17f6b48 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java @@ -95,31 +95,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstra protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_ALL(int cntDevices) throws Exception { long ts = calendar.getTimeInMillis(); List devices = connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices(cntDevices, ts); - - TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(messageName(STATE), ONLINE.name())); - await(alias + messageName(STATE) + ", device: " + savedGateway.getName()) - .atMost(40, TimeUnit.SECONDS) - .until(() -> { - var foundEntry = tsService.findAllLatest(tenantId, savedGateway.getId()).get().stream() - .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) - .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) - .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) - .findFirst(); - return foundEntry.isPresent(); - }); - - for (Device device : devices) { - await(alias + messageName(STATE) + ", device: " + device.getName()) - .atMost(40, TimeUnit.SECONDS) - .until(() -> { - var foundEntry = tsService.findAllLatest(tenantId, device.getId()).get().stream() - .filter(tsKv -> tsKv.getKey().equals(tsKvEntry.getKey())) - .filter(tsKv -> tsKv.getValue().equals(tsKvEntry.getValue())) - .filter(tsKv -> tsKv.getTs() == tsKvEntry.getTs()) - .findFirst(); - return foundEntry.isPresent(); - }); - } + state_ONLINE_ALL (devices, ts); } protected void processConnectClientWithCorrectAccessTokenWithNDEATH_State_ONLINE_All_Then_OneDeviceOFFLINE(int cntDevices, int indexDeviceDisconnect) throws Exception { diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java new file mode 100644 index 0000000000..797b1beea5 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java @@ -0,0 +1,45 @@ +/** + * Copyright © 2016-2026 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.sparkplug.connection; + +import org.eclipse.paho.mqttv5.common.MqttException; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.thingsboard.server.dao.service.DaoSqlTest; + +/** + * Created by nickAS21 on 12.01.23 + */ +@DaoSqlTest +public class MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest extends AbstractMqttV5ClientSparkplugConnectionTest { + + @Before + public void beforeTest() throws Exception { + beforeSparkplugTest(true); + } + + @After + public void afterTest() throws MqttException { + if (client.isConnected()) { + client.disconnect(); + } + } + + @Test + public void testClientWithCorrectAccessTokenWithNDEATHTwoDevicesCreatingBeforeFirstNameDeviceIdSecondNameFull() throws Exception { + connectClientWithCorrectAccessTokenWithNDEATHDevicesCreatingBefore_Test(2); + }} diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java index 01038c21f1..0d5e73e2f0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java @@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBConnectionTest extends AbstractMqttV5ClientSp @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java index e36b4e6260..08ade0274d 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java @@ -28,7 +28,7 @@ public class MqttV5RpcSparkplugTest extends AbstractMqttV5RpcSparkplugTest { @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After @@ -47,6 +47,7 @@ public class MqttV5RpcSparkplugTest extends AbstractMqttV5RpcSparkplugTest { public void testClientDeviceWithCorrectAccessTokenPublish_TwoWayRpc_Success() throws Exception { processClientDeviceWithCorrectAccessTokenPublish_TwoWayRpc_Success(); } + @Test public void testClientDeviceWithCorrectAccessTokenPublishWithAlias_TwoWayRpc_Success() throws Exception { processClientDeviceWithCorrectAccessTokenPublishWithAlias_TwoWayRpc_Success(); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java index bafe2d81d2..5d9f453415 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java @@ -29,7 +29,7 @@ public class MqttV5ClientSparkplugBTelemetryTest extends AbstractMqttV5ClientSpa @Before public void beforeTest() throws Exception { - beforeSparkplugTest(); + beforeSparkplugTest(false); } @After From 8f550475a026717b05199ba2d2906dca41e20cd3 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Sun, 15 Feb 2026 12:13:22 +0200 Subject: [PATCH 4/8] sparkplug - add tests with devices when devices are already created --- .../transport/DefaultTransportApiService.java | 210 ++++++++++++------ .../AbstractMqttV5ClientSparkplugTest.java | 98 ++++---- ...ctMqttV5ClientSparkplugConnectionTest.java | 4 +- ...actMqttV5ClientSparkplugTelemetryTest.java | 5 +- common/proto/src/main/proto/queue.proto | 1 + .../AbstractGatewaySessionHandler.java | 13 +- .../SparkplugDeviceSessionContext.java | 1 - .../mqtt/util/sparkplug/SparkplugTopic.java | 5 +- .../util/sparkplug/SparkplugTopicService.java | 1 + 9 files changed, 197 insertions(+), 141 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index 48bd497c6e..1baa3cdccf 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -114,6 +114,7 @@ import java.util.stream.Collectors; import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.PASSWORD_MISMATCH; import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.VALID; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_REGEXP; /** * Created by ashvayka on 05.10.18. @@ -330,86 +331,153 @@ public class DefaultTransportApiService implements TransportApiService { } private TransportApiResponseMsg handle(GetOrCreateDeviceFromGatewayRequestMsg requestMsg) { - DeviceId gatewayId = new DeviceId(new UUID(requestMsg.getGatewayIdMSB(), requestMsg.getGatewayIdLSB())); + DeviceId gatewayId = toDeviceId(requestMsg); Device gateway = deviceService.findDeviceById(TenantId.SYS_TENANT_ID, gatewayId); - Lock deviceCreationLock = deviceCreationLocks.computeIfAbsent(requestMsg.getDeviceName(), id -> new ReentrantLock()); - deviceCreationLock.lock(); + String deviceName = requestMsg.getDeviceName(); + Lock lock = deviceCreationLocks.computeIfAbsent(deviceName, k -> new ReentrantLock()); + lock.lock(); try { - Device device = deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), requestMsg.getDeviceName()); - if (device == null) { - TenantId tenantId = gateway.getTenantId(); - device = new Device(); - device.setTenantId(tenantId); - device.setName(requestMsg.getDeviceName()); - device.setType(requestMsg.getDeviceType()); - device.setCustomerId(gateway.getCustomerId()); - DeviceProfile deviceProfile = deviceProfileCache.findOrCreateDeviceProfile(gateway.getTenantId(), requestMsg.getDeviceType()); - - device.setDeviceProfileId(deviceProfile.getId()); - ObjectNode additionalInfo = JacksonUtil.newObjectNode(); - additionalInfo.put(DataConstants.LAST_CONNECTED_GATEWAY, gatewayId.toString()); - device.setAdditionalInfo(additionalInfo); - device = deviceService.saveDevice(device); - - relationService.saveRelation(tenantId, new EntityRelation(gateway.getId(), device.getId(), "Created")); - - TbMsgMetaData metaData = new TbMsgMetaData(); - CustomerId customerId = gateway.getCustomerId(); - if (customerId != null && !customerId.isNullUid()) { - metaData.putValue("customerId", customerId.toString()); - } - metaData.putValue("gatewayId", gatewayId.toString()); - - DeviceId deviceId = device.getId(); - JsonNode entityNode = JacksonUtil.valueToTree(device); - TbMsg tbMsg = TbMsg.newMsg() - .type(TbMsgType.ENTITY_CREATED) - .originator(deviceId) - .customerId(customerId) - .copyMetaData(metaData) - .dataType(TbMsgDataType.JSON) - .data(JacksonUtil.toString(entityNode)) - .build(); - tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, null); - } else { - JsonNode deviceAdditionalInfo = device.getAdditionalInfo(); - if (deviceAdditionalInfo == null) { - deviceAdditionalInfo = JacksonUtil.newObjectNode(); - } - if (deviceAdditionalInfo.isObject() && - (!deviceAdditionalInfo.has(DataConstants.LAST_CONNECTED_GATEWAY) - || !gatewayId.toString().equals(deviceAdditionalInfo.get(DataConstants.LAST_CONNECTED_GATEWAY).asText()))) { - ObjectNode newDeviceAdditionalInfo = (ObjectNode) deviceAdditionalInfo; - newDeviceAdditionalInfo.put(DataConstants.LAST_CONNECTED_GATEWAY, gatewayId.toString()); - deviceService.saveDevice(device); - } - } - GetOrCreateDeviceFromGatewayResponseMsg.Builder builder = GetOrCreateDeviceFromGatewayResponseMsg.newBuilder() - .setDeviceInfo(ProtoUtils.toDeviceInfoProto(device)); - DeviceProfile deviceProfile = deviceProfileCache.get(device.getTenantId(), device.getDeviceProfileId()); - if (deviceProfile != null) { - builder.setDeviceProfile(ProtoUtils.toProto(deviceProfile)); - } else { - log.warn("[{}] Failed to find device profile [{}] for device. ", device.getId(), device.getDeviceProfileId()); - } - return TransportApiResponseMsg.newBuilder() - .setGetOrCreateDeviceResponseMsg(builder.build()) - .build(); + Device device = findOrCreateDevice(requestMsg, gateway, gatewayId); + updateLastConnectedGateway(device, gatewayId); + return buildResponse(device); } catch (JsonProcessingException e) { - log.warn("[{}] Failed to lookup device by gateway id and name: [{}]", gatewayId, requestMsg.getDeviceName(), e); + log.warn("[{}] Failed to process device [{}]", gatewayId, deviceName, e); throw new RuntimeException(e); } catch (EntitiesLimitExceededException e) { - log.warn("[{}][{}] API limit exception: [{}]", e.getTenantId(), gatewayId, e.getMessage()); - return TransportApiResponseMsg.newBuilder() - .setGetOrCreateDeviceResponseMsg( - GetOrCreateDeviceFromGatewayResponseMsg.newBuilder() - .setError(TransportProtos.TransportApiRequestErrorCode.ENTITY_LIMIT)) - .build(); + return buildLimitErrorResponse(e, gatewayId); } finally { - deviceCreationLock.unlock(); + lock.unlock(); + deviceCreationLocks.remove(deviceName, lock); + } + } + + private DeviceId toDeviceId(GetOrCreateDeviceFromGatewayRequestMsg requestMsg) { + return new DeviceId(new UUID( + requestMsg.getGatewayIdMSB(), + requestMsg.getGatewayIdLSB() + )); + } + + private Device findOrCreateDevice(GetOrCreateDeviceFromGatewayRequestMsg requestMsg, + Device gateway, + DeviceId gatewayId) throws JsonProcessingException { + TenantId tenantId = gateway.getTenantId(); + String deviceName = requestMsg.getDeviceName(); + Device device = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); + if (device != null) { + return device; + } + device = tryRenameSparkplugDevice(requestMsg, gateway); + if (device != null) { + return device; + } + device = createNewDevice(requestMsg, gateway, gatewayId); + pushCreatedEvent(device, gateway); + return device; + } + + private Device tryRenameSparkplugDevice(GetOrCreateDeviceFromGatewayRequestMsg requestMsg, + Device gateway) { + if (!requestMsg.getIsSparkplug()) { + return null; + } + String[] topicPath = requestMsg.getDeviceName().split(DEVICE_NAME_SPLIT_REGEXP); + if (topicPath.length != 3) { + return null; + } + String deviceId = topicPath[2]; + Device existingDevice = + deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), deviceId); + if (existingDevice == null) { + return null; + } + existingDevice.setName(requestMsg.getDeviceName()); + return deviceService.saveDevice(existingDevice); + } + + private Device createNewDevice(GetOrCreateDeviceFromGatewayRequestMsg requestMsg, + Device gateway, + DeviceId gatewayId) { + TenantId tenantId = gateway.getTenantId(); + Device device = new Device(); + device.setTenantId(tenantId); + device.setName(requestMsg.getDeviceName()); + device.setType(requestMsg.getDeviceType()); + device.setCustomerId(gateway.getCustomerId()); + DeviceProfile profile = + deviceProfileCache.findOrCreateDeviceProfile(tenantId, requestMsg.getDeviceType()); + device.setDeviceProfileId(profile.getId()); + ObjectNode additionalInfo = JacksonUtil.newObjectNode(); + additionalInfo.put(DataConstants.LAST_CONNECTED_GATEWAY, gatewayId.toString()); + device.setAdditionalInfo(additionalInfo); + device = deviceService.saveDevice(device); + relationService.saveRelation( + tenantId, + new EntityRelation(gateway.getId(), device.getId(), "Created") + ); + return device; + } + + private void updateLastConnectedGateway(Device device, DeviceId gatewayId) { + String gatewayIdStr = gatewayId.toString(); + JsonNode info = device.getAdditionalInfo(); + ObjectNode objectNode = (info instanceof ObjectNode) + ? (ObjectNode) info + : JacksonUtil.newObjectNode(); + if (!objectNode.has(DataConstants.LAST_CONNECTED_GATEWAY) + || !gatewayIdStr.equals(objectNode.get(DataConstants.LAST_CONNECTED_GATEWAY).asText())) { + objectNode.put(DataConstants.LAST_CONNECTED_GATEWAY, gatewayIdStr); + device.setAdditionalInfo(objectNode); + deviceService.saveDevice(device); } } + private void pushCreatedEvent(Device device, Device gateway) { + TenantId tenantId = gateway.getTenantId(); + CustomerId customerId = gateway.getCustomerId(); + TbMsgMetaData metaData = new TbMsgMetaData(); + metaData.putValue("gatewayId", gateway.getId().toString()); + if (customerId != null && !customerId.isNullUid()) { + metaData.putValue("customerId", customerId.toString()); + } + JsonNode entityNode = JacksonUtil.valueToTree(device); + TbMsg msg = TbMsg.newMsg() + .type(TbMsgType.ENTITY_CREATED) + .originator(device.getId()) + .customerId(customerId) + .copyMetaData(metaData) + .dataType(TbMsgDataType.JSON) + .data(JacksonUtil.toString(entityNode)) + .build(); + tbClusterService.pushMsgToRuleEngine(tenantId, device.getId(), msg, null); + } + + private TransportApiResponseMsg buildResponse(Device device) throws JsonProcessingException { + GetOrCreateDeviceFromGatewayResponseMsg.Builder builder = + GetOrCreateDeviceFromGatewayResponseMsg.newBuilder() + .setDeviceInfo(ProtoUtils.toDeviceInfoProto(device)); + DeviceProfile profile = + deviceProfileCache.get(device.getTenantId(), device.getDeviceProfileId()); + if (profile != null) { + builder.setDeviceProfile(ProtoUtils.toProto(profile)); + } + return TransportApiResponseMsg.newBuilder() + .setGetOrCreateDeviceResponseMsg(builder.build()) + .build(); + } + + private TransportApiResponseMsg buildLimitErrorResponse(EntitiesLimitExceededException e, + DeviceId gatewayId) { + log.warn("[{}][{}] API limit exception: [{}]", + e.getTenantId(), gatewayId, e.getMessage()); + return TransportApiResponseMsg.newBuilder() + .setGetOrCreateDeviceResponseMsg( + GetOrCreateDeviceFromGatewayResponseMsg.newBuilder() + .setError(TransportProtos.TransportApiRequestErrorCode.ENTITY_LIMIT) + ) + .build(); + } + private TransportApiResponseMsg handle(ProvisionDeviceRequestMsg requestMsg) { ProvisionResponse provisionResponse; try { 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 7d8770bbcf..9cca1a1a54 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 @@ -71,7 +71,9 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConn import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_REGEXP; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; /** * Created by nickAS21 on 12.01.23 @@ -103,37 +105,18 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte protected Set sparkplugAttributesMetricNames; public void beforeSparkplugTest(boolean isCreateDevices) throws Exception { + MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() + .gatewayName(edgeNodeDeviceName) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) + .transportPayloadType(TransportPayloadType.PROTOBUF) + .build(); + processBeforeTest(configProperties); if (isCreateDevices) { - MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() - .gatewayName(edgeNodeDeviceName) - .isSparkplug(true) - .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) - .transportPayloadType(TransportPayloadType.PROTOBUF) - .build(); - processBeforeTest(configProperties); - configProperties = MqttTestConfigProperties.builder() - .gatewayName(deviceId) - .isSparkplug(true) - .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) - .transportPayloadType(TransportPayloadType.PROTOBUF) - .build(); - processBeforeTest(configProperties); - configProperties = MqttTestConfigProperties.builder() - .gatewayName(groupId + ":" + edgeNode + ":" + deviceId) - .isSparkplug(true) - .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) - .transportPayloadType(TransportPayloadType.PROTOBUF) - .build(); - processBeforeTest(configProperties); - - } else { - MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() - .gatewayName(edgeNodeDeviceName) - .isSparkplug(true) - .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) - .transportPayloadType(TransportPayloadType.PROTOBUF) - .build(); - processBeforeTest(configProperties); + String deviceName = deviceId + "_1"; + createDevice(deviceName, deviceProfile.getName(), false); + deviceName = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceId + "_2"; + createDevice(deviceName, deviceProfile.getName(), false); } } @@ -175,7 +158,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte options.setSessionExpiryInterval(0L); options.setUserName(gatewayAccessToken); String nameSpace = nameSpaceBad.length == 0 ? TOPIC_ROOT_SPB_V_1_0 : nameSpaceBad[0]; - String topic = nameSpace + "/" + groupId + "/" + SparkplugMessageType.NDEATH.name() + "/" + edgeNode; + String topic = nameSpace + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NDEATH.name() + TOPIC_SPLIT_REGEXP + edgeNode; // The NDEATH message MUST set the MQTT Will QoS to 1 and Retained flag to false MqttMessage msg = new MqttMessage(); msg.setId(0); @@ -198,7 +181,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte payloadBirthNode.addMetrics(metric); payloadBirthNode.setTimestamp(ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } @@ -212,7 +195,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; payloadBirthDevice.addMetrics(metric); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceIdName, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceIdName, payloadBirthDevice.build().toByteArray(), 0, false); AtomicReference device = new AtomicReference<>(); await(alias + "find device [" + deviceIdName + "] after created") @@ -243,52 +226,53 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte payloadBirthNode.addMetrics(metric); payloadBirthNode.setTimestamp(ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } valueDeviceInt32 = 4024; metric = createMetric(valueDeviceInt32, ts, metricBirthName_Int32, metricBirthDataType_Int32, -1L); // as old device name -> deviceId - String deviceName = deviceId; - AtomicReference device1 = new AtomicReference<>(); - String finalDeviceName1 = deviceName; - await(alias + "find device [" + deviceId + "] before connecting") - .atMost(200, TimeUnit.SECONDS) - .until(() -> { - device1.set(doGet("/api/tenant/devices?deviceName=" + finalDeviceName1, Device.class)); - return device1.get() != null; - }); + String deviceiDName1 = deviceId + "_1"; if (client.isConnected()) { SparkplugBProto.Payload.Builder payloadBirthDevice1 = SparkplugBProto.Payload.newBuilder() .setTimestamp(ts) .setSeq(getSeqNum()); payloadBirthDevice1.addMetrics(metric); - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceId, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceiDName1, payloadBirthDevice1.build().toByteArray(), 0, false); - devices.add(device1.get()); + } - // as new device name -> groupId + ":" + edgeNode + ":" + deviceId; - deviceName = groupId + ":" + edgeNode + ":" + deviceId; - AtomicReference device2 = new AtomicReference<>(); - String finalDeviceName2 = deviceName; - await(alias + "find device [" + deviceName + "] before connecting") + String deviceName1 = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceiDName1;; + AtomicReference device1 = new AtomicReference<>(); + await(alias + "find device [" + deviceName1 + "] before connecting") .atMost(200, TimeUnit.SECONDS) .until(() -> { - device2.set(doGet("/api/tenant/devices?deviceName=" + finalDeviceName2, Device.class)); - return device2.get() != null; + device1.set(doGet("/api/tenant/devices?deviceName=" + deviceName1, Device.class)); + return device1.get() != null; }); + devices.add(device1.get()); + // as new device name -> groupId + ":" + edgeNode + ":" + deviceId; + String deviceiDName2 = deviceId + "_2"; if (client.isConnected()) { SparkplugBProto.Payload.Builder payloadBirthDevice2 = SparkplugBProto.Payload.newBuilder() .setTimestamp(ts) .setSeq(getSeqNum()); payloadBirthDevice2.addMetrics(metric); - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceId, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceiDName2, payloadBirthDevice2.build().toByteArray(), 0, false); - devices.add(device2.get()); } + String deviceName2 = groupId + ":" + edgeNode + ":" + deviceiDName2; + AtomicReference device2 = new AtomicReference<>(); + await(alias + "find device [" + deviceName2 + "] before connecting") + .atMost(200, TimeUnit.SECONDS) + .until(() -> { + device2.set(doGet("/api/tenant/devices?deviceName=" + deviceName2, Device.class)); + return device2.get() != null; + }); + devices.add(device2.get()); Assert.assertEquals(cntDevices, devices.size()); state_ONLINE_ALL (devices, calendar.getTimeInMillis()); } @@ -334,7 +318,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte payloadBirthNode.addMetrics(metric); payloadBirthNode.setTimestamp(ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } @@ -348,7 +332,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte payloadBirthDevice.addMetrics(metric); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceIdName, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceIdName, payloadBirthDevice.build().toByteArray(), 0, false); AtomicReference device = new AtomicReference<>(); await(alias + "find device [" + deviceName + "] after created") @@ -397,7 +381,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte listKeys.add(metricKey); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.NBIRTH.name() + "/" + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } return listKeys; diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java index 89b17f6b48..d0f3430fe7 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java @@ -37,10 +37,10 @@ import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConnectionState.OFFLINE; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConnectionState.ONLINE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; /** * Created by nickAS21 on 12.01.23 @@ -111,7 +111,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstra if (client.isConnected()) { List devicesList = new ArrayList<>(devices); Device device = devicesList.get(indexDeviceDisconnect); - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + SparkplugMessageType.DDEATH.name() + "/" + edgeNode + "/" + device.getName(), + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DDEATH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + device.getName(), payloadDeathDevice.build().toByteArray(), 0, false); await(alias + messageName(STATE) + ", device: " + device.getName()) .atMost(40, TimeUnit.SECONDS) diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java index 319b61b7f3..dc513421e6 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java @@ -30,6 +30,7 @@ import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; /** * Created by nickAS21 on 12.01.23 @@ -67,7 +68,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac createdAddMetricValuePrimitiveTsKv(listTsKvEntry, listKeys, ndataPayload, ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + messageTypeName + "/" + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + messageTypeName + TOPIC_SPLIT_REGEXP + edgeNode, ndataPayload.build().toByteArray(), 0, false); } @@ -96,7 +97,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac createdAddMetricValueArraysPrimitiveTsKv(listTsKvEntry, listKeys, ndataPayload, ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/" + messageTypeName + "/" + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + messageTypeName + TOPIC_SPLIT_REGEXP + edgeNode, ndataPayload.build().toByteArray(), 0, false); } diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 8a44a19e17..70b623f6c6 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -471,6 +471,7 @@ message GetOrCreateDeviceFromGatewayRequestMsg { int64 gatewayIdLSB = 2; string deviceName = 3; string deviceType = 4; + bool isSparkplug = 5; } message GetOrCreateDeviceFromGatewayResponseMsg { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index cbb805731e..799a0fdf21 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java @@ -261,7 +261,7 @@ public abstract class AbstractGatewaySessionHandler { ack(msg, MqttReasonCodes.PubAck.SUCCESS); log.trace("[{}][{}][{}] onDeviceConnectOk: [{}]", gateway.getTenantId(), gateway.getDeviceId(), sessionId, deviceName); @@ -277,7 +277,7 @@ public abstract class AbstractGatewaySessionHandler onDeviceConnect(String deviceName, String deviceType) { + ListenableFuture onDeviceConnect(String deviceName, String deviceType, boolean isSparkplug) { T result = devices.get(deviceName); if (result == null) { Lock deviceCreationLock = deviceCreationLockMap.computeIfAbsent(deviceName, s -> new ReentrantLock()); @@ -285,7 +285,7 @@ public abstract class AbstractGatewaySessionHandler onDeviceConnectSparkplug(SparkplugTopic topic, String deviceType) { T result = devices.get(topic.getNodeDeviceName()); if (result == null) { - return onDeviceConnect(topic.getNodeDeviceNameAllPath(), deviceType); + return onDeviceConnect(topic.getNodeDeviceNameAllPath(), deviceType, true); } else { return Futures.immediateFuture(result); } } - private ListenableFuture getDeviceCreationFuture(String deviceName, String deviceType) { + private ListenableFuture getDeviceCreationFuture(String deviceName, String deviceType, boolean isSparkplug) { final SettableFuture futureToSet = SettableFuture.create(); ListenableFuture future = deviceFutures.putIfAbsent(deviceName, futureToSet); if (future != null) { @@ -319,6 +319,7 @@ public abstract class AbstractGatewaySessionHandler() { @Override @@ -882,7 +883,7 @@ public abstract class AbstractGatewaySessionHandler onSuccess, Consumer onFailure) { - ListenableFuture deviceCtxFuture = onDeviceConnect(deviceName, DEFAULT_DEVICE_TYPE); + ListenableFuture deviceCtxFuture = onDeviceConnect(deviceName, DEFAULT_DEVICE_TYPE, false); process(deviceCtxFuture, onSuccess, onFailure); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java index 2bd0d7702f..2653ed8166 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugDeviceSessionContext.java @@ -128,5 +128,4 @@ public class SparkplugDeviceSessionContext extends AbstractGatewayDeviceSessionC rpcRequest.getMethodName() + ". " + e.getMessage()); } } - } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java index 910fed0f3a..b27b36e837 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java @@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.parseMessageType; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_REGEXP; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; @@ -332,9 +333,9 @@ public class SparkplugTopic { public String getNodeDeviceNameAllPath() { StringBuilder sb = new StringBuilder(); if (hostApplicationId == null) { - sb.append(getGroupId()).append(":").append(getEdgeNodeId()); + sb.append(getGroupId()).append(DEVICE_NAME_SPLIT_REGEXP).append(getEdgeNodeId()); if (getDeviceId() != null) { - sb.append(":").append(getDeviceId()); + sb.append(DEVICE_NAME_SPLIT_REGEXP).append(getDeviceId()); } } return sb.toString(); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java index 74ea6858fe..6558ba618f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java @@ -35,6 +35,7 @@ public class SparkplugTopicService { public static final String TOPIC_ROOT_SPB_V_1_0 = "spBv1.0"; public static final String TOPIC_ROOT_CERT_SP = "$sparkplug/certificates/"; public static final String TOPIC_SPLIT_REGEXP = "/"; + public static final String DEVICE_NAME_SPLIT_REGEXP = ":"; public static final String TOPIC_STATE_REGEXP = TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + STATE.name() + TOPIC_SPLIT_REGEXP; public static SparkplugTopic getSplitTopic(String topic) throws ThingsboardException { From a1134829772f3e1abfa66fd4bf858edf57515a51 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Sun, 15 Feb 2026 16:08:37 +0200 Subject: [PATCH 5/8] sparkplug - change date created new tets --- ...ttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java index 797b1beea5..e2adc4c11f 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java @@ -22,7 +22,7 @@ import org.junit.Test; import org.thingsboard.server.dao.service.DaoSqlTest; /** - * Created by nickAS21 on 12.01.23 + * Created by nickAS21 on 16.02.26 */ @DaoSqlTest public class MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest extends AbstractMqttV5ClientSparkplugConnectionTest { From aecf580e04d34b1a6394c31ed703bf168aa7a066 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Thu, 9 Apr 2026 17:56:08 +0300 Subject: [PATCH 6/8] spark[lug - add deviceId to label --- .../transport/DefaultTransportApiService.java | 5 +++++ .../AbstractMqttV5ClientSparkplugTest.java | 18 +++++++++++------- ...ugBConnectionDevicesCreatingBeforeTest.java | 5 +++++ 3 files changed, 21 insertions(+), 7 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index 1baa3cdccf..ad2d66d66a 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -392,6 +392,7 @@ public class DefaultTransportApiService implements TransportApiService { return null; } existingDevice.setName(requestMsg.getDeviceName()); + existingDevice.setLabel(deviceId); return deviceService.saveDevice(existingDevice); } @@ -402,6 +403,10 @@ public class DefaultTransportApiService implements TransportApiService { Device device = new Device(); device.setTenantId(tenantId); device.setName(requestMsg.getDeviceName()); + if (requestMsg.getIsSparkplug()){ + String [] topicDevice = requestMsg.getDeviceName().split(DEVICE_NAME_SPLIT_REGEXP); + if (topicDevice.length == 3) device.setLabel(topicDevice[2]); + } device.setType(requestMsg.getDeviceType()); device.setCustomerId(gateway.getCustomerId()); DeviceProfile profile = 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 9b6ef16e52..b52f38f55e 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 @@ -234,18 +234,18 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte valueDeviceInt32 = 4024; metric = createMetric(valueDeviceInt32, ts, metricBirthName_Int32, metricBirthDataType_Int32, -1L); // as old device name -> deviceId - String deviceiDName1 = deviceId + "_1"; + String deviceIdNameLabel1 = deviceId + "_1"; if (client.isConnected()) { SparkplugBProto.Payload.Builder payloadBirthDevice1 = SparkplugBProto.Payload.newBuilder() .setTimestamp(ts) .setSeq(getSeqNum()); payloadBirthDevice1.addMetrics(metric); - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceiDName1, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceIdNameLabel1, payloadBirthDevice1.build().toByteArray(), 0, false); } - String deviceName1 = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceiDName1;; + String deviceName1 = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceIdNameLabel1;; AtomicReference device1 = new AtomicReference<>(); await(alias + "find device [" + deviceName1 + "] before connecting") .atMost(200, TimeUnit.SECONDS) @@ -256,16 +256,16 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte devices.add(device1.get()); // as new device name -> groupId + ":" + edgeNode + ":" + deviceId; - String deviceiDName2 = deviceId + "_2"; + String deviceIdName2 = deviceId + "_2"; if (client.isConnected()) { SparkplugBProto.Payload.Builder payloadBirthDevice2 = SparkplugBProto.Payload.newBuilder() .setTimestamp(ts) .setSeq(getSeqNum()); payloadBirthDevice2.addMetrics(metric); - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceiDName2, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceIdName2, payloadBirthDevice2.build().toByteArray(), 0, false); } - String deviceName2 = groupId + ":" + edgeNode + ":" + deviceiDName2; + String deviceName2 = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceIdName2; AtomicReference device2 = new AtomicReference<>(); await(alias + "find device [" + deviceName2 + "] before connecting") .atMost(200, TimeUnit.SECONDS) @@ -276,6 +276,10 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte devices.add(device2.get()); Assert.assertEquals(cntDevices, devices.size()); 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 + Assert.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. + Assert.assertNull(device2.get().getLabel()); } protected void state_ONLINE_ALL (List devices, long ts) { @@ -328,7 +332,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() .setTimestamp(ts) .setSeq(getSeqNum()); - String deviceIdName = deviceId + "_" + 1; + String deviceIdName = deviceId + "_1"; String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; payloadBirthDevice.addMetrics(metric); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java index e2adc4c11f..a28b04d70e 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java @@ -27,6 +27,11 @@ import org.thingsboard.server.dao.service.DaoSqlTest; @DaoSqlTest public class MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest extends AbstractMqttV5ClientSparkplugConnectionTest { + /** + * String deviceName_1 = deviceId + "_1"; Only name device. Without a complete topic: how it was in the old version. + * String deviceName_2 = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceId + "_2"; With complete topic: how it was in the new version. + * @throws Exception + */ @Before public void beforeTest() throws Exception { beforeSparkplugTest(true); From 2ee1a1d3b56b51ec3d1403f8bc2ce406ecf834a3 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Mon, 11 May 2026 14:28:54 +0300 Subject: [PATCH 7/8] spark[lug - refactoring review - 01 (Without test) --- .../transport/DefaultTransportApiService.java | 74 +++++++++++++++---- .../AbstractMqttV5ClientSparkplugTest.java | 40 +++++----- ...ctMqttV5ClientSparkplugConnectionTest.java | 4 +- ...gBConnectionDevicesCreatingBeforeTest.java | 3 +- ...actMqttV5ClientSparkplugTelemetryTest.java | 6 +- .../AbstractGatewaySessionHandler.java | 29 ++++++-- .../SparkplugDeviceSessionContext.java | 1 + .../session/SparkplugNodeSessionHandler.java | 18 +++-- .../mqtt/util/sparkplug/SparkplugTopic.java | 10 +-- .../util/sparkplug/SparkplugTopicService.java | 6 +- 10 files changed, 129 insertions(+), 62 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index ad2d66d66a..09be48c641 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -22,6 +22,7 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; import com.google.protobuf.ByteString; +import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.exception.EntitiesLimitExceededException; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; @@ -114,7 +115,7 @@ import java.util.stream.Collectors; import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.PASSWORD_MISMATCH; import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.VALID; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_REGEXP; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_SEPARATOR; /** * Created by ashvayka on 05.10.18. @@ -347,7 +348,6 @@ public class DefaultTransportApiService implements TransportApiService { return buildLimitErrorResponse(e, gatewayId); } finally { lock.unlock(); - deviceCreationLocks.remove(deviceName, lock); } } @@ -367,45 +367,87 @@ public class DefaultTransportApiService implements TransportApiService { if (device != null) { return device; } - device = tryRenameSparkplugDevice(requestMsg, gateway); + String[] topicPath = requestMsg.getDeviceName().split(DEVICE_NAME_SPLIT_SEPARATOR); + device = tryRenameSparkplugDevice(requestMsg, gateway, topicPath); if (device != null) { return device; } - device = createNewDevice(requestMsg, gateway, gatewayId); + device = createNewDevice(requestMsg, gateway, gatewayId, topicPath); pushCreatedEvent(device, gateway); return device; } - private Device tryRenameSparkplugDevice(GetOrCreateDeviceFromGatewayRequestMsg requestMsg, - Device gateway) { + private Device tryRenameSparkplugDevice(GetOrCreateDeviceFromGatewayRequestMsg requestMsg, Device gateway, String[] topicPath) { if (!requestMsg.getIsSparkplug()) { return null; } - String[] topicPath = requestMsg.getDeviceName().split(DEVICE_NAME_SPLIT_REGEXP); + if (topicPath.length != 3) { return null; } + String deviceId = topicPath[2]; - Device existingDevice = - deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), deviceId); + Device existingDevice = deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), deviceId); + if (existingDevice == null) { return null; } - existingDevice.setName(requestMsg.getDeviceName()); - existingDevice.setLabel(deviceId); - return deviceService.saveDevice(existingDevice); + + // Security check: verify that the device was created by this gateway + try { + boolean isRelated = relationService.checkRelation( + gateway.getTenantId(), + gateway.getId(), + existingDevice.getId(), + "Created", + RelationTypeGroup.COMMON + ); + + if (!isRelated) { + log.warn("[{}] Security breach attempt! Gateway tried to rename device [{}] without 'Created' relation.", + gateway.getId(), existingDevice.getId()); + return null; + } + } catch (Exception e) { + log.error("[{}] Error checking relation for device {}", gateway.getId(), existingDevice.getId(), e); + return null; + } + + boolean changed = false; + String newName = requestMsg.getDeviceName(); + + 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); + + // Update label only if it's empty to avoid overwriting user changes + if (existingDevice.getLabel() == null || existingDevice.getLabel().isEmpty()) { + existingDevice.setLabel(deviceId); + } + + changed = true; + } + + return changed ? deviceService.saveDevice(existingDevice) : existingDevice; } private Device createNewDevice(GetOrCreateDeviceFromGatewayRequestMsg requestMsg, Device gateway, - DeviceId gatewayId) { + DeviceId gatewayId, String[] topicPath) { TenantId tenantId = gateway.getTenantId(); Device device = new Device(); device.setTenantId(tenantId); device.setName(requestMsg.getDeviceName()); - if (requestMsg.getIsSparkplug()){ - String [] topicDevice = requestMsg.getDeviceName().split(DEVICE_NAME_SPLIT_REGEXP); - if (topicDevice.length == 3) device.setLabel(topicDevice[2]); + if (requestMsg.getIsSparkplug()) { + if (topicPath.length == 3) device.setLabel(topicPath[2]); } device.setType(requestMsg.getDeviceType()); device.setCustomerId(gateway.getCustomerId()); 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 b52f38f55e..1ba5d1ba3a 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 @@ -71,9 +71,9 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConn import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_REGEXP; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_SEPARATOR; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_SEPARATOR; /** * Created by nickAS21 on 12.01.23 @@ -115,7 +115,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte if (isCreateDevices) { String deviceName = deviceId + "_1"; createDevice(deviceName, deviceProfile.getName(), false); - deviceName = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceId + "_2"; + deviceName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceId + "_2"; createDevice(deviceName, deviceProfile.getName(), false); } } @@ -158,7 +158,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte options.setSessionExpiryInterval(0L); options.setUserName(gatewayAccessToken); String nameSpace = nameSpaceBad.length == 0 ? TOPIC_ROOT_SPB_V_1_0 : nameSpaceBad[0]; - String topic = nameSpace + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NDEATH.name() + TOPIC_SPLIT_REGEXP + edgeNode; + String topic = nameSpace + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.NDEATH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode; // The NDEATH message MUST set the MQTT Will QoS to 1 and Retained flag to false MqttMessage msg = new MqttMessage(); msg.setId(0); @@ -181,7 +181,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte payloadBirthNode.addMetrics(metric); payloadBirthNode.setTimestamp(ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } @@ -192,14 +192,14 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte .setTimestamp(ts) .setSeq(getSeqNum()); String deviceIdName = deviceId + "_" + i; - String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; + String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; payloadBirthDevice.addMetrics(metric); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + 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, payloadBirthDevice.build().toByteArray(), 0, false); AtomicReference device = new AtomicReference<>(); await(alias + "find device [" + deviceIdName + "] after created") - .atMost(200, TimeUnit.SECONDS) + .atMost(40, TimeUnit.SECONDS) .ignoreExceptions() .until(() -> { device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); @@ -227,7 +227,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte payloadBirthNode.addMetrics(metric); payloadBirthNode.setTimestamp(ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } @@ -241,14 +241,14 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte .setTimestamp(ts) .setSeq(getSeqNum()); payloadBirthDevice1.addMetrics(metric); - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceIdNameLabel1, + 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 + deviceIdNameLabel1, payloadBirthDevice1.build().toByteArray(), 0, false); } - String deviceName1 = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceIdNameLabel1;; + String deviceName1 = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceIdNameLabel1; AtomicReference device1 = new AtomicReference<>(); await(alias + "find device [" + deviceName1 + "] before connecting") - .atMost(200, TimeUnit.SECONDS) + .atMost(40, TimeUnit.SECONDS) .until(() -> { device1.set(doGet("/api/tenant/devices?deviceName=" + deviceName1, Device.class)); return device1.get() != null; @@ -262,13 +262,13 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte .setTimestamp(ts) .setSeq(getSeqNum()); payloadBirthDevice2.addMetrics(metric); - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + deviceIdName2, + 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 + deviceIdName2, payloadBirthDevice2.build().toByteArray(), 0, false); } - String deviceName2 = groupId + DEVICE_NAME_SPLIT_REGEXP + edgeNode + DEVICE_NAME_SPLIT_REGEXP + deviceIdName2; + String deviceName2 = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceIdName2; AtomicReference device2 = new AtomicReference<>(); await(alias + "find device [" + deviceName2 + "] before connecting") - .atMost(200, TimeUnit.SECONDS) + .atMost(40, TimeUnit.SECONDS) .until(() -> { device2.set(doGet("/api/tenant/devices?deviceName=" + deviceName2, Device.class)); return device2.get() != null; @@ -323,7 +323,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte payloadBirthNode.addMetrics(metric); payloadBirthNode.setTimestamp(ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } @@ -333,15 +333,15 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte .setTimestamp(ts) .setSeq(getSeqNum()); String deviceIdName = deviceId + "_1"; - String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; + String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; payloadBirthDevice.addMetrics(metric); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + 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, payloadBirthDevice.build().toByteArray(), 0, false); AtomicReference device = new AtomicReference<>(); await(alias + "find device [" + deviceName + "] after created") - .atMost(200, TimeUnit.SECONDS) + .atMost(40, TimeUnit.SECONDS) .ignoreExceptions() .until(() -> { device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); @@ -387,7 +387,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte listKeys.add(metricKey); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_REGEXP + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.NBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode, payloadBirthNode.build().toByteArray(), 0, false); } return listKeys; diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java index d0f3430fe7..1402e5c4c3 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/AbstractMqttV5ClientSparkplugConnectionTest.java @@ -40,7 +40,7 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConn import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_SEPARATOR; /** * Created by nickAS21 on 12.01.23 @@ -111,7 +111,7 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstra if (client.isConnected()) { List devicesList = new ArrayList<>(devices); Device device = devicesList.get(indexDeviceDisconnect); - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + SparkplugMessageType.DDEATH.name() + TOPIC_SPLIT_REGEXP + edgeNode + TOPIC_SPLIT_REGEXP + device.getName(), + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.DDEATH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode + TOPIC_SPLIT_SEPARATOR + device.getName(), payloadDeathDevice.build().toByteArray(), 0, false); await(alias + messageName(STATE) + ", device: " + device.getName()) .atMost(40, TimeUnit.SECONDS) diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java index a28b04d70e..720e829371 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java @@ -47,4 +47,5 @@ public class MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest extends A @Test public void testClientWithCorrectAccessTokenWithNDEATHTwoDevicesCreatingBeforeFirstNameDeviceIdSecondNameFull() throws Exception { connectClientWithCorrectAccessTokenWithNDEATHDevicesCreatingBefore_Test(2); - }} + } +} diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java index dc513421e6..8f2d0659af 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/AbstractMqttV5ClientSparkplugTelemetryTest.java @@ -30,7 +30,7 @@ import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_SEPARATOR; /** * Created by nickAS21 on 12.01.23 @@ -68,7 +68,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac createdAddMetricValuePrimitiveTsKv(listTsKvEntry, listKeys, ndataPayload, ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + messageTypeName + TOPIC_SPLIT_REGEXP + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + messageTypeName + TOPIC_SPLIT_SEPARATOR + edgeNode, ndataPayload.build().toByteArray(), 0, false); } @@ -97,7 +97,7 @@ public abstract class AbstractMqttV5ClientSparkplugTelemetryTest extends Abstrac createdAddMetricValueArraysPrimitiveTsKv(listTsKvEntry, listKeys, ndataPayload, ts); if (client.isConnected()) { - client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + groupId + TOPIC_SPLIT_REGEXP + messageTypeName + TOPIC_SPLIT_REGEXP + edgeNode, + client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + messageTypeName + TOPIC_SPLIT_SEPARATOR + edgeNode, ndataPayload.build().toByteArray(), 0, false); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index 799a0fdf21..94c24f5a1d 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java @@ -298,11 +298,22 @@ public abstract class AbstractGatewaySessionHandler onDeviceConnectSparkplug(SparkplugTopic topic, String deviceType) { - T result = devices.get(topic.getNodeDeviceName()); + String fullPath = topic.getNodeDeviceNameAllPath(); + // Primary lookup: try to find the device by its full-path name (standard for new devices) + T result = devices.get(fullPath); + if (result == null) { - return onDeviceConnect(topic.getNodeDeviceNameAllPath(), deviceType, true); - } else { + // Secondary lookup (Legacy Fallback): check for the short name if full path is not found. + // This supports devices migrated from older versions. + String shortName = topic.getNodeDeviceName(); + result = devices.get(shortName); + } + + if (result != null) { return Futures.immediateFuture(result); + } else { + // If not found in cache at all, proceed with connection/creation using full path + return onDeviceConnect(fullPath, deviceType, true); } } @@ -855,8 +866,16 @@ public abstract class AbstractGatewaySessionHandler deviceCtx = this.onDeviceConnectProto(topic); - deviceName = checkDeviceName(deviceCtx.get().getDeviceInfo().getDeviceName()); String finalDeviceName = deviceName; contextListenableFuture = Futures.transform(deviceCtx, ctx -> { if (topic.isType(DBIRTH)) { sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), finalDeviceName, ONLINE, sparkplugBProto.getTimestamp()); + try { ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); + } catch (IllegalArgumentException | DuplicateKeyException e) { + log.error("[{}] Failed to set birth metrics", finalDeviceName, e); + throw new RuntimeException(e); + } } return ctx; }, MoreExecutors.directExecutor()); - } catch (IllegalArgumentException | DuplicateKeyException | ExecutionException | InterruptedException e) { + } catch (IllegalArgumentException | DuplicateKeyException e) { throw new RuntimeException(e); } } @@ -200,7 +204,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler= 4 && splitTopic.length <= 5 && splitTopic[0].equals(this.sparkplugTopicNode.getNamespace()) && splitTopic[1].equals(this.sparkplugTopicNode.getGroupId()) && diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java index b27b36e837..9b0cb62ca6 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopic.java @@ -21,9 +21,9 @@ import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.parseMessageType; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_REGEXP; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.DEVICE_NAME_SPLIT_SEPARATOR; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_ROOT_SPB_V_1_0; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_SEPARATOR; /** * Created by nickAS21 on 12.12.22 @@ -197,7 +197,7 @@ public class SparkplugTopic { try { if (isValidIdElementToUTF8(topicString)) { SparkplugMessageType messageType; - String[] splitTopic = topicString.split(TOPIC_SPLIT_REGEXP); + String[] splitTopic = topicString.split(TOPIC_SPLIT_SEPARATOR); if (TOPIC_ROOT_SPB_V_1_0.equals(splitTopic[0])) { if (splitTopic.length == 3) { messageType = parseMessageType(splitTopic[1]); @@ -333,9 +333,9 @@ public class SparkplugTopic { public String getNodeDeviceNameAllPath() { StringBuilder sb = new StringBuilder(); if (hostApplicationId == null) { - sb.append(getGroupId()).append(DEVICE_NAME_SPLIT_REGEXP).append(getEdgeNodeId()); + sb.append(getGroupId()).append(DEVICE_NAME_SPLIT_SEPARATOR).append(getEdgeNodeId()); if (getDeviceId() != null) { - sb.append(DEVICE_NAME_SPLIT_REGEXP).append(getDeviceId()); + sb.append(DEVICE_NAME_SPLIT_SEPARATOR).append(getDeviceId()); } } return sb.toString(); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java index 6558ba618f..377b759957 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugTopicService.java @@ -34,9 +34,9 @@ public class SparkplugTopicService { private static final Map SPLIT_TOPIC_CACHE = new HashMap<>(); public static final String TOPIC_ROOT_SPB_V_1_0 = "spBv1.0"; public static final String TOPIC_ROOT_CERT_SP = "$sparkplug/certificates/"; - public static final String TOPIC_SPLIT_REGEXP = "/"; - public static final String DEVICE_NAME_SPLIT_REGEXP = ":"; - public static final String TOPIC_STATE_REGEXP = TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_REGEXP + STATE.name() + TOPIC_SPLIT_REGEXP; + public static final String TOPIC_SPLIT_SEPARATOR = "/"; + public static final String DEVICE_NAME_SPLIT_SEPARATOR = ":"; + public static final String TOPIC_STATE_SEPARATOR = TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + STATE.name() + TOPIC_SPLIT_SEPARATOR; public static SparkplugTopic getSplitTopic(String topic) throws ThingsboardException { SparkplugTopic sparkplugTopic = SPLIT_TOPIC_CACHE.get(topic); From 6570882b84e21cd896cd079df2f1d131b1fb1116 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Mon, 11 May 2026 17:49:40 +0300 Subject: [PATCH 8/8] spark[lug - refactoring review - 01 (With test) --- .../transport/DefaultTransportApiService.java | 20 +- .../AbstractMqttV5ClientSparkplugTest.java | 186 +++++++++++++++++- ...gBConnectionDevicesCreatingBeforeTest.java | 20 ++ 3 files changed, 215 insertions(+), 11 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index 09be48c641..ed43c71027 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -394,25 +394,31 @@ public class DefaultTransportApiService implements TransportApiService { } // Security check: verify that the device was created by this gateway + boolean isRelated = false; try { - boolean isRelated = relationService.checkRelation( + // Security check: verify that the device was originally created by this gateway + isRelated = relationService.checkRelation( gateway.getTenantId(), gateway.getId(), existingDevice.getId(), "Created", RelationTypeGroup.COMMON ); - - if (!isRelated) { - log.warn("[{}] Security breach attempt! Gateway tried to rename device [{}] without 'Created' relation.", - gateway.getId(), existingDevice.getId()); - return null; - } } 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 + 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 boolean changed = false; String newName = requestMsg.getDeviceName(); 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 1ba5d1ba3a..2adc7b08a5 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 @@ -31,7 +31,9 @@ import org.junit.Assert; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.StringUtils; 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.id.DeviceId; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry; @@ -39,6 +41,7 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties; @@ -57,6 +60,7 @@ import java.util.concurrent.atomic.AtomicReference; import static org.awaitility.Awaitility.await; import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.thingsboard.common.util.JacksonUtil.newArrayNode; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Bytes; import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int16; @@ -113,10 +117,22 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte .build(); processBeforeTest(configProperties); if (isCreateDevices) { - String deviceName = deviceId + "_1"; - createDevice(deviceName, deviceProfile.getName(), false); - deviceName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceId + "_2"; - createDevice(deviceName, deviceProfile.getName(), false); + // 1. Create the first device with a short name (legacy style) + String deviceName1 = deviceId + "_1"; + Device device1 = createDevice(deviceName1, deviceProfile.getName(), false); + + // 2. Establish 'Created' relation so the transport identifies this gateway as the owner + String relationType = "Created"; + EntityRelation relation1 = createFromRelation(savedGateway, device1, relationType); + 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"; + Device device2 = createDevice(deviceName2, deviceProfile.getName(), false); + + // 4. Establish 'Created' relation for the second device as well + EntityRelation relation2 = createFromRelation(savedGateway, device2, relationType); + doPost("/api/relation", relation2).andExpect(status().isOk()); } } @@ -282,6 +298,111 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte Assert.assertNull(device2.get().getLabel()); } + /** + * Coverage: Rename when a device with the target full-path name already exists (collision). + */ + protected void renameCollisionWhenTargetNameAlreadyExists_Test() throws Exception { + long ts = calendar.getTimeInMillis(); + String shortName = deviceId + "_1"; // Created in beforeTest + String fullPathName = groupId + ":" + edgeNode + ":" + shortName; + + // Manually create a device that already has the "new" full-path name to trigger a collision + createDevice(fullPathName, deviceProfile.getName(), false); + + clientWithCorrectNodeAccessTokenWithNDEATH(); + + SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts) + .setSeq(getSeqNum()); + payload.addMetrics(createMetric(123, ts, metricBirthName_Int32, metricBirthDataType_Int32, -1L)); + + // Gateway sends DBIRTH for the short name. + // 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, + payload.build().toByteArray(), 0, false); + + await("Checking stability after collision") + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + Device oldDevice = doGet("/api/tenant/devices?deviceName=" + shortName, Device.class); + Device conflictDevice = doGet("/api/tenant/devices?deviceName=" + fullPathName, Device.class); + // Both devices must still exist, proving no exception crashed the process + return oldDevice != null && conflictDevice != null; + }); + } + + /** + * Coverage: The privilege concern — attempt to rename a device not owned by the gateway. + * This test verifies that the original device's ID remains unchanged, meaning it was not hijacked. + */ + protected void unauthorizedRenameAttemptBad_Test() throws Exception { + long ts = calendar.getTimeInMillis(); + String strangerName = "unauthorized_device_rename"; + + // 1. Create a "stranger" device via API (it has no 'Created' relation to the gateway) + Device stranger = new Device(); + stranger.setName(strangerName); + stranger.setType("default"); + doPost("/api/device", stranger); + final DeviceId originalStrangerId = stranger.getId(); + + clientWithCorrectNodeAccessTokenWithNDEATH(); + + SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts).setSeq(getSeqNum()); + + // 2. Unauthorized gateway attempts to rename this device via Sparkplug topic path + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + strangerName, + payload.build().toByteArray(), 0, false); + + String expectedFullPath = groupId + ":" + edgeNode + ":" + strangerName; + + // 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") + .atMost(40, TimeUnit.SECONDS) + .pollDelay(2, TimeUnit.SECONDS) + .untilAsserted(() -> { + // Check if the original device still exists with its original ID + Device currentStranger = doGet("/api/tenant/devices?deviceName=" + strangerName, Device.class); + Assert.assertNotNull("Original device disappeared!", currentStranger); + Assert.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 + Device newDevice = doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class); + if (newDevice != null) { + Assert.assertNotEquals("Stranger device was successfully hijacked (IDs match)!", originalStrangerId, newDevice.getId()); + } + }); + } + + /** + * Coverage: The privilege concern — attempt to rename a device not owned by the gateway. + */ + protected void unauthorizedRenameAttempt_Test() throws Exception { + long ts = calendar.getTimeInMillis(); + String strangerName = "unauthorized_device_rename"; + + // Create a device without a "Created" relation to the gateway + Device stranger = new Device(); + stranger.setName(strangerName); + stranger.setType("default"); + doPost("/api/device", stranger); + + clientWithCorrectNodeAccessTokenWithNDEATH(); + + SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts).setSeq(getSeqNum()); + + // Unauthorized gateway attempts to rename the device via Sparkplug topic + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + strangerName, + payload.build().toByteArray(), 0, false); + + String expectedFullPath = groupId + ":" + edgeNode + ":" + strangerName; + await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> + doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class, status().isNotFound()) + ); + } + protected void state_ONLINE_ALL (List devices, long ts) { TsKvEntry tsKvEntry = new BasicTsKvEntry(ts, new StringDataEntry(messageName(STATE), ONLINE.name())); await(alias + messageName(STATE) + ", device: " + savedGateway.getName()) @@ -309,6 +430,58 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte } } + /** + * Coverage: Concurrent first-message registration with the lock mechanism. + */ + protected void concurrentFirstMessageRegistration_Test() throws Exception { + int threadCount = 5; + String concurrentDeviceName = "concurrent_device"; + clientWithCorrectNodeAccessTokenWithNDEATH(); + + java.util.concurrent.ExecutorService executor = java.util.concurrent.Executors.newFixedThreadPool(threadCount); + long ts = calendar.getTimeInMillis(); + + for (int i = 0; i < threadCount; i++) { + executor.submit(() -> { + try { + SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder() + .setTimestamp(ts).setSeq(0); + client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + concurrentDeviceName, + payload.build().toByteArray(), 0, false); + } catch (Exception e) { + log.error("Concurrent publish failed", e); + } + }); + } + + String expectedName = groupId + ":" + edgeNode + ":" + concurrentDeviceName; + await("Wait for concurrent registration result") + .atMost(40, TimeUnit.SECONDS) // Restored to 40s as requested + .until(() -> doGet("/api/tenant/devices?deviceName=" + expectedName, Device.class) != null); + + executor.shutdown(); + } + + /** + * Coverage: Sparkplug-message handling when msgId <= 0 (#7). + * Verifies that the transport does not close the session for Sparkplug clients using msgId 0. + */ + protected void sparkplugSessionStaysAliveWithZeroMsgId_Test() throws Exception { + // clientMqttV5ConnectWithNDEATH internally sets msgId = 0 for the Will message. + // This validates that the connection is accepted despite msgId being 0. + IMqttToken connectionResult = clientMqttV5ConnectWithNDEATH(calendar.getTimeInMillis(), 0, -1L); + Assert.assertTrue("Sparkplug connection should be successful with msgId=0", client.isConnected()); + + // Publish NBIRTH message which usually goes through the aggregate callback. + // This verifies that msgId=0 in the callback does not trigger closeDeviceSession. + connectionWithNBirth(Int32, "test_metric_msgId_0", 555); + + // Awaitility to ensure the session remains open after processing. + await("Verify Sparkplug session remains open after receiving msgId=0") + .atMost(40, TimeUnit.SECONDS) + .until(() -> client.isConnected()); + } + protected List connectClientWithCorrectAccessTokenWithNDEATHWithAliasCreatedDevices(long ts) throws Exception { List devices = new ArrayList<>(); Long alias = 0L; @@ -630,4 +803,9 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte } } + private EntityRelation createFromRelation(Device mainDevice, Device device, String relationType) { + return new EntityRelation(mainDevice.getId(), device.getId(), relationType); + } + + } diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java index 720e829371..92583e655c 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java @@ -48,4 +48,24 @@ public class MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest extends A public void testClientWithCorrectAccessTokenWithNDEATHTwoDevicesCreatingBeforeFirstNameDeviceIdSecondNameFull() throws Exception { connectClientWithCorrectAccessTokenWithNDEATHDevicesCreatingBefore_Test(2); } + + @Test + public void testRenameWhenDeviceFullPathAlreadyExists_Collision() throws Exception { + renameCollisionWhenTargetNameAlreadyExists_Test(); + } + + @Test + public void testUnauthorizedRenameAttempt() throws Exception { + unauthorizedRenameAttempt_Test(); + } + + @Test + public void testConcurrentFirstMessageRegistration() throws Exception { + concurrentFirstMessageRegistration_Test(); + } + + @Test + public void testSparkplugSessionStaysAliveWithZeroMsgId() throws Exception { + sparkplugSessionStaysAliveWithZeroMsgId_Test(); + } }