Browse Source

added ability to create custom attrs subscribe topic

pull/6986/head
thingsboard 4 years ago
parent
commit
e09aef8d22
  1. 2
      application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java
  2. 5
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java
  3. 5
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  4. 14
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java

2
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()); mqttDeviceProfileTransportConfiguration.setDeviceTelemetryTopic(config.getTelemetryTopicFilter());
} }
if (StringUtils.hasLength(config.getAttributesTopicFilter())) { if (StringUtils.hasLength(config.getAttributesTopicFilter())) {
mqttDeviceProfileTransportConfiguration.setDeviceAttributesTopic(config.getAttributesTopicFilter()); mqttDeviceProfileTransportConfiguration.setDeviceAttributesPublishTopic(config.getAttributesTopicFilter());
} }
mqttDeviceProfileTransportConfiguration.setSendAckOnValidationException(config.isSendAckOnValidationException()); mqttDeviceProfileTransportConfiguration.setSendAckOnValidationException(config.isSendAckOnValidationException());
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration; TransportPayloadTypeConfiguration transportPayloadTypeConfiguration;

5
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java

@ -25,7 +25,10 @@ public class MqttDeviceProfileTransportConfiguration implements DeviceProfileTra
@NoXss @NoXss
private String deviceTelemetryTopic = MqttTopics.DEVICE_TELEMETRY_TOPIC; private String deviceTelemetryTopic = MqttTopics.DEVICE_TELEMETRY_TOPIC;
@NoXss @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 TransportPayloadTypeConfiguration transportPayloadTypeConfiguration;
private boolean sendAckOnValidationException; private boolean sendAckOnValidationException;

5
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()) { for (MqttTopicSubscription subscription : mqttMsg.payload().topicSubscriptions()) {
String topic = subscription.topicName(); String topic = subscription.topicName();
MqttQoS reqQoS = subscription.qualityOfService(); MqttQoS reqQoS = subscription.qualityOfService();
if (deviceSessionCtx.isDeviceSubscriptionAttributesTopic(topic)){
processAttributesSubscribe(grantedQoSList, topic, reqQoS, TopicType.V1);
activityReported = true;
continue;
}
try { try {
switch (topic) { switch (topic) {
case MqttTopics.DEVICE_ATTRIBUTES_TOPIC: { case MqttTopics.DEVICE_ATTRIBUTES_TOPIC: {

14
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 boolean provisionOnly = false;
private volatile MqttTopicFilter telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); 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 TransportPayloadType payloadType = TransportPayloadType.JSON;
private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor; private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor;
private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor; private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor;
@ -105,7 +106,11 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
} }
public boolean isDeviceAttributesTopic(String topicName) { 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() { public MqttTransportAdaptor getPayloadAdaptor() {
@ -156,7 +161,8 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration(); TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration();
payloadType = transportPayloadTypeConfiguration.getTransportPayloadType(); payloadType = transportPayloadTypeConfiguration.getTransportPayloadType();
telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic());
attributesTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); attributesPublishTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesPublishTopic());
attributesSubscribeTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesSubscribeTopic());
sendAckOnValidationException = mqttConfig.isSendAckOnValidationException(); sendAckOnValidationException = mqttConfig.isSendAckOnValidationException();
if (TransportPayloadType.PROTOBUF.equals(payloadType)) { if (TransportPayloadType.PROTOBUF.equals(payloadType)) {
ProtoTransportPayloadConfiguration protoTransportPayloadConfig = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; ProtoTransportPayloadConfiguration protoTransportPayloadConfig = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration;
@ -166,7 +172,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
} }
} else { } else {
telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter();
attributesTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter();
payloadType = TransportPayloadType.JSON; payloadType = TransportPayloadType.JSON;
sendAckOnValidationException = false; sendAckOnValidationException = false;
} }

Loading…
Cancel
Save