From 3477ddac0fc35793d14ad704d6daa0a1ffa08936 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Wed, 8 Mar 2023 17:39:39 +0200 Subject: [PATCH] Refactoring due to comments --- .../transport/DefaultTransportApiService.java | 12 ++--- .../mqtt/AbstractMqttIntegrationTest.java | 4 +- .../mqtt/MqttTestConfigProperties.java | 4 +- .../AbstractMqttV5ClientSparkplugTest.java | 6 +-- ...ctMqttV5ClientSparkplugAttributesTest.java | 5 -- ...ientSparkplugBAttributesInProfileTest.java | 4 +- .../server/common/data/DataConstants.java | 1 - ...ttDeviceProfileTransportConfiguration.java | 4 +- .../transport/mqtt/MqttTransportHandler.java | 40 ++-------------- .../AbstractGatewaySessionHandler.java | 48 +++++++++---------- .../MqttDeviceAwareSessionContext.java | 14 ------ .../SparkplugDeviceSessionContext.java | 29 ++++++----- .../session/SparkplugNodeSessionHandler.java | 37 +++++++------- .../session/DeviceAwareSessionContext.java | 2 +- ...ile-transport-configuration.component.html | 12 ++--- ...ofile-transport-configuration.component.ts | 22 ++++----- ui-ngx/src/app/shared/models/device.models.ts | 4 +- 17 files changed, 98 insertions(+), 150 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 23f6a4a08c..0c0ff3e0c9 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 @@ -47,8 +47,6 @@ import org.thingsboard.server.common.data.device.data.CoapDeviceTransportConfigu import org.thingsboard.server.common.data.device.data.Lwm2mDeviceTransportConfiguration; import org.thingsboard.server.common.data.device.data.PowerMode; import org.thingsboard.server.common.data.device.data.PowerSavingConfiguration; -import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; -import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.ProvisionDeviceProfileCredentials; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; @@ -308,8 +306,7 @@ public class DefaultTransportApiService implements TransportApiService { if (customerId != null && !customerId.isNullUid()) { metaData.putValue("customerId", customerId.toString()); } - String deviceIdStr = requestMsg.getSparkplug() ? "sparkplugId" : "gatewayId"; - metaData.putValue(deviceIdStr, gatewayId.toString()); + metaData.putValue("gatewayId", gatewayId.toString()); DeviceId deviceId = device.getId(); ObjectNode entityNode = mapper.valueToTree(device); @@ -320,12 +317,11 @@ public class DefaultTransportApiService implements TransportApiService { if (deviceAdditionalInfo == null) { deviceAdditionalInfo = JacksonUtil.newObjectNode(); } - String lastConnectedStr = requestMsg.getSparkplug() ? DataConstants.LAST_CONNECTED_SPARKPLUG : DataConstants.LAST_CONNECTED_GATEWAY; if (deviceAdditionalInfo.isObject() && - (!deviceAdditionalInfo.has(lastConnectedStr) - || !gatewayId.toString().equals(deviceAdditionalInfo.get(lastConnectedStr).asText()))) { + (!deviceAdditionalInfo.has(DataConstants.LAST_CONNECTED_GATEWAY) + || !gatewayId.toString().equals(deviceAdditionalInfo.get(DataConstants.LAST_CONNECTED_GATEWAY).asText()))) { ObjectNode newDeviceAdditionalInfo = (ObjectNode) deviceAdditionalInfo; - newDeviceAdditionalInfo.put(lastConnectedStr, gatewayId.toString()); + newDeviceAdditionalInfo.put(DataConstants.LAST_CONNECTED_GATEWAY, gatewayId.toString()); Device savedDevice = deviceService.saveDevice(device); tbClusterService.onDeviceUpdated(savedDevice, device); } 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 781c95d743..f4fe8b168e 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 @@ -103,8 +103,8 @@ public abstract class AbstractMqttIntegrationTest extends AbstractTransportInteg if (StringUtils.hasLength(config.getAttributesTopicFilter())) { mqttDeviceProfileTransportConfiguration.setDeviceAttributesTopic(config.getAttributesTopicFilter()); } - mqttDeviceProfileTransportConfiguration.setSparkPlug(config.isSparkPlug()); - mqttDeviceProfileTransportConfiguration.setSparkPlugAttributesMetricNames(config.sparkPlugAttributesMetricNames); + mqttDeviceProfileTransportConfiguration.setSparkplug(config.isSparkplug()); + mqttDeviceProfileTransportConfiguration.setSparkplugAttributesMetricNames(config.sparkplugAttributesMetricNames); mqttDeviceProfileTransportConfiguration.setSendAckOnValidationException(config.isSendAckOnValidationException()); TransportPayloadTypeConfiguration transportPayloadTypeConfiguration; if (TransportPayloadType.JSON.equals(transportPayloadType)) { diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java index f8d224bf41..5815adae76 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTestConfigProperties.java @@ -28,8 +28,8 @@ public class MqttTestConfigProperties { String deviceName; String gatewayName; - boolean isSparkPlug; - Set sparkPlugAttributesMetricNames; + boolean isSparkplug; + Set sparkplugAttributesMetricNames; TransportPayloadType transportPayloadType; 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 28b0830b39..99134be5d4 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 @@ -95,13 +95,13 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte //*BIRTH protected static final MetricDataType metricBirthDataType_Int32 = Int32; protected static final String metricBirthName_Int32 = "Device Metric int32"; - protected Set sparkPlugAttributesMetricNames; + protected Set sparkplugAttributesMetricNames; public void beforeSparkplugTest() throws Exception { MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() .gatewayName("Test Connect Sparkplug client node") - .isSparkPlug(true) - .sparkPlugAttributesMetricNames(sparkPlugAttributesMetricNames) + .isSparkplug(true) + .sparkplugAttributesMetricNames(sparkplugAttributesMetricNames) .transportPayloadType(TransportPayloadType.PROTOBUF) .build(); processBeforeTest(configProperties); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java index 6f0181cbf1..04349d0672 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java @@ -41,11 +41,6 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopi @Slf4j public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends AbstractMqttV5ClientSparkplugTest { - /** - * "sparkPlugAttributesMetricNames": ["SN node", "SN device", "Firmware version", "Date version", "Last date update"] - * @throws Exception - */ - protected void processClientWithCorrectAccessTokenPublishNCMDReBirth() throws Exception { clientWithCorrectNodeAccessTokenWithNDEATH(); List listKeys = connectionWithNBirth(metricBirthDataType_Int32, metricBirthName_Int32, nextInt32()); 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 4988a7850d..1d5a2c127f 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 @@ -31,8 +31,8 @@ public class MqttV5ClientSparkplugBAttributesInProfileTest extends AbstractMqttV @Before public void beforeTest() throws Exception { - sparkPlugAttributesMetricNames = new HashSet<>(); - sparkPlugAttributesMetricNames.add(metricBirthName_Int32); + sparkplugAttributesMetricNames = new HashSet<>(); + sparkplugAttributesMetricNames.add(metricBirthName_Int32); beforeSparkplugTest(); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index e6d9eeef64..fd35c0a67a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -120,7 +120,6 @@ public class DataConstants { public static final String MSG_SOURCE_KEY = "source"; public static final String LAST_CONNECTED_GATEWAY = "lastConnectedGateway"; - public static final String LAST_CONNECTED_SPARKPLUG = "lastConnectedSparkplug"; public static final String MAIN_QUEUE_NAME = "Main"; public static final String MAIN_QUEUE_TOPIC = "tb_rule_engine.main"; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java index 10634cecc9..d0b83e56d6 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java @@ -29,8 +29,8 @@ public class MqttDeviceProfileTransportConfiguration implements DeviceProfileTra @NoXss private String deviceAttributesTopic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC; private TransportPayloadTypeConfiguration transportPayloadTypeConfiguration; - private boolean sparkPlug; - private Set sparkPlugAttributesMetricNames; + private boolean sparkplug; + private Set sparkplugAttributesMetricNames; private boolean sendAckOnValidationException; @Override diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index f8804bc6fd..0d43fc6c93 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -392,21 +392,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement int msgId = mqttMsg.variableHeader().packetId(); try { SparkplugTopic sparkplugTopic = parseTopicPublish(topicName); - String deviceName = sparkplugTopic.isNode() ? deviceSessionCtx.getDeviceInfo().getDeviceName() : sparkplugTopic.getDeviceId(); if (sparkplugTopic.isNode()) { // A node topic SparkplugBProto.Payload sparkplugBProtoNode = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(mqttMsg.payload())); switch (sparkplugTopic.getType()) { - case STATE: - // TODO - break; case NBIRTH: case NCMD: case NDATA: - sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoNode, deviceName, sparkplugTopic); - break; - case NRECORD: - // TODO + sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoNode, deviceSessionCtx.getDeviceInfo().getDeviceName(), sparkplugTopic); break; default: } @@ -414,30 +407,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement // A device topic SparkplugBProto.Payload sparkplugBProtoDevice = SparkplugBProto.Payload.parseFrom(ProtoMqttAdaptor.toBytes(mqttMsg.payload())); switch (sparkplugTopic.getType()) { - case STATE: - // TODO - break; case DBIRTH: case DCMD: case DDATA: - sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoDevice, deviceName, sparkplugTopic); + sparkplugSessionHandler.onAttributesTelemetryProto(msgId, sparkplugBProtoDevice, sparkplugTopic.getDeviceId(), sparkplugTopic); break; - /** - * TODO - * 7.3.2. Device Death Certificate (DDEATH) - * The Sparkplug™ Topic Namespace for a device Death Certificate is: - * namespace/group_id/DDEATH/edge_node_id/device_id - * It is the responsibility of the MQTT EoN node to indicate the real-time state of either physical legacy device using - * poll/response protocols and/or local logical devices. If the device becomes unavailable for any reason (no - * response, CRC error, etc.) it is the responsibility of the EoN node to publish a DDEATH on behalf of the end device. - * Immediately upon reception of a DDEATH, any MQTT client subscribed to this device should set the data quality of - * all metrics to “STALE” and should note the time stamp when the DDEATH message was received. - */ case DDEATH: - sparkplugSessionHandler.onDeviceDisconnect(mqttMsg, deviceName); - break; - case DRECORD: - // TODO + sparkplugSessionHandler.onDeviceDisconnect(mqttMsg, sparkplugTopic.getDeviceId()); break; default: } @@ -1084,16 +1060,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - /** - * Sparkplug™ Specification Version 2.2 - * 7.1.1. EoN Node Death Certificate (NDEATH) - * The Death Certificate topic for an MQTT EoN node is: - * namespace/group_id/NDEATH/edge_node_id - * The Death Certificate topic and payload described here are not “published” as an MQTT message by a client, but - * provided as parameters within the MQTT CONNECT control packet when this MQTT EoN node first establishes the - * MQTT Client session. - */ - private void checkSparkplugNodeSession(MqttConnectMessage connectMessage, ChannelHandlerContext ctx) { try { if (sparkplugSessionHandler == null) { 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 28c25c8133..2bb59dc9e8 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 @@ -83,7 +83,7 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMess * Created by ashvayka on 19.01.17. */ @Slf4j -public abstract class AbstractGatewaySessionHandler { +public abstract class AbstractGatewaySessionHandler { protected static final String DEFAULT_DEVICE_TYPE = "default"; private static final String CAN_T_PARSE_VALUE = "Can't parse value: "; @@ -94,8 +94,8 @@ public abstract class AbstractGatewaySessionHandler { protected final TransportDeviceInfo gateway; protected final UUID sessionId; private final ConcurrentMap deviceCreationLockMap; - private final ConcurrentMap devices; - private final ConcurrentMap> deviceFutures; + private final ConcurrentMap devices; + private final ConcurrentMap> deviceFutures; protected final ConcurrentMap mqttQoSMap; protected final ChannelHandlerContext channel; protected final DeviceSessionCtx deviceSessionCtx; @@ -206,9 +206,9 @@ public abstract class AbstractGatewaySessionHandler { protected void processOnConnect(MqttPublishMessage msg, String deviceName, String deviceType) { log.trace("[{}] onDeviceConnect: {}", sessionId, deviceName); - Futures.addCallback(onDeviceConnect(deviceName, deviceType), new FutureCallback() { + Futures.addCallback(onDeviceConnect(deviceName, deviceType), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext result) { + public void onSuccess(@Nullable T result) { ack(msg, ReturnCode.SUCCESS); log.trace("[{}] onDeviceConnectOk: {}", sessionId, deviceName); } @@ -221,8 +221,8 @@ public abstract class AbstractGatewaySessionHandler { }, context.getExecutor()); } - ListenableFuture onDeviceConnect(String deviceName, String deviceType) { - MqttDeviceAwareSessionContext result = devices.get(deviceName); + ListenableFuture onDeviceConnect(String deviceName, String deviceType) { + T result = devices.get(deviceName); if (result == null) { Lock deviceCreationLock = deviceCreationLockMap.computeIfAbsent(deviceName, s -> new ReentrantLock()); deviceCreationLock.lock(); @@ -241,9 +241,9 @@ public abstract class AbstractGatewaySessionHandler { } } - private ListenableFuture getDeviceCreationFuture(String deviceName, String deviceType) { - final SettableFuture futureToSet = SettableFuture.create(); - ListenableFuture future = deviceFutures.putIfAbsent(deviceName, futureToSet); + private ListenableFuture getDeviceCreationFuture(String deviceName, String deviceType) { + final SettableFuture futureToSet = SettableFuture.create(); + ListenableFuture future = deviceFutures.putIfAbsent(deviceName, futureToSet); if (future != null) { return future; } @@ -258,7 +258,7 @@ public abstract class AbstractGatewaySessionHandler { new TransportServiceCallback<>() { @Override public void onSuccess(GetOrCreateDeviceFromGatewayResponse msg) { - AbstractGatewayDeviceSessionContext deviceSessionCtx = newDeviceSessionCtx(msg); + T deviceSessionCtx = newDeviceSessionCtx(msg); if (devices.putIfAbsent(deviceName, deviceSessionCtx) == null) { log.trace("[{}] First got or created device [{}], type [{}] for the gateway session", sessionId, deviceName, deviceType); SessionInfoProto deviceSessionInfo = deviceSessionCtx.getSessionInfo(); @@ -288,7 +288,7 @@ public abstract class AbstractGatewaySessionHandler { } } - protected abstract AbstractGatewayDeviceSessionContext newDeviceSessionCtx(GetOrCreateDeviceFromGatewayResponse msg); + protected abstract T newDeviceSessionCtx(GetOrCreateDeviceFromGatewayResponse msg); protected int getMsgId(MqttPublishMessage mqttMsg) { return mqttMsg.variableHeader().packetId(); @@ -341,7 +341,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { if (!deviceEntry.getValue().isJsonArray()) { throw new JsonSyntaxException(CAN_T_PARSE_VALUE + json); } @@ -375,7 +375,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { TransportProtos.PostTelemetryMsg msg = telemetryMsg.getMsg(); try { TransportProtos.PostTelemetryMsg postTelemetryMsg = ProtoConverter.validatePostTelemetryMsg(msg.toByteArray()); @@ -425,7 +425,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { if (!deviceEntry.getValue().isJsonObject()) { throw new JsonSyntaxException(CAN_T_PARSE_VALUE + json); } @@ -459,7 +459,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { TransportApiProtos.ClaimDevice claimRequest = claimDeviceMsg.getClaimRequest(); if (claimRequest == null) { throw new IllegalArgumentException("Claim request for device: " + deviceName + " is null!"); @@ -501,7 +501,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { if (!deviceEntry.getValue().isJsonObject()) { throw new JsonSyntaxException(CAN_T_PARSE_VALUE + json); } @@ -530,7 +530,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { TransportProtos.PostAttributeMsg kvListProto = attributesMsg.getMsg(); if (kvListProto == null) { throw new IllegalArgumentException("Attributes List for device: " + deviceName + " is empty!"); @@ -609,7 +609,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { Integer requestId = jsonObj.get("id").getAsInt(); String data = jsonObj.get("data").toString(); TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder() @@ -634,7 +634,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { Integer requestId = gatewayRpcResponseMsg.getId(); String data = gatewayRpcResponseMsg.getData(); TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder() @@ -661,7 +661,7 @@ public abstract class AbstractGatewaySessionHandler { Futures.addCallback(checkDeviceConnected(deviceName), new FutureCallback<>() { @Override - public void onSuccess(@Nullable MqttDeviceAwareSessionContext deviceCtx) { + public void onSuccess(@Nullable T deviceCtx) { transportService.process(deviceCtx.getSessionInfo(), requestMsg, getPubAckCallback(channel, deviceName, msgId, requestMsg)); } @@ -685,8 +685,8 @@ public abstract class AbstractGatewaySessionHandler { return result.build(); } - protected ListenableFuture checkDeviceConnected(String deviceName) { - MqttDeviceAwareSessionContext ctx = devices.get(deviceName); + protected ListenableFuture checkDeviceConnected(String deviceName) { + T ctx = devices.get(deviceName); if (ctx == null) { log.debug("[{}] Missing device [{}] for the gateway session", sessionId, deviceName); return onDeviceConnect(deviceName, DEFAULT_DEVICE_TYPE); @@ -729,7 +729,6 @@ public abstract class AbstractGatewaySessionHandler { private void deregisterSession(String deviceName, MqttDeviceAwareSessionContext deviceSessionCtx) { if (this.deviceSessionCtx.isSparkplug()) { - // add Msg Telemetry: key STATE type: String value: OFFLINE ts: sparkplugBProto.getTimestamp() sendSparkplugStateOnTelemetry(deviceSessionCtx.getSessionInfo(), deviceSessionCtx.getDeviceInfo().getDeviceName(), OFFLINE, new Date().getTime()); } @@ -744,7 +743,6 @@ public abstract class AbstractGatewaySessionHandler { keyValueProtoBuilder.setType(TransportProtos.KeyValueType.STRING_V); keyValueProtoBuilder.setStringV(connectionState.name()); TransportProtos.PostTelemetryMsg postTelemetryMsg = postTelemetryMsgCreated(keyValueProtoBuilder.build(), ts); - transportService.process(sessionInfo, postTelemetryMsg, getPubAckCallback(channel, deviceName, -1, postTelemetryMsg)); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java index 1fe5c7edc5..0498eadd9f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java @@ -32,30 +32,16 @@ import java.util.stream.Collectors; public abstract class MqttDeviceAwareSessionContext extends DeviceAwareSessionContext { private final ConcurrentMap mqttQoSMap; - private Map deviceBirthMetrics; public MqttDeviceAwareSessionContext(UUID sessionId, ConcurrentMap mqttQoSMap) { super(sessionId); this.mqttQoSMap = mqttQoSMap; - this.deviceBirthMetrics = null; } public ConcurrentMap getMqttQoSMap() { return mqttQoSMap; } - public Map getDeviceBirthMetrics() { - return deviceBirthMetrics; - } - - public void setDeviceBirthMetrics(java.util.List metrics) { - if (this.deviceBirthMetrics == null) { - this.deviceBirthMetrics = new ConcurrentHashMap<>(); - } - this.deviceBirthMetrics.putAll(metrics.stream() - .collect(Collectors.toMap(metric -> metric.getName(), metric -> metric))); - } - public MqttQoS getQoSForTopic(String topic) { List qosList = mqttQoSMap.entrySet() .stream() 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 b5bf28b36f..fb4ab488a8 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 @@ -23,19 +23,25 @@ import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugRpcRequestHeader; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic; import java.util.Date; +import java.util.Map; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.stream.Collectors; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.getTsKvProto; @Slf4j public class SparkplugDeviceSessionContext extends AbstractGatewayDeviceSessionContext { + private final Map deviceBirthMetrics = new ConcurrentHashMap<>(); + public SparkplugDeviceSessionContext(SparkplugNodeSessionHandler parent, TransportDeviceInfo deviceInfo, DeviceProfile deviceProfile, @@ -45,6 +51,16 @@ public class SparkplugDeviceSessionContext extends AbstractGatewayDeviceSessionC super(parent, deviceInfo, deviceProfile, mqttQoSMap, transportService); } + public Map getDeviceBirthMetrics() { + return deviceBirthMetrics; + } + + public void setDeviceBirthMetrics(java.util.List metrics) { + this.deviceBirthMetrics.putAll(metrics.stream() + .collect(Collectors.toMap(SparkplugBProto.Payload.Metric::getName, metric -> metric))); + } + + @Override public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg notification) { log.trace("[{}] Received attributes update notification to sparkplug device", sessionId); @@ -64,20 +80,7 @@ public class SparkplugDeviceSessionContext extends AbstractGatewayDeviceSessionC public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { log.trace("[{}] Received RPC Request notification to sparkplug device", sessionId); try { - /** - * DCMD {"metricName":"MyDeviceMetricText","value":"MyNodeMetric05_String_Value"} - * DCMD {"metricName":"MyNodeMetric02_LongInt64","value":2814119464032075444} - * DCMD {"metricName":"MyNodeMetric03_Double","value":6336935578763180333} - * DCMD {"metricName":"MyNodeMetric04_Float","value":413.18222} - * DCMD {"metricName":"Node Control/Rebirth","value":false} - * DCMD {"metricName":"MyNodeMetric06_Json_Bytes", "value":[40,47,-49]} - */ SparkplugMessageType messageType = SparkplugMessageType.parseMessageType(rpcRequest.getMethodName()); - if (messageType == null) { - parent.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), - ThingsboardErrorCode.INVALID_ARGUMENTS, "Unsupported SparkplugMessageType: " + rpcRequest.getMethodName() + rpcRequest.getParams()); - return; - } SparkplugRpcRequestHeader header = JacksonUtil.fromString(rpcRequest.getParams(), SparkplugRpcRequestHeader.class); header.setMessageType(messageType.name()); TransportProtos.TsKvProto tsKvProto = getTsKvProto(header.getMetricName(), header.getValue(), new Date().getTime()); 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 70533fea1b..fc480fa882 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 @@ -18,6 +18,7 @@ package org.thingsboard.server.transport.mqtt.session; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; import com.google.gson.JsonParser; import com.google.gson.JsonSyntaxException; import com.google.protobuf.Descriptors; @@ -65,7 +66,7 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopi * Created by nickAS21 on 12.12.22 */ @Slf4j -public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { +public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { private final SparkplugTopic sparkplugTopicNode; private final Map nodeBirthMetrics; @@ -81,7 +82,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { public void setNodeBirthMetrics(java.util.List metrics) { this.nodeBirthMetrics.putAll(metrics.stream() - .collect(Collectors.toMap(metric -> metric.getName(), metric -> metric))); + .collect(Collectors.toMap(SparkplugBProto.Payload.Metric::getName, metric -> metric))); } public Map getNodeBirthMetrics() { @@ -102,24 +103,28 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { public void onAttributesTelemetryProto(int msgId, SparkplugBProto.Payload sparkplugBProto, String deviceName, SparkplugTopic topic) throws AdaptorException, ThingsboardException { checkDeviceName(deviceName); - ListenableFuture contextListenableFuture = topic.isNode() ? - Futures.immediateFuture(this.deviceSessionCtx) : onDeviceConnectProto(deviceName); - try { - if (topic.isType(NBIRTH) || topic.isType(DBIRTH)) { - // add Msg Telemetry: key STATE type: String value: ONLINE ts: sparkplugBProto.getTimestamp() - sendSparkplugStateOnTelemetry(contextListenableFuture.get().getSessionInfo(), deviceName, ONLINE, - sparkplugBProto.getTimestamp()); - } + + ListenableFuture contextListenableFuture; + if (topic.isNode()) { if (topic.isType(NBIRTH)) { + sendSparkplugStateOnTelemetry(this.deviceSessionCtx.getSessionInfo(), deviceName, ONLINE, + sparkplugBProto.getTimestamp()); setNodeBirthMetrics(sparkplugBProto.getMetricsList()); - } else if (topic.isType(DBIRTH)) { - contextListenableFuture.get().setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); } - } catch (InterruptedException | ExecutionException e) { - log.error("Failed add Metrics or change SparkplugConnectionState. MessageType *BIRTH.", e); + contextListenableFuture = Futures.immediateFuture(this.deviceSessionCtx); + } else { + ListenableFuture deviceCtx = onDeviceConnectProto(deviceName); + contextListenableFuture = Futures.transform(deviceCtx, ctx -> { + if (topic.isType(DBIRTH)) { + sendSparkplugStateOnTelemetry(ctx.getSessionInfo(), deviceName, ONLINE, + sparkplugBProto.getTimestamp()); + ctx.setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); + } + return ctx; + }, MoreExecutors.directExecutor()); } Set attributesMetricNames = ((MqttDeviceProfileTransportConfiguration) deviceSessionCtx - .getDeviceProfile().getProfileData().getTransportConfiguration()).getSparkPlugAttributesMetricNames(); + .getDeviceProfile().getProfileData().getTransportConfiguration()).getSparkplugAttributesMetricNames(); if (attributesMetricNames != null) { List attributesMsgList = convertToPostAttributes(sparkplugBProto, attributesMetricNames, deviceName); onDeviceAttributesProto(contextListenableFuture, msgId, attributesMsgList, deviceName); @@ -213,7 +218,7 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { } } - private ListenableFuture onDeviceConnectProto(String deviceName) throws ThingsboardException { + private ListenableFuture onDeviceConnectProto(String deviceName) throws ThingsboardException { try { String deviceType = this.gateway.getDeviceType() + "-node"; return onDeviceConnect(deviceName, deviceType); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java index f2628b12b2..4870cfc3ca 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java @@ -87,7 +87,7 @@ public abstract class DeviceAwareSessionContext implements SessionContext { public boolean isSparkplug() { DeviceProfileTransportConfiguration transportConfiguration = this.deviceProfile.getProfileData().getTransportConfiguration(); if (transportConfiguration instanceof MqttDeviceProfileTransportConfiguration) { - return ((MqttDeviceProfileTransportConfiguration) transportConfiguration).isSparkPlug(); + return ((MqttDeviceProfileTransportConfiguration) transportConfiguration).isSparkplug(); } else { return false; } diff --git a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html index 3e48aaf011..1ef7892ecb 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.html @@ -16,17 +16,17 @@ -->
- + {{ 'device-profile.mqtt-device-topic-filters-spark-plug' | translate }} -
+ *ngIf="mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplug').value"> device-profile.mqtt-device-topic-filters-spark-plug-attribute-metric-names - + {{name}} close @@ -40,7 +40,7 @@ -
+
device-profile.mqtt-device-topic-filters
diff --git a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts index 832e72c95d..cf931a7fed 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts +++ b/ui-ngx/src/app/modules/home/components/profile/device/mqtt-device-profile-transport-configuration.component.ts @@ -96,8 +96,8 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control this.mqttDeviceProfileTransportConfigurationFormGroup = this.fb.group({ deviceAttributesTopic: [null, [Validators.required, this.validationMQTTTopic()]], deviceTelemetryTopic: [null, [Validators.required, this.validationMQTTTopic()]], - sparkPlug: [false], - sparkPlugAttributesMetricNames: [null], + sparkplug: [false], + sparkplugAttributesMetricNames: [null], sendAckOnValidationException: [false, Validators.required], transportPayloadTypeConfiguration: this.fb.group({ transportPayloadType: [TransportPayloadType.JSON, Validators.required], @@ -123,13 +123,13 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control .patchValue(false, {emitEvent: false}); } }); - this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlug').valueChanges.pipe( + this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplug').valueChanges.pipe( takeUntil(this.destroy$) ).subscribe((value) => { if (value) { this.mqttDeviceProfileTransportConfigurationFormGroup.disable({emitEvent: false}); - this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlug').enable({emitEvent: false}); - this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlugAttributesMetricNames').enable({emitEvent: false}); + this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplug').enable({emitEvent: false}); + this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplugAttributesMetricNames').enable({emitEvent: false}); } else { this.mqttDeviceProfileTransportConfigurationFormGroup.enable({emitEvent: false}); } @@ -152,7 +152,7 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control this.mqttDeviceProfileTransportConfigurationFormGroup.disable({emitEvent: false}); } else { this.mqttDeviceProfileTransportConfigurationFormGroup.enable({emitEvent: false}); - this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlug').updateValueAndValidity({onlySelf: true}); + this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplug').updateValueAndValidity({onlySelf: true}); } } @@ -170,17 +170,17 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control this.mqttDeviceProfileTransportConfigurationFormGroup.patchValue(value, {emitEvent: false}); this.updateTransportPayloadBasedControls(value.transportPayloadTypeConfiguration?.transportPayloadType); if (!this.disabled) { - this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlug').updateValueAndValidity({onlySelf: true}); + this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplug').updateValueAndValidity({onlySelf: true}); } } } removeAttributeMetricName(name: string): void { - const names: string[] = this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlugAttributesMetricNames').value; + const names: string[] = this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplugAttributesMetricNames').value; const index = names.indexOf(name); if (index >= 0) { names.splice(index, 1); - this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlugAttributesMetricNames').setValue(names); + this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplugAttributesMetricNames').setValue(names); } } @@ -189,13 +189,13 @@ export class MqttDeviceProfileTransportConfigurationComponent implements Control let value = event.value; if ((value || '').trim()) { value = value.trim(); - let names: string[] = this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlugAttributesMetricNames').value; + let names: string[] = this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplugAttributesMetricNames').value; if (!names || names.indexOf(value) === -1) { if (!names) { names = []; } names.push(value); - this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkPlugAttributesMetricNames').setValue(names, {emitEvent: true}); + this.mqttDeviceProfileTransportConfigurationFormGroup.get('sparkplugAttributesMetricNames').setValue(names, {emitEvent: true}); } } if (input) { diff --git a/ui-ngx/src/app/shared/models/device.models.ts b/ui-ngx/src/app/shared/models/device.models.ts index 9ac0d894bf..df7ba701c7 100644 --- a/ui-ngx/src/app/shared/models/device.models.ts +++ b/ui-ngx/src/app/shared/models/device.models.ts @@ -242,7 +242,7 @@ export interface DefaultDeviceProfileTransportConfiguration { export interface MqttDeviceProfileTransportConfiguration { deviceTelemetryTopic?: string; deviceAttributesTopic?: string; - sparkPlug?: boolean; + sparkplug?: boolean; sendAckOnValidationException?: boolean; transportPayloadTypeConfiguration?: { transportPayloadType?: TransportPayloadType; @@ -360,7 +360,7 @@ export function createDeviceProfileTransportConfiguration(type: DeviceTransportT const mqttTransportConfiguration: MqttDeviceProfileTransportConfiguration = { deviceTelemetryTopic: 'v1/devices/me/telemetry', deviceAttributesTopic: 'v1/devices/me/attributes', - sparkPlug: false, + sparkplug: false, sendAckOnValidationException: false, transportPayloadTypeConfiguration: { transportPayloadType: TransportPayloadType.JSON,