From e09aef8d224781d5ab1267c1d2909a76c1721724 Mon Sep 17 00:00:00 2001 From: thingsboard Date: Wed, 20 Jul 2022 17:43:46 +0300 Subject: [PATCH] added ability to create custom attrs subscribe topic --- .../mqtt/AbstractMqttIntegrationTest.java | 2 +- .../MqttDeviceProfileTransportConfiguration.java | 5 ++++- .../transport/mqtt/MqttTransportHandler.java | 5 +++++ .../transport/mqtt/session/DeviceSessionCtx.java | 14 ++++++++++---- 4 files changed, 20 insertions(+), 6 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 9fbdff3fa5..ec64ce4af3 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 @@ -106,7 +106,7 @@ public abstract class AbstractMqttIntegrationTest extends AbstractTransportInteg mqttDeviceProfileTransportConfiguration.setDeviceTelemetryTopic(config.getTelemetryTopicFilter()); } if (StringUtils.hasLength(config.getAttributesTopicFilter())) { - mqttDeviceProfileTransportConfiguration.setDeviceAttributesTopic(config.getAttributesTopicFilter()); + mqttDeviceProfileTransportConfiguration.setDeviceAttributesPublishTopic(config.getAttributesTopicFilter()); } mqttDeviceProfileTransportConfiguration.setSendAckOnValidationException(config.isSendAckOnValidationException()); TransportPayloadTypeConfiguration transportPayloadTypeConfiguration; 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 1a4d2c72ad..8c1493bce7 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 @@ -25,7 +25,10 @@ public class MqttDeviceProfileTransportConfiguration implements DeviceProfileTra @NoXss private String deviceTelemetryTopic = MqttTopics.DEVICE_TELEMETRY_TOPIC; @NoXss - private String deviceAttributesTopic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC; + private String deviceAttributesPublishTopic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC; + @NoXss + private String deviceAttributesSubscribeTopic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC;//todo + private TransportPayloadTypeConfiguration transportPayloadTypeConfiguration; private boolean sendAckOnValidationException; 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 f37715efc8..288f32ec4e 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 @@ -622,6 +622,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement for (MqttTopicSubscription subscription : mqttMsg.payload().topicSubscriptions()) { String topic = subscription.topicName(); MqttQoS reqQoS = subscription.qualityOfService(); + if (deviceSessionCtx.isDeviceSubscriptionAttributesTopic(topic)){ + processAttributesSubscribe(grantedQoSList, topic, reqQoS, TopicType.V1); + activityReported = true; + continue; + } try { switch (topic) { case MqttTopics.DEVICE_ATTRIBUTES_TOPIC: { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index 2163610b38..4b69f75f06 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -75,7 +75,8 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { private boolean provisionOnly = false; private volatile MqttTopicFilter telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); - private volatile MqttTopicFilter attributesTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); + private volatile MqttTopicFilter attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); + private volatile MqttTopicFilter attributesSubscribeTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); private volatile TransportPayloadType payloadType = TransportPayloadType.JSON; private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor; private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor; @@ -105,7 +106,11 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { } public boolean isDeviceAttributesTopic(String topicName) { - return attributesTopicFilter.filter(topicName); + return attributesPublishTopicFilter.filter(topicName); + } + + public boolean isDeviceSubscriptionAttributesTopic(String topicName) { + return attributesSubscribeTopicFilter.filter(topicName); } public MqttTransportAdaptor getPayloadAdaptor() { @@ -156,7 +161,8 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration(); payloadType = transportPayloadTypeConfiguration.getTransportPayloadType(); telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); - attributesTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); + attributesPublishTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesPublishTopic()); + attributesSubscribeTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesSubscribeTopic()); sendAckOnValidationException = mqttConfig.isSendAckOnValidationException(); if (TransportPayloadType.PROTOBUF.equals(payloadType)) { ProtoTransportPayloadConfiguration protoTransportPayloadConfig = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; @@ -166,7 +172,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { } } else { telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); - attributesTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); + attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); payloadType = TransportPayloadType.JSON; sendAckOnValidationException = false; }