Browse Source

sparkplug-unique-device-names - comments -2

pull/15753/head
nickAS21 4 months ago
parent
commit
4bfc810d75
  1. 38
      application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
  2. 67
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java
  3. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesInProfileTest.java
  4. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/MqttV5ClientSparkplugBAttributesTest.java
  5. 3
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionDevicesCreatingBeforeTest.java
  6. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/connection/MqttV5ClientSparkplugBConnectionTest.java
  7. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/MqttV5RpcSparkplugTest.java
  8. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/timeseries/MqttV5ClientSparkplugBTelemetryTest.java
  9. 8
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java
  10. 35
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java

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

@ -396,7 +396,6 @@ public class DefaultTransportApiService implements TransportApiService {
// Security check: verify that the device was created by this gateway // Security check: verify that the device was created by this gateway
boolean isRelated = false; boolean isRelated = false;
try { try {
// Security check: verify that the device was originally created by this gateway
isRelated = relationService.checkRelation( isRelated = relationService.checkRelation(
gateway.getTenantId(), gateway.getTenantId(),
gateway.getId(), gateway.getId(),
@ -405,37 +404,34 @@ public class DefaultTransportApiService implements TransportApiService {
RelationTypeGroup.COMMON RelationTypeGroup.COMMON
); );
} catch (Exception e) { } catch (Exception e) {
// Log the error from the relation service but return null to allow potential recovery log.error("[{}] Failed checking relation for device {}",
log.error("[{}] Error checking relation for device {}", gateway.getId(), existingDevice.getId(), e); gateway.getId(),
return null; existingDevice.getId(),
e);
throw new RuntimeException(
"Failed checking relation for device " + existingDevice.getId(),
e
);
} }
// If the device is found but not related to this gateway, it's a security breach // If the device is found but not related to this gateway
if (!isRelated) { if (!isRelated) {
log.error("[{}] Security breach attempt! Gateway tried to rename device [{}] without 'Created' relation.", log.debug("[{}] Device [{}] exists but is not related to gateway. " +
gateway.getId(), existingDevice.getId()); "Skipping rename and allowing creation of Sparkplug device [{}].",
// Throwing exception to halt the entire connection process gateway.getId(),
throw new RuntimeException("Security breach attempt! Unauthorized device rename."); existingDevice.getId(),
requestMsg.getDeviceName());
return null;
} }
// Logic for renaming the device if it's related and no naming conflicts exist // Logic for renaming the device if it's related and no naming conflicts exist
boolean changed = false; boolean changed = false;
String newName = requestMsg.getDeviceName(); String newName = requestMsg.getDeviceName();
if (!newName.equals(existingDevice.getName())) { if (!newName.equals(existingDevice.getName())) {
// Check if the new name is already taken by another device
Device conflictDevice = deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), newName);
if (conflictDevice != null) {
log.warn("[{}] Cannot rename device [{}] to [{}]: name already exists!",
gateway.getId(), existingDevice.getId(), newName);
return existingDevice;
}
existingDevice.setName(newName); existingDevice.setName(newName);
// Update label only if it's empty to avoid overwriting user changes if (StringUtils.isEmpty(existingDevice.getLabel())) {
if (existingDevice.getLabel() == null || existingDevice.getLabel().isEmpty()) {
existingDevice.setLabel(deviceId); existingDevice.setLabel(deviceId);
} }

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

