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 a3533c202d..5f0f7dceac 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 @@ -23,8 +23,6 @@ public class MqttDeviceProfileTransportConfiguration implements DeviceProfileTra private String deviceTelemetryTopic = MqttTopics.DEVICE_TELEMETRY_TOPIC; private String deviceAttributesTopic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC; - private String deviceRpcRequestTopic = MqttTopics.DEVICE_RPC_REQUESTS_TOPIC; - private String deviceRpcResponseTopic = MqttTopics.DEVICE_RPC_RESPONSE_TOPIC; @Override public DeviceTransportType getType() { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentialsType.java b/common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentialsType.java index 8d647dcb1a..e7faf1b5ce 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentialsType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentialsType.java @@ -18,6 +18,7 @@ package org.thingsboard.server.common.data.security; public enum DeviceCredentialsType { ACCESS_TOKEN, - X509_CERTIFICATE + X509_CERTIFICATE, + MQTT_BASIC } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/session/SessionContext.java b/common/message/src/main/java/org/thingsboard/server/common/msg/session/SessionContext.java index b27b957500..d124fd48b0 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/session/SessionContext.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/session/SessionContext.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.common.msg.session; +import org.thingsboard.server.common.data.DeviceProfile; + import java.util.UUID; public interface SessionContext { @@ -22,4 +24,6 @@ public interface SessionContext { UUID getSessionId(); int nextMsgId(); + + void onProfileUpdate(DeviceProfile deviceProfile); } 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 b9be909b40..57ed954a9c 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 @@ -99,9 +99,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private final SslHandler sslHandler; private final ConcurrentMap mqttQoSMap; - private volatile SessionInfoProto sessionInfo; + private final DeviceSessionCtx deviceSessionCtx; private volatile InetSocketAddress address; - private volatile DeviceSessionCtx deviceSessionCtx; private volatile GatewaySessionHandler gatewaySessionHandler; MqttTransportHandler(MqttTransportContext context, SslHandler sslHandler) { @@ -152,7 +151,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement case PINGREQ: if (checkConnected(ctx, msg)) { ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); - transportService.reportActivity(sessionInfo); + transportService.reportActivity(deviceSessionCtx.getSessionInfo()); } break; case DISCONNECT: @@ -176,7 +175,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) { if (gatewaySessionHandler != null) { handleGatewayPublishMsg(topicName, msgId, mqttMsg); - transportService.reportActivity(sessionInfo); + transportService.reportActivity(deviceSessionCtx.getSessionInfo()); } } else { processDevicePublish(ctx, mqttMsg, topicName, msgId); @@ -215,26 +214,26 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) { try { - if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_TOPIC)) { + if (deviceSessionCtx.isDeviceTelemetryTopic(topicName)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = adaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(sessionInfo, postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg)); - } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_TOPIC)) { + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + } else if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) { TransportProtos.PostAttributeMsg postAttributeMsg = adaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(sessionInfo, postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { TransportProtos.GetAttributeRequestMsg getAttributeMsg = adaptor.convertToGetAttributes(deviceSessionCtx, mqttMsg); - transportService.process(sessionInfo, getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_RPC_RESPONSE_TOPIC)) { TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = adaptor.convertToDeviceRpcResponse(deviceSessionCtx, mqttMsg); - transportService.process(sessionInfo, rpcResponseMsg, getPubAckCallback(ctx, msgId, rpcResponseMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), rpcResponseMsg, getPubAckCallback(ctx, msgId, rpcResponseMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_RPC_REQUESTS_TOPIC)) { TransportProtos.ToServerRpcRequestMsg rpcRequestMsg = adaptor.convertToServerRpcRequest(deviceSessionCtx, mqttMsg); - transportService.process(sessionInfo, rpcRequestMsg, getPubAckCallback(ctx, msgId, rpcRequestMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequestMsg, getPubAckCallback(ctx, msgId, rpcRequestMsg)); } else if (topicName.equals(MqttTopics.DEVICE_CLAIM_TOPIC)) { TransportProtos.ClaimDeviceMsg claimDeviceMsg = adaptor.convertToClaimDevice(deviceSessionCtx, mqttMsg); - transportService.process(sessionInfo, claimDeviceMsg, getPubAckCallback(ctx, msgId, claimDeviceMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), claimDeviceMsg, getPubAckCallback(ctx, msgId, claimDeviceMsg)); } else { - transportService.reportActivity(sessionInfo); + transportService.reportActivity(deviceSessionCtx.getSessionInfo()); } } catch (AdaptorException e) { log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); @@ -274,13 +273,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { switch (topic) { case MqttTopics.DEVICE_ATTRIBUTES_TOPIC: { - transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null); + transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null); registerSubQoS(topic, grantedQoSList, reqQoS); activityReported = true; break; } case MqttTopics.DEVICE_RPC_REQUESTS_SUB_TOPIC: { - transportService.process(sessionInfo, TransportProtos.SubscribeToRPCMsg.newBuilder().build(), null); + transportService.process(deviceSessionCtx.getSessionInfo(), TransportProtos.SubscribeToRPCMsg.newBuilder().build(), null); registerSubQoS(topic, grantedQoSList, reqQoS); activityReported = true; break; @@ -303,7 +302,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } if (!activityReported) { - transportService.reportActivity(sessionInfo); + transportService.reportActivity(deviceSessionCtx.getSessionInfo()); } ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList)); } @@ -324,12 +323,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { switch (topicName) { case MqttTopics.DEVICE_ATTRIBUTES_TOPIC: { - transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().setUnsubscribe(true).build(), null); + transportService.process(deviceSessionCtx.getSessionInfo(), + TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().setUnsubscribe(true).build(), null); activityReported = true; break; } case MqttTopics.DEVICE_RPC_REQUESTS_SUB_TOPIC: { - transportService.process(sessionInfo, TransportProtos.SubscribeToRPCMsg.newBuilder().setUnsubscribe(true).build(), null); + transportService.process(deviceSessionCtx.getSessionInfo(), + TransportProtos.SubscribeToRPCMsg.newBuilder().setUnsubscribe(true).build(), null); activityReported = true; break; } @@ -339,7 +340,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } if (!activityReported) { - transportService.reportActivity(sessionInfo); + transportService.reportActivity(deviceSessionCtx.getSessionInfo()); } ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId())); } @@ -499,8 +500,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void doDisconnect() { if (deviceSessionCtx.isConnected()) { - transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.CLOSED), null); - transportService.deregisterSession(sessionInfo); + transportService.process(deviceSessionCtx.getSessionInfo(), DefaultTransportService.getSessionEventMsg(SessionEvent.CLOSED), null); + transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); if (gatewaySessionHandler != null) { gatewaySessionHandler.onGatewayDisconnect(); } @@ -515,11 +516,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } else { deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo()); deviceSessionCtx.setDeviceProfile(msg.getDeviceProfile()); - sessionInfo = SessionInfoCreator.create(msg, context, sessionId); - transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.OPEN), new TransportServiceCallback() { + deviceSessionCtx.setSessionInfo(SessionInfoCreator.create(msg, context, sessionId)); + transportService.process(deviceSessionCtx.getSessionInfo(), DefaultTransportService.getSessionEventMsg(SessionEvent.OPEN), new TransportServiceCallback() { @Override public void onSuccess(Void msg) { - transportService.registerAsyncSession(sessionInfo, MqttTransportHandler.this); + transportService.registerAsyncSession(deviceSessionCtx.getSessionInfo(), MqttTransportHandler.this); checkGatewaySession(); ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); log.info("[{}] Client connected!", sessionId); @@ -581,7 +582,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onProfileUpdate(DeviceProfile deviceProfile) { - deviceSessionCtx.getDeviceInfo().setDeviceType(deviceProfile.getName()); - sessionInfo = SessionInfoProto.newBuilder().mergeFrom(sessionInfo).setDeviceType(deviceProfile.getName()).build(); + deviceSessionCtx.onProfileUpdate(deviceProfile); } } 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 d8732802db..3f5e0dc7ad 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 @@ -18,6 +18,13 @@ package org.thingsboard.server.transport.mqtt.session; import io.netty.channel.ChannelHandlerContext; import lombok.Getter; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.DeviceTransportType; +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.MqttTopics; +import org.thingsboard.server.transport.mqtt.util.MqttTopicFilter; +import org.thingsboard.server.transport.mqtt.util.MqttTopicFilterFactory; import java.util.UUID; import java.util.concurrent.ConcurrentMap; @@ -31,7 +38,11 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { @Getter private ChannelHandlerContext channel; - private AtomicInteger msgIdSeq = new AtomicInteger(0); + private final AtomicInteger msgIdSeq = new AtomicInteger(0); + + private volatile MqttTopicFilter telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); + private volatile MqttTopicFilter attributesTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); + public DeviceSessionCtx(UUID sessionId, ConcurrentMap mqttQoSMap) { super(sessionId, mqttQoSMap); @@ -44,4 +55,37 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { public int nextMsgId() { return msgIdSeq.incrementAndGet(); } + + public boolean isDeviceTelemetryTopic(String topicName) { + return telemetryTopicFilter.filter(topicName); + } + + public boolean isDeviceAttributesTopic(String topicName) { + return attributesTopicFilter.filter(topicName); + } + + @Override + public void setDeviceProfile(DeviceProfile deviceProfile) { + super.setDeviceProfile(deviceProfile); + updateTopicFilters(deviceProfile); + } + + @Override + public void onProfileUpdate(DeviceProfile deviceProfile) { + super.onProfileUpdate(deviceProfile); + updateTopicFilters(deviceProfile); + } + + private void updateTopicFilters(DeviceProfile deviceProfile) { + DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); + if (transportConfiguration.getType().equals(DeviceTransportType.MQTT) && + transportConfiguration instanceof MqttDeviceProfileTransportConfiguration) { + MqttDeviceProfileTransportConfiguration mqttConfig = (MqttDeviceProfileTransportConfiguration) transportConfiguration; + telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); + attributesTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); + } else { + telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); + attributesTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); + } + } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java index c58d6e1744..dc2217f346 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java @@ -33,12 +33,12 @@ import java.util.concurrent.ConcurrentMap; public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext implements SessionMsgListener { private final GatewaySessionHandler parent; - private volatile SessionInfoProto sessionInfo; - public GatewayDeviceSessionCtx(GatewaySessionHandler parent, TransportDeviceInfo deviceInfo, DeviceProfile deviceProfile, ConcurrentMap mqttQoSMap) { + public GatewayDeviceSessionCtx(GatewaySessionHandler parent, TransportDeviceInfo deviceInfo, + DeviceProfile deviceProfile, ConcurrentMap mqttQoSMap) { super(UUID.randomUUID(), mqttQoSMap); this.parent = parent; - this.sessionInfo = SessionInfoProto.newBuilder() + setSessionInfo(SessionInfoProto.newBuilder() .setNodeId(parent.getNodeId()) .setSessionIdMSB(sessionId.getMostSignificantBits()) .setSessionIdLSB(sessionId.getLeastSignificantBits()) @@ -52,7 +52,7 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple .setGwSessionIdLSB(parent.getSessionId().getLeastSignificantBits()) .setDeviceProfileIdMSB(deviceInfo.getDeviceProfileId().getId().getMostSignificantBits()) .setDeviceProfileIdLSB(deviceInfo.getDeviceProfileId().getId().getLeastSignificantBits()) - .build(); + .build()); setDeviceInfo(deviceInfo); setDeviceProfile(deviceProfile); } @@ -67,10 +67,6 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple return parent.nextMsgId(); } - SessionInfoProto getSessionInfo() { - return sessionInfo; - } - @Override public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg response) { try { @@ -107,10 +103,4 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) { // This feature is not supported in the TB IoT Gateway yet. } - - @Override - public void onProfileUpdate(DeviceProfile deviceProfile) { - deviceInfo.setDeviceType(deviceProfile.getName()); - sessionInfo = SessionInfoProto.newBuilder().mergeFrom(sessionInfo).setDeviceType(deviceProfile.getName()).build(); - } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/EqualsTopicFilter.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/EqualsTopicFilter.java new file mode 100644 index 0000000000..539c8f2b7e --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/EqualsTopicFilter.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2020 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.util; + +import lombok.Data; + +@Data +public class EqualsTopicFilter implements MqttTopicFilter { + + private final String filter; + + @Override + public boolean filter(String topic) { + return filter.equals(topic); + } +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilter.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilter.java new file mode 100644 index 0000000000..005deb5d44 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilter.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2020 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.util; + +public interface MqttTopicFilter { + + boolean filter(String topic); + +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java new file mode 100644 index 0000000000..51545a42f0 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java @@ -0,0 +1,57 @@ +/** + * Copyright © 2016-2020 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.util; + +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.device.profile.MqttTopics; + +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.regex.Pattern; + +@Slf4j +public class MqttTopicFilterFactory { + + private static final ConcurrentMap filters = new ConcurrentHashMap<>(); + private static final MqttTopicFilter DEFAULT_TELEMETRY_TOPIC_FILTER = toFilter(MqttTopics.DEVICE_TELEMETRY_TOPIC); + private static final MqttTopicFilter DEFAULT_ATTRIBUTES_TOPIC_FILTER = toFilter(MqttTopics.DEVICE_ATTRIBUTES_TOPIC); + + public static MqttTopicFilter toFilter(String topicFilter) { + if (topicFilter == null || topicFilter.isEmpty()) { + throw new IllegalArgumentException("Topic filter can't be empty!"); + } + return filters.computeIfAbsent(topicFilter, filter -> { + if (filter.contains("+") || filter.contains("#")) { + String regex = filter + .replace("\\", "\\\\") + .replace("+", "[^/]+") + .replace("/#", "($|/.*)"); + log.debug("Converting [{}] to [{}]", filter, regex); + return new RegexTopicFilter(regex); + } else { + return new EqualsTopicFilter(filter); + } + }); + } + + public static MqttTopicFilter getDefaultTelemetryFilter() { + return DEFAULT_TELEMETRY_TOPIC_FILTER; + } + + public static MqttTopicFilter getDefaultAttributesFilter() { + return DEFAULT_ATTRIBUTES_TOPIC_FILTER; + } +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtil.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/RegexTopicFilter.java similarity index 64% rename from common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtil.java rename to common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/RegexTopicFilter.java index 1c3c2ff680..d5f50ae8c9 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtil.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/RegexTopicFilter.java @@ -15,20 +15,21 @@ */ package org.thingsboard.server.transport.mqtt.util; -import lombok.extern.slf4j.Slf4j; +import lombok.Data; import java.util.regex.Pattern; -@Slf4j -public class MqttTopicRegexUtil { +@Data +public class RegexTopicFilter implements MqttTopicFilter { - public static Pattern toRegex(String topicFilter) { - String regex = topicFilter - .replace("\\", "\\\\") - .replace("+", "[^/]+") - .replace("/#", "($|/.*)"); - log.debug("Converting [{}] to [{}]", topicFilter, regex); - return Pattern.compile(regex); + private final Pattern regex; + + public RegexTopicFilter(String regex) { + this.regex = Pattern.compile(regex); } + @Override + public boolean filter(String topic) { + return regex.matcher(topic).matches(); + } } diff --git a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtilTest.java b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java similarity index 58% rename from common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtilTest.java rename to common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java index 2cbfd22d5a..0b854d51ef 100644 --- a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtilTest.java +++ b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java @@ -26,7 +26,7 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @RunWith(MockitoJUnitRunner.class) -public class MqttTopicRegexUtilTest { +public class MqttTopicFilterFactoryTest { private static String TEST_STR_1 = "Sensor/Temperature/House/48"; private static String TEST_STR_2 = "Sensor/Temperature"; @@ -34,23 +34,23 @@ public class MqttTopicRegexUtilTest { @Test public void metadataCanBeUpdated() throws ScriptException { - Pattern filter = MqttTopicRegexUtil.toRegex("Sensor/Temperature/House/+"); - assertTrue(filter.matcher(TEST_STR_1).matches()); - assertFalse(filter.matcher(TEST_STR_2).matches()); - - filter = MqttTopicRegexUtil.toRegex("Sensor/+/House/#"); - assertTrue(filter.matcher(TEST_STR_1).matches()); - assertFalse(filter.matcher(TEST_STR_2).matches()); - - filter = MqttTopicRegexUtil.toRegex("Sensor/#"); - assertTrue(filter.matcher(TEST_STR_1).matches()); - assertTrue(filter.matcher(TEST_STR_2).matches()); - assertTrue(filter.matcher(TEST_STR_3).matches()); - - filter = MqttTopicRegexUtil.toRegex("Sensor/Temperature/#"); - assertTrue(filter.matcher(TEST_STR_1).matches()); - assertTrue(filter.matcher(TEST_STR_2).matches()); - assertFalse(filter.matcher(TEST_STR_3).matches()); + MqttTopicFilter filter = MqttTopicFilterFactory.toFilter("Sensor/Temperature/House/+"); + assertTrue(filter.filter(TEST_STR_1)); + assertFalse(filter.filter(TEST_STR_2)); + + filter = MqttTopicFilterFactory.toFilter("Sensor/+/House/#"); + assertTrue(filter.filter(TEST_STR_1)); + assertFalse(filter.filter(TEST_STR_2)); + + filter = MqttTopicFilterFactory.toFilter("Sensor/#"); + assertTrue(filter.filter(TEST_STR_1)); + assertTrue(filter.filter(TEST_STR_2)); + assertTrue(filter.filter(TEST_STR_3)); + + filter = MqttTopicFilterFactory.toFilter("Sensor/Temperature/#"); + assertTrue(filter.filter(TEST_STR_1)); + assertTrue(filter.filter(TEST_STR_2)); + assertFalse(filter.filter(TEST_STR_3)); } } 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 a454a4e06d..2f7f2d69e6 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 @@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.msg.session.SessionContext; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto; import java.util.UUID; @@ -41,6 +42,9 @@ public abstract class DeviceAwareSessionContext implements SessionContext { @Getter @Setter protected volatile DeviceProfile deviceProfile; + @Getter + @Setter + private volatile TransportProtos.SessionInfoProto sessionInfo; private volatile boolean connected; @@ -54,6 +58,13 @@ public abstract class DeviceAwareSessionContext implements SessionContext { this.deviceId = deviceInfo.getDeviceId(); } + @Override + public void onProfileUpdate(DeviceProfile deviceProfile) { + this.deviceProfile = deviceProfile; + this.deviceInfo.setDeviceType(deviceProfile.getName()); + this.sessionInfo = TransportProtos.SessionInfoProto.newBuilder().mergeFrom(sessionInfo).setDeviceType(deviceProfile.getName()).build(); + } + public boolean isConnected() { return connected; }