Browse Source

Custom MQTT Topic filters in device profile

pull/3477/head
Andrii Shvaika 6 years ago
parent
commit
ad273deace
  1. 2
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java
  2. 3
      common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentialsType.java
  3. 4
      common/message/src/main/java/org/thingsboard/server/common/msg/session/SessionContext.java
  4. 52
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  5. 46
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
  6. 18
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java
  7. 29
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/EqualsTopicFilter.java
  8. 22
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilter.java
  9. 57
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java
  10. 21
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/RegexTopicFilter.java
  11. 36
      common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java
  12. 11
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java

2
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() {

3
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
}

4
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);
}

52
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<MqttTopicMatcher, Integer> 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<Void>() {
deviceSessionCtx.setSessionInfo(SessionInfoCreator.create(msg, context, sessionId));
transportService.process(deviceSessionCtx.getSessionInfo(), DefaultTransportService.getSessionEventMsg(SessionEvent.OPEN), new TransportServiceCallback<Void>() {
@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);
}
}

46
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<MqttTopicMatcher, Integer> 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();
}
}
}

18
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<MqttTopicMatcher, Integer> mqttQoSMap) {
public GatewayDeviceSessionCtx(GatewaySessionHandler parent, TransportDeviceInfo deviceInfo,
DeviceProfile deviceProfile, ConcurrentMap<MqttTopicMatcher, Integer> 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();
}
}

29
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);
}
}

22
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);
}

57
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<String, MqttTopicFilter> 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;
}
}

21
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtil.java → 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();
}
}

36
common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicRegexUtilTest.java → 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));
}
}

11
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;
}

Loading…
Cancel
Save