@ -31,7 +31,6 @@ import org.junit.Assert;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.common.data.asset.AssetInfo;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
@ -58,6 +57,7 @@ import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
import static java.util.concurrent.Executors.newFixedThreadPool;
import static org.awaitility.Awaitility.await; import static org.awaitility.Awaitility.await;
import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@ -108,7 +108,8 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
protected static final String metricBirthName_Int32 = "Device Metric int32"; protected static final String metricBirthName_Int32 = "Device Metric int32";
protected Set<String> sparkplugAttributesMetricNames; protected Set<String> sparkplugAttributesMetricNames;
public void beforeSparkplugTest(boolean isCreateDevices) throws Exception { public void beforeSparkplugTest() throws Exception {
MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder()
.gatewayName(edgeNodeDeviceName) .gatewayName(edgeNodeDeviceName)
.isSparkplug(true) .isSparkplug(true)
@ -116,24 +117,26 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.transportPayloadType(TransportPayloadType.PROTOBUF) .transportPayloadType(TransportPayloadType.PROTOBUF)
.build(); .build();
processBeforeTest(configProperties); processBeforeTest(configProperties);
if (isCreateDevices) { }
// 1. Create the first device with a short name (legacy style)
String deviceName1 = deviceId + "_1"; public void seedLegacyAndFullPathDevices() throws Exception {
Device device1 = createDevice(deviceName1, deviceProfile.getName(), false); // 1. Create the first device with a short name (legacy style)
String deviceName1 = deviceId + "_1";
// 2. Establish 'Created' relation so the transport identifies this gateway as the owner Device device1 = createDevice(deviceName1, deviceProfile.getName(), false);
String relationType = "Created";
EntityRelation relation1 = createFromRelation(savedGateway, device1, relationType); // 2. Establish 'Created' relation so the transport identifies this gateway as the owner
doPost("/api/relation", relation1).andExpect(status().isOk()); String relationType = "Created";
EntityRelation relation1 = createFromRelation(savedGateway, device1, relationType);
// 3. Create the second device with a full-path name doPost("/api/relation", relation1).andExpect(status().isOk());
String deviceName2 = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceId + "_2";
Device device2 = createDevice(deviceName2, deviceProfile.getName(), false); // 3. Create the second device with a full-path name
String deviceName2 = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceId + "_2";
// 4. Establish 'Created' relation for the second device as well Device device2 = createDevice(deviceName2, deviceProfile.getName(), false);
EntityRelation relation2 = createFromRelation(savedGateway, device2, relationType);
doPost("/api/relation", relation2).andExpect(status().isOk()); // 4. Establish 'Created' relation for the second device as well
} EntityRelation relation2 = createFromRelation(savedGateway, device2, relationType);
doPost("/api/relation", relation2).andExpect(status().isOk());
} }
public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception { public void clientWithCorrectNodeAccessTokenWithNDEATH() throws Exception {
@ -208,7 +211,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts) .setTimestamp(ts)
.setSeq(getSeqNum()); .setSeq(getSeqNum());
String deviceIdName = deviceId + "_" + i; String deviceIdName = deviceId + "_" + i;
String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; String deviceName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceIdName;
payloadBirthDevice.addMetrics(metric); payloadBirthDevice.addMetrics(metric);
if (client.isConnected()) { if (client.isConnected()) {
client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode + TOPIC_SPLIT_SEPARATOR + deviceIdName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + TOPIC_SPLIT_SEPARATOR + SparkplugMessageType.DBIRTH.name() + TOPIC_SPLIT_SEPARATOR + edgeNode + TOPIC_SPLIT_SEPARATOR + deviceIdName,
@ -304,7 +307,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
protected void renameCollisionWhenTargetNameAlreadyExists_Test() throws Exception { protected void renameCollisionWhenTargetNameAlreadyExists_Test() throws Exception {
long ts = calendar.getTimeInMillis(); long ts = calendar.getTimeInMillis();
String shortName = deviceId + "_1"; // Created in beforeTest String shortName = deviceId + "_1"; // Created in beforeTest
String fullPathName = groupId + ":" + edgeNode + ":" + shortName; String fullPathName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + shortName;
// Manually create a device that already has the "new" full-path name to trigger a collision // Manually create a device that already has the "new" full-path name to trigger a collision
createDevice(fullPathName, deviceProfile.getName(), false); createDevice(fullPathName, deviceProfile.getName(), false);
@ -318,7 +321,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
// Gateway sends DBIRTH for the short name. // Gateway sends DBIRTH for the short name.
// Transport will try to rename it but should find a conflict and handle it gracefully. // Transport will try to rename it but should find a conflict and handle it gracefully.
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + shortName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + shortName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
await("Checking stability after collision") await("Checking stability after collision")
@ -352,10 +355,10 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts).setSeq(getSeqNum()); .setTimestamp(ts).setSeq(getSeqNum());
// 2. Unauthorized gateway attempts to rename this device via Sparkplug topic path // 2. Unauthorized gateway attempts to rename this device via Sparkplug topic path
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + strangerName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + strangerName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
String expectedFullPath = groupId + ":" + edgeNode + ":" + strangerName; String expectedFullPath = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + strangerName;
// 3. Verify security: the original device must still be linked to its short name with the same ID // 3. Verify security: the original device must still be linked to its short name with the same ID
await("Verify original device was not hijacked") await("Verify original device was not hijacked")
@ -394,10 +397,10 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts).setSeq(getSeqNum()); .setTimestamp(ts).setSeq(getSeqNum());
// Unauthorized gateway attempts to rename the device via Sparkplug topic // Unauthorized gateway attempts to rename the device via Sparkplug topic
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + strangerName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + strangerName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
String expectedFullPath = groupId + ":" + edgeNode + ":" + strangerName; String expectedFullPath = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + strangerName;
await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> await().atMost(30, TimeUnit.SECONDS).untilAsserted(() ->
doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class, status().isNotFound()) doGet("/api/tenant/devices?deviceName=" + expectedFullPath, Device.class, status().isNotFound())
); );
@ -438,7 +441,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
String concurrentDeviceName = "concurrent_device"; String concurrentDeviceName = "concurrent_device";
clientWithCorrectNodeAccessTokenWithNDEATH(); clientWithCorrectNodeAccessTokenWithNDEATH();
java.util.concurrent.ExecutorService executor = java.util.concurrent.Executors.newFixedThreadPool(threadCount); java.util.concurrent.ExecutorService executor = newFixedThreadPool(threadCount);
long ts = calendar.getTimeInMillis(); long ts = calendar.getTimeInMillis();
for (int i = 0; i < threadCount; i++) { for (int i = 0; i < threadCount; i++) {
@ -446,7 +449,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
try { try {
SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder() SparkplugBProto.Payload.Builder payload = SparkplugBProto.Payload.newBuilder()
.setTimestamp(ts).setSeq(0); .setTimestamp(ts).setSeq(0);
client.publish(TOPIC_ROOT_SPB_V_1_0 + "/" + groupId + "/DBIRTH/" + edgeNode + "/" + concurrentDeviceName, client.publish(TOPIC_ROOT_SPB_V_1_0 + TOPIC_SPLIT_SEPARATOR + groupId + "/DBIRTH/" + edgeNode + TOPIC_SPLIT_SEPARATOR + concurrentDeviceName,
payload.build().toByteArray(), 0, false); payload.build().toByteArray(), 0, false);
} catch (Exception e) { } catch (Exception e) {
log.error("Concurrent publish failed", e); log.error("Concurrent publish failed", e);
@ -454,9 +457,9 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
}); });
} }
String expectedName = groupId + ":" + edgeNode + ":" + concurrentDeviceName; String expectedName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + concurrentDeviceName;
await("Wait for concurrent registration result") await("Wait for concurrent registration result")
.atMost(40, TimeUnit.SECONDS) // Restored to 40s as requested .atMost(40, TimeUnit.SECONDS)
.until(() -> doGet("/api/tenant/devices?deviceName=" + expectedName, Device.class) != null); .until(() -> doGet("/api/tenant/devices?deviceName=" + expectedName, Device.class) != null);
executor.shutdown(); executor.shutdown();
@ -506,7 +509,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte
.setTimestamp(ts) .setTimestamp(ts)
.setSeq(getSeqNum()); .setSeq(getSeqNum());
String deviceIdName = deviceId + "_1"; String deviceIdName = deviceId + "_1";
String deviceName = groupId + ":" + edgeNode + ":" + deviceIdName; String deviceName = groupId + DEVICE_NAME_SPLIT_SEPARATOR + edgeNode + DEVICE_NAME_SPLIT_SEPARATOR + deviceIdName;
payloadBirthDevice.addMetrics(metric); payloadBirthDevice.addMetrics(metric);
if (client.isConnected()) { if (client.isConnected()) {

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Loading…
Cancel
Save