From b650bf67a499fa7ab38067ee9d432de1648298fe Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Thu, 19 Jan 2017 17:30:11 +0200 Subject: [PATCH 1/3] TB-33: Implementation --- .../session/DeviceAwareSessionContext.java | 14 ++ .../server/dao/device/DeviceService.java | 4 + .../server/dao/device/DeviceServiceImpl.java | 13 ++ .../server/transport/mqtt/MqttTopics.java | 43 ++++ .../transport/mqtt/MqttTransportHandler.java | 205 ++++++++++++------ .../mqtt/MqttTransportServerInitializer.java | 7 +- .../transport/mqtt/MqttTransportService.java | 11 +- .../mqtt/adaptors/JsonMqttAdaptor.java | 65 ++++-- .../mqtt/adaptors/JsonMqttGatewayAdaptor.java | 49 +++++ .../mqtt/adaptors/MqttGatewayAdaptor.java | 36 +++ .../mqtt/adaptors/MqttTransportAdaptor.java | 4 +- ...tSessionCtx.java => DeviceSessionCtx.java} | 4 +- .../mqtt/session/GatewayDeviceSessionCtx.java | 79 +++++++ .../mqtt/session/GatewaySessionCtx.java | 69 ++++++ 14 files changed, 499 insertions(+), 104 deletions(-) create mode 100644 transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTopics.java create mode 100644 transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java create mode 100644 transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttGatewayAdaptor.java rename transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/{MqttSessionCtx.java => DeviceSessionCtx.java} (95%) create mode 100644 transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java create mode 100644 transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java index 8321c0468f..89debed458 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java @@ -42,6 +42,12 @@ public abstract class DeviceAwareSessionContext implements SessionContext { this.authService = authService; } + public DeviceAwareSessionContext(SessionMsgProcessor processor, DeviceAuthService authService, Device device) { + this(processor, authService); + this.device = device; + } + + public boolean login(DeviceCredentialsFilter credentials) { DeviceAuthResult result = authService.process(credentials); if (result.isSuccess()) { @@ -56,6 +62,14 @@ public abstract class DeviceAwareSessionContext implements SessionContext { } } + public DeviceAuthService getAuthService() { + return authService; + } + + public SessionMsgProcessor getProcessor() { + return processor; + } + public Device getDevice() { return device; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java index ad4c338f13..8d780b6aa1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java @@ -22,10 +22,14 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; +import java.util.Optional; + public interface DeviceService { Device findDeviceById(DeviceId deviceId); + Optional findDeviceByTenantIdAndName(TenantId tenantId, String name); + Device saveDevice(Device device); Device assignDeviceToCustomer(DeviceId deviceId, CustomerId customerId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java index 3a3c018d84..681188e6be 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java @@ -38,6 +38,7 @@ import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.tenant.TenantDao; import java.util.List; +import java.util.Optional; import static org.thingsboard.server.dao.DaoUtil.convertDataList; import static org.thingsboard.server.dao.DaoUtil.getData; @@ -69,6 +70,18 @@ public class DeviceServiceImpl implements DeviceService { return getData(deviceEntity); } + @Override + public Optional findDeviceByTenantIdAndName(TenantId tenantId, String name) { + log.trace("Executing findDeviceByTenantIdAndName [{}][{}]", tenantId, name); + validateId(tenantId, "Incorrect tenantId " + tenantId); + Optional deviceEntityOpt = deviceDao.findDevicesByTenantIdAndName(tenantId.getId(), name); + if (deviceEntityOpt.isPresent()) { + return Optional.of(getData(deviceEntityOpt.get())); + } else { + return Optional.empty(); + } + } + @Override public Device saveDevice(Device device) { log.trace("Executing saveDevice [{}]", device); diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTopics.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTopics.java new file mode 100644 index 0000000000..4f91e1add4 --- /dev/null +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTopics.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2017 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; + +/** + * Created by ashvayka on 19.01.17. + */ +public class MqttTopics { + + public static final String BASE_DEVICE_API_TOPIC = "v1/devices/me"; + public static final String DEVICE_RPC_RESPONSE_TOPIC = BASE_DEVICE_API_TOPIC + "/rpc/response/"; + public static final String DEVICE_RPC_RESPONSE_SUB_TOPIC = DEVICE_RPC_RESPONSE_TOPIC + "+"; + public static final String DEVICE_RPC_REQUESTS_TOPIC = BASE_DEVICE_API_TOPIC + "/rpc/request/"; + public static final String DEVICE_RPC_REQUESTS_SUB_TOPIC = DEVICE_RPC_REQUESTS_TOPIC + "+"; + public static final String DEVICE_ATTRIBUTES_RESPONSE_TOPIC_PREFIX = BASE_DEVICE_API_TOPIC + "/attributes/response/"; + public static final String DEVICE_ATTRIBUTES_RESPONSES_TOPIC = DEVICE_ATTRIBUTES_RESPONSE_TOPIC_PREFIX + "+"; + public static final String DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX = BASE_DEVICE_API_TOPIC + "/attributes/request/"; + public static final String DEVICE_TELEMETRY_TOPIC = BASE_DEVICE_API_TOPIC + "/telemetry"; + public static final String DEVICE_ATTRIBUTES_TOPIC = BASE_DEVICE_API_TOPIC + "/attributes"; + + public static final String BASE_GATEWAY_API_TOPIC = "v1/gateway"; + public static final String GATEWAY_CONNECT_TOPIC = "v1/gateway/connect"; + public static final String GATEWAY_DISCONNECT_TOPIC = "v1/gateway/disconnect"; + public static final String GATEWAY_ATTRIBUTES_TOPIC = "v1/gateway/attributes"; + public static final String GATEWAY_TELEMETRY_TOPIC = "v1/gateway/telemetry"; + + + private MqttTopics() { + } +} diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 30b37e3ff6..e74495572f 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -15,26 +15,33 @@ */ package org.thingsboard.server.transport.mqtt; +import com.fasterxml.jackson.databind.JsonNode; +import com.google.gson.JsonElement; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; import io.netty.handler.codec.mqtt.*; +import io.netty.handler.codec.mqtt.MqttConnectReturnCode; +import io.netty.handler.codec.mqtt.MqttQoS; import io.netty.handler.ssl.SslHandler; import io.netty.util.concurrent.Future; import io.netty.util.concurrent.GenericFutureListener; import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; +import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.security.DeviceTokenCredentials; import org.thingsboard.server.common.data.security.DeviceX509Credentials; import org.thingsboard.server.common.msg.session.AdaptorToSessionActorMsg; import org.thingsboard.server.common.msg.session.BasicToDeviceActorSessionMsg; -import org.thingsboard.server.common.msg.session.MsgType; import org.thingsboard.server.common.msg.session.ctrl.SessionCloseMsg; import org.thingsboard.server.common.transport.SessionMsgProcessor; import org.thingsboard.server.common.transport.adaptor.AdaptorException; import org.thingsboard.server.common.transport.auth.DeviceAuthService; import org.thingsboard.server.dao.EncryptionUtil; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; -import org.thingsboard.server.transport.mqtt.session.MqttSessionCtx; +import org.thingsboard.server.transport.mqtt.session.GatewaySessionCtx; +import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; import org.thingsboard.server.transport.mqtt.util.SslUtil; import javax.net.ssl.SSLPeerUnverifiedException; @@ -42,35 +49,38 @@ import javax.security.cert.X509Certificate; import java.util.ArrayList; import java.util.List; +import static io.netty.handler.codec.mqtt.MqttConnectReturnCode.*; +import static io.netty.handler.codec.mqtt.MqttMessageType.*; +import static io.netty.handler.codec.mqtt.MqttQoS.*; +import static org.thingsboard.server.common.msg.session.MsgType.*; +import static org.thingsboard.server.transport.mqtt.MqttTopics.*; + /** * @author Andrew Shvayka */ @Slf4j public class MqttTransportHandler extends ChannelInboundHandlerAdapter implements GenericFutureListener> { - public static final MqttQoS MAX_SUPPORTED_QOS_LVL = MqttQoS.AT_LEAST_ONCE; - public static final String BASE_TOPIC = "v1/devices/me"; - public static final String ATTRIBUTES_TOPIC = BASE_TOPIC + "/attributes"; - public static final String TELEMETRY_TOPIC = BASE_TOPIC + "/telemetry"; - public static final String ATTRIBUTES_REQUEST_TOPIC_PREFIX = BASE_TOPIC + "/attributes/request/"; - public static final String ATTRIBUTES_RESPONSE_TOPIC_PREFIX = BASE_TOPIC + "/attributes/response/"; - public static final String ATTRIBUTES_RESPONSES_TOPIC = ATTRIBUTES_RESPONSE_TOPIC_PREFIX + "+"; - public static final String RPC_REQUESTS_TOPIC = BASE_TOPIC + "/rpc/request/"; - public static final String RPC_REQUESTS_SUB_TOPIC = RPC_REQUESTS_TOPIC + "+"; - public static final String RPC_RESPONSE_TOPIC = BASE_TOPIC + "/rpc/response/"; - public static final String RPC_RESPONSE_SUB_TOPIC = RPC_RESPONSE_TOPIC + "+"; - private final MqttSessionCtx sessionCtx; + public static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; + + private final DeviceSessionCtx deviceSessionCtx; private final String sessionId; private final MqttTransportAdaptor adaptor; private final SessionMsgProcessor processor; + private final DeviceService deviceService; + private final DeviceAuthService authService; private final SslHandler sslHandler; + private volatile boolean connected; + private volatile GatewaySessionCtx gatewaySessionCtx; - public MqttTransportHandler(SessionMsgProcessor processor, DeviceAuthService authService, + public MqttTransportHandler(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, MqttTransportAdaptor adaptor, SslHandler sslHandler) { this.processor = processor; + this.deviceService = deviceService; + this.authService = authService; this.adaptor = adaptor; - this.sessionCtx = new MqttSessionCtx(processor, authService, adaptor); - this.sessionId = sessionCtx.getSessionId().toUidStr(); + this.deviceSessionCtx = new DeviceSessionCtx(processor, authService, adaptor); + this.sessionId = deviceSessionCtx.getSessionId().toUidStr(); this.sslHandler = sslHandler; } @@ -83,7 +93,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } private void processMqttMsg(ChannelHandlerContext ctx, MqttMessage msg) { - sessionCtx.setChannel(ctx); + deviceSessionCtx.setChannel(ctx); switch (msg.fixedHeader().messageType()) { case CONNECT: processConnect(ctx, (MqttConnectMessage) msg); @@ -98,36 +108,68 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement processUnsubscribe(ctx, (MqttUnsubscribeMessage) msg); break; case PINGREQ: - ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(MqttMessageType.PINGRESP, false, MqttQoS.AT_MOST_ONCE, false, 0))); + if (checkConnected(ctx)) { + ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); + } break; case DISCONNECT: - processDisconnect(ctx); + if (checkConnected(ctx)) { + processDisconnect(ctx); + } break; } } private void processPublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg) { + if (!checkConnected(ctx)) { + return; + } String topicName = mqttMsg.variableHeader().topicName(); int msgId = mqttMsg.variableHeader().messageId(); log.trace("[{}] Processing publish msg [{}][{}]!", sessionId, topicName, msgId); + + if (topicName.startsWith(BASE_GATEWAY_API_TOPIC)) { + AdaptorToSessionActorMsg msg = null; + if (gatewaySessionCtx != null) { + try { + if (topicName.equals(GATEWAY_CONNECT_TOPIC)) { + gatewaySessionCtx.connect(getDeviceName(mqttMsg)); + } else if (topicName.equals(GATEWAY_DISCONNECT_TOPIC)) { + gatewaySessionCtx.disconnect(getDeviceName(mqttMsg)); + } + } catch (RuntimeException | AdaptorException e) { + log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); + } + } + } else { + processDevicePublish(ctx, mqttMsg, topicName, msgId); + } + } + + private String getDeviceName(MqttPublishMessage mqttMsg) throws AdaptorException { + JsonElement json = JsonMqttAdaptor.validateJsonPayload(deviceSessionCtx.getSessionId(), mqttMsg.payload()); + return json.getAsJsonObject().get("device").getAsString(); + } + + private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) { AdaptorToSessionActorMsg msg = null; try { - if (topicName.equals(ATTRIBUTES_TOPIC)) { - msg = adaptor.convertToActorMsg(sessionCtx, MsgType.POST_ATTRIBUTES_REQUEST, mqttMsg); - } else if (topicName.equals(TELEMETRY_TOPIC)) { - msg = adaptor.convertToActorMsg(sessionCtx, MsgType.POST_TELEMETRY_REQUEST, mqttMsg); - } else if (topicName.startsWith(ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { - msg = adaptor.convertToActorMsg(sessionCtx, MsgType.GET_ATTRIBUTES_REQUEST, mqttMsg); + if (topicName.equals(DEVICE_TELEMETRY_TOPIC)) { + msg = adaptor.convertToActorMsg(deviceSessionCtx, POST_TELEMETRY_REQUEST, mqttMsg); + } else if (topicName.equals(DEVICE_ATTRIBUTES_TOPIC)) { + msg = adaptor.convertToActorMsg(deviceSessionCtx, POST_ATTRIBUTES_REQUEST, mqttMsg); + } else if (topicName.startsWith(DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { + msg = adaptor.convertToActorMsg(deviceSessionCtx, GET_ATTRIBUTES_REQUEST, mqttMsg); if (msgId >= 0) { ctx.writeAndFlush(createMqttPubAckMsg(msgId)); } - } else if (topicName.startsWith(RPC_RESPONSE_TOPIC)) { - msg = adaptor.convertToActorMsg(sessionCtx, MsgType.TO_DEVICE_RPC_RESPONSE, mqttMsg); + } else if (topicName.startsWith(DEVICE_RPC_RESPONSE_TOPIC)) { + msg = adaptor.convertToActorMsg(deviceSessionCtx, TO_DEVICE_RPC_RESPONSE, mqttMsg); if (msgId >= 0) { ctx.writeAndFlush(createMqttPubAckMsg(msgId)); } - } else if (topicName.startsWith(RPC_REQUESTS_TOPIC)) { - msg = adaptor.convertToActorMsg(sessionCtx, MsgType.TO_SERVER_RPC_REQUEST, mqttMsg); + } else if (topicName.startsWith(DEVICE_RPC_REQUESTS_TOPIC)) { + msg = adaptor.convertToActorMsg(deviceSessionCtx, TO_SERVER_RPC_REQUEST, mqttMsg); if (msgId >= 0) { ctx.writeAndFlush(createMqttPubAckMsg(msgId)); } @@ -135,60 +177,65 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } catch (AdaptorException e) { log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); } - if (msg != null) { - processor.process(new BasicToDeviceActorSessionMsg(sessionCtx.getDevice(), msg)); + processor.process(new BasicToDeviceActorSessionMsg(deviceSessionCtx.getDevice(), msg)); } else { - log.warn("[{}] Closing current session due to invalid publish msg [{}][{}]", sessionId, topicName, msgId); + log.info("[{}] Closing current session due to invalid publish msg [{}][{}]", sessionId, topicName, msgId); ctx.close(); } } private void processSubscribe(ChannelHandlerContext ctx, MqttSubscribeMessage mqttMsg) { - log.info("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId()); + if (!checkConnected(ctx)) { + return; + } + log.trace("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId()); List grantedQoSList = new ArrayList<>(); for (MqttTopicSubscription subscription : mqttMsg.payload().topicSubscriptions()) { String topicName = subscription.topicName(); //TODO: handle this qos level. MqttQoS reqQoS = subscription.qualityOfService(); try { - if (topicName.equals(ATTRIBUTES_TOPIC)) { - AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(sessionCtx, MsgType.SUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg); - processor.process(new BasicToDeviceActorSessionMsg(sessionCtx.getDevice(), msg)); + if (topicName.equals(DEVICE_ATTRIBUTES_TOPIC)) { + AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, SUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg); + processor.process(new BasicToDeviceActorSessionMsg(deviceSessionCtx.getDevice(), msg)); grantedQoSList.add(getMinSupportedQos(reqQoS)); - } else if (topicName.equals(RPC_REQUESTS_SUB_TOPIC)) { - AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(sessionCtx, MsgType.SUBSCRIBE_RPC_COMMANDS_REQUEST, mqttMsg); - processor.process(new BasicToDeviceActorSessionMsg(sessionCtx.getDevice(), msg)); + } else if (topicName.equals(DEVICE_RPC_REQUESTS_SUB_TOPIC)) { + AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, SUBSCRIBE_RPC_COMMANDS_REQUEST, mqttMsg); + processor.process(new BasicToDeviceActorSessionMsg(deviceSessionCtx.getDevice(), msg)); grantedQoSList.add(getMinSupportedQos(reqQoS)); - } else if (topicName.equals(RPC_RESPONSE_SUB_TOPIC)) { + } else if (topicName.equals(DEVICE_RPC_RESPONSE_SUB_TOPIC)) { grantedQoSList.add(getMinSupportedQos(reqQoS)); - } else if (topicName.equals(ATTRIBUTES_RESPONSES_TOPIC)) { - sessionCtx.setAllowAttributeResponses(); + } else if (topicName.equals(DEVICE_ATTRIBUTES_RESPONSES_TOPIC)) { + deviceSessionCtx.setAllowAttributeResponses(); grantedQoSList.add(getMinSupportedQos(reqQoS)); } else { log.warn("[{}] Failed to subscribe to [{}][{}]", sessionId, topicName, reqQoS); - grantedQoSList.add(MqttQoS.FAILURE.value()); + grantedQoSList.add(FAILURE.value()); } } catch (AdaptorException e) { log.warn("[{}] Failed to subscribe to [{}][{}]", sessionId, topicName, reqQoS); - grantedQoSList.add(MqttQoS.FAILURE.value()); + grantedQoSList.add(FAILURE.value()); } } ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList)); } private void processUnsubscribe(ChannelHandlerContext ctx, MqttUnsubscribeMessage mqttMsg) { - log.info("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId()); + if (!checkConnected(ctx)) { + return; + } + log.trace("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId()); for (String topicName : mqttMsg.payload().topics()) { try { - if (topicName.equals(ATTRIBUTES_TOPIC)) { - AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(sessionCtx, MsgType.UNSUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg); - processor.process(new BasicToDeviceActorSessionMsg(sessionCtx.getDevice(), msg)); - } else if (topicName.equals(RPC_REQUESTS_SUB_TOPIC)) { - AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(sessionCtx, MsgType.UNSUBSCRIBE_RPC_COMMANDS_REQUEST, mqttMsg); - processor.process(new BasicToDeviceActorSessionMsg(sessionCtx.getDevice(), msg)); - } else if (topicName.equals(ATTRIBUTES_RESPONSES_TOPIC)) { - sessionCtx.setDisallowAttributeResponses(); + if (topicName.equals(DEVICE_ATTRIBUTES_TOPIC)) { + AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, UNSUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg); + processor.process(new BasicToDeviceActorSessionMsg(deviceSessionCtx.getDevice(), msg)); + } else if (topicName.equals(DEVICE_RPC_REQUESTS_SUB_TOPIC)) { + AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, UNSUBSCRIBE_RPC_COMMANDS_REQUEST, mqttMsg); + processor.process(new BasicToDeviceActorSessionMsg(deviceSessionCtx.getDevice(), msg)); + } else if (topicName.equals(DEVICE_ATTRIBUTES_RESPONSES_TOPIC)) { + deviceSessionCtx.setDisallowAttributeResponses(); } } catch (AdaptorException e) { log.warn("[{}] Failed to process unsubscription [{}] to [{}]", sessionId, mqttMsg.variableHeader().messageId(), topicName); @@ -199,7 +246,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private MqttMessage createUnSubAckMessage(int msgId) { MqttFixedHeader mqttFixedHeader = - new MqttFixedHeader(MqttMessageType.SUBACK, false, MqttQoS.AT_LEAST_ONCE, false, 0); + new MqttFixedHeader(SUBACK, false, AT_LEAST_ONCE, false, 0); MqttMessageIdVariableHeader mqttMessageIdVariableHeader = MqttMessageIdVariableHeader.from(msgId); return new MqttMessage(mqttFixedHeader, mqttMessageIdVariableHeader); } @@ -217,13 +264,15 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void processAuthTokenConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) { String userName = msg.payload().userName(); if (StringUtils.isEmpty(userName)) { - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD)); ctx.close(); - } else if (sessionCtx.login(new DeviceTokenCredentials(msg.payload().userName()))) { - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_ACCEPTED)); - } else { - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_NOT_AUTHORIZED)); + } else if (!deviceSessionCtx.login(new DeviceTokenCredentials(msg.payload().userName()))) { + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); ctx.close(); + } else { + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); + connected = true; + checkGatewaySession(); } } @@ -231,14 +280,16 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { String strCert = SslUtil.getX509CertificateString(cert); String sha3Hash = EncryptionUtil.getSha3Hash(strCert); - if (sessionCtx.login(new DeviceX509Credentials(sha3Hash))) { - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_ACCEPTED)); + if (deviceSessionCtx.login(new DeviceX509Credentials(sha3Hash))) { + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); + connected = true; + checkGatewaySession(); } else { - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_NOT_AUTHORIZED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); ctx.close(); } } catch (Exception e) { - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_NOT_AUTHORIZED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); ctx.close(); } } @@ -262,7 +313,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode) { MqttFixedHeader mqttFixedHeader = - new MqttFixedHeader(MqttMessageType.CONNACK, false, MqttQoS.AT_MOST_ONCE, false, 0); + new MqttFixedHeader(CONNACK, false, AT_MOST_ONCE, false, 0); MqttConnAckVariableHeader mqttConnAckVariableHeader = new MqttConnAckVariableHeader(returnCode, true); return new MqttConnAckMessage(mqttFixedHeader, mqttConnAckVariableHeader); @@ -281,7 +332,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private static MqttSubAckMessage createSubAckMessage(Integer msgId, List grantedQoSList) { MqttFixedHeader mqttFixedHeader = - new MqttFixedHeader(MqttMessageType.SUBACK, false, MqttQoS.AT_LEAST_ONCE, false, 0); + new MqttFixedHeader(SUBACK, false, AT_LEAST_ONCE, false, 0); MqttMessageIdVariableHeader mqttMessageIdVariableHeader = MqttMessageIdVariableHeader.from(msgId); MqttSubAckPayload mqttSubAckPayload = new MqttSubAckPayload(grantedQoSList); return new MqttSubAckMessage(mqttFixedHeader, mqttMessageIdVariableHeader, mqttSubAckPayload); @@ -293,14 +344,32 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public static MqttPubAckMessage createMqttPubAckMsg(int requestId) { MqttFixedHeader mqttFixedHeader = - new MqttFixedHeader(MqttMessageType.PUBACK, false, MqttQoS.AT_LEAST_ONCE, false, 0); + new MqttFixedHeader(PUBACK, false, AT_LEAST_ONCE, false, 0); MqttMessageIdVariableHeader mqttMsgIdVariableHeader = MqttMessageIdVariableHeader.from(requestId); return new MqttPubAckMessage(mqttFixedHeader, mqttMsgIdVariableHeader); } + private boolean checkConnected(ChannelHandlerContext ctx) { + if (connected) { + return true; + } else { + log.info("[{}] Closing current session due to invalid msg order [{}][{}]", sessionId); + ctx.close(); + return false; + } + } + + private void checkGatewaySession() { + Device device = deviceSessionCtx.getDevice(); + JsonNode gatewayNode = device.getAdditionalInfo().get("gateway"); + if (gatewayNode != null && gatewayNode.asBoolean()) { + gatewaySessionCtx = new GatewaySessionCtx(processor, deviceService, authService, device); + } + } + @Override public void operationComplete(Future future) throws Exception { - processor.process(SessionCloseMsg.onError(sessionCtx.getSessionId())); + processor.process(SessionCloseMsg.onError(deviceSessionCtx.getSessionId())); } } diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java index 323ed1e4b8..9444cdd1c0 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java @@ -29,6 +29,7 @@ import io.netty.handler.ssl.util.SelfSignedCertificate; import org.springframework.beans.factory.annotation.Value; import org.thingsboard.server.common.transport.SessionMsgProcessor; import org.thingsboard.server.common.transport.auth.DeviceAuthService; +import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; import javax.net.ssl.SSLException; @@ -40,13 +41,15 @@ import java.security.cert.CertificateException; public class MqttTransportServerInitializer extends ChannelInitializer { private final SessionMsgProcessor processor; + private final DeviceService deviceService; private final DeviceAuthService authService; private final MqttTransportAdaptor adaptor; private final MqttSslHandlerProvider sslHandlerProvider; - public MqttTransportServerInitializer(SessionMsgProcessor processor, DeviceAuthService authService, MqttTransportAdaptor adaptor, + public MqttTransportServerInitializer(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, MqttTransportAdaptor adaptor, MqttSslHandlerProvider sslHandlerProvider) { this.processor = processor; + this.deviceService = deviceService; this.authService = authService; this.adaptor = adaptor; this.sslHandlerProvider = sslHandlerProvider; @@ -63,7 +66,7 @@ public class MqttTransportServerInitializer extends ChannelInitializer convertToAdaptorMsg(MqttSessionCtx ctx, SessionActorToAdaptorMsg sessionMsg) throws AdaptorException { + public Optional convertToAdaptorMsg(DeviceSessionCtx ctx, SessionActorToAdaptorMsg sessionMsg) throws AdaptorException { MqttMessage result = null; ToDeviceMsg msg = sessionMsg.getMsg(); switch (msg.getMsgType()) { @@ -100,7 +108,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { GetAttributesResponse response = (GetAttributesResponse) msg; if (response.isSuccess()) { result = createMqttPublishMsg(ctx, - MqttTransportHandler.ATTRIBUTES_RESPONSE_TOPIC_PREFIX + requestId, + MqttTopics.DEVICE_ATTRIBUTES_RESPONSE_TOPIC_PREFIX + requestId, response.getData().get(), true); } else { throw new AdaptorException(response.getError().get()); @@ -115,16 +123,16 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { break; case ATTRIBUTES_UPDATE_NOTIFICATION: AttributesUpdateNotification notification = (AttributesUpdateNotification) msg; - result = createMqttPublishMsg(ctx, MqttTransportHandler.ATTRIBUTES_TOPIC, notification.getData(), false); + result = createMqttPublishMsg(ctx, MqttTopics.DEVICE_ATTRIBUTES_TOPIC, notification.getData(), false); break; case TO_DEVICE_RPC_REQUEST: ToDeviceRpcRequestMsg rpcRequest = (ToDeviceRpcRequestMsg) msg; - result = createMqttPublishMsg(ctx, MqttTransportHandler.RPC_REQUESTS_TOPIC + rpcRequest.getRequestId(), + result = createMqttPublishMsg(ctx, MqttTopics.DEVICE_RPC_REQUESTS_TOPIC + rpcRequest.getRequestId(), rpcRequest); break; case TO_SERVER_RPC_RESPONSE: ToServerRpcResponseMsg rpcResponse = (ToServerRpcResponseMsg) msg; - result = createMqttPublishMsg(ctx, MqttTransportHandler.RPC_REQUESTS_TOPIC + rpcResponse.getRequestId(), + result = createMqttPublishMsg(ctx, MqttTopics.DEVICE_RPC_REQUESTS_TOPIC + rpcResponse.getRequestId(), rpcResponse); break; case RULE_ENGINE_ERROR: @@ -135,19 +143,19 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { return Optional.ofNullable(result); } - private MqttPublishMessage createMqttPublishMsg(MqttSessionCtx ctx, String topic, AttributesKVMsg msg, boolean asMap) { + private MqttPublishMessage createMqttPublishMsg(DeviceSessionCtx ctx, String topic, AttributesKVMsg msg, boolean asMap) { return createMqttPublishMsg(ctx, topic, JsonConverter.toJson(msg, asMap)); } - private MqttPublishMessage createMqttPublishMsg(MqttSessionCtx ctx, String topic, ToDeviceRpcRequestMsg msg) { + private MqttPublishMessage createMqttPublishMsg(DeviceSessionCtx ctx, String topic, ToDeviceRpcRequestMsg msg) { return createMqttPublishMsg(ctx, topic, JsonConverter.toJson(msg, false)); } - private MqttPublishMessage createMqttPublishMsg(MqttSessionCtx ctx, String topic, ToServerRpcResponseMsg msg) { + private MqttPublishMessage createMqttPublishMsg(DeviceSessionCtx ctx, String topic, ToServerRpcResponseMsg msg) { return createMqttPublishMsg(ctx, topic, JsonConverter.toJson(msg)); } - private MqttPublishMessage createMqttPublishMsg(MqttSessionCtx ctx, String topic, JsonElement json) { + private MqttPublishMessage createMqttPublishMsg(DeviceSessionCtx ctx, String topic, JsonElement json) { MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, false, MqttQoS.AT_LEAST_ONCE, false, 0); MqttPublishVariableHeader header = new MqttPublishVariableHeader(topic, ctx.nextMsgId()); @@ -156,10 +164,10 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { return new MqttPublishMessage(mqttFixedHeader, header, payload); } - private FromDeviceMsg convertToGetAttributesRequest(MqttSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException { + private FromDeviceMsg convertToGetAttributesRequest(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException { String topicName = inbound.variableHeader().topicName(); try { - Integer requestId = Integer.valueOf(topicName.substring(MqttTransportHandler.ATTRIBUTES_REQUEST_TOPIC_PREFIX.length())); + Integer requestId = Integer.valueOf(topicName.substring(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX.length())); String payload = inbound.payload().toString(UTF8); JsonElement requestBody = new JsonParser().parse(payload); Set clientKeys = toStringSet(requestBody, "clientKeys"); @@ -175,10 +183,10 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } } - private FromDeviceMsg convertToRpcCommandResponse(MqttSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException { + private FromDeviceMsg convertToRpcCommandResponse(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException { String topicName = inbound.variableHeader().topicName(); try { - Integer requestId = Integer.valueOf(topicName.substring(MqttTransportHandler.RPC_RESPONSE_TOPIC.length())); + Integer requestId = Integer.valueOf(topicName.substring(MqttTopics.DEVICE_RPC_RESPONSE_TOPIC.length())); String payload = inbound.payload().toString(UTF8); return new ToDeviceRpcResponseMsg( requestId, @@ -199,7 +207,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } private UpdateAttributesRequest convertToUpdateAttributesRequest(SessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { - String payload = validatePayload(ctx, inbound.payload()); + String payload = validatePayload(ctx.getSessionId(), inbound.payload()); try { return JsonConverter.convertToAttributes(new JsonParser().parse(payload), inbound.variableHeader().messageId()); } catch (IllegalStateException | JsonSyntaxException ex) { @@ -208,7 +216,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } private TelemetryUploadRequest convertToTelemetryUploadRequest(SessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { - String payload = validatePayload(ctx, inbound.payload()); + String payload = validatePayload(ctx.getSessionId(), inbound.payload()); try { return JsonConverter.convertToTelemetry(new JsonParser().parse(payload), inbound.variableHeader().messageId()); } catch (IllegalStateException | JsonSyntaxException ex) { @@ -216,22 +224,31 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } } - private FromDeviceMsg convertToServerRpcRequest(MqttSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException { + private FromDeviceMsg convertToServerRpcRequest(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException { String topicName = inbound.variableHeader().topicName(); - String payload = validatePayload(ctx, inbound.payload()); + String payload = validatePayload(ctx.getSessionId(), inbound.payload()); try { - Integer requestId = Integer.valueOf(topicName.substring(MqttTransportHandler.RPC_REQUESTS_TOPIC.length())); + Integer requestId = Integer.valueOf(topicName.substring(MqttTopics.DEVICE_RPC_REQUESTS_TOPIC.length())); return JsonConverter.convertToServerRpcRequest(new JsonParser().parse(payload), requestId); } catch (IllegalStateException | JsonSyntaxException ex) { throw new AdaptorException(ex); } } - private String validatePayload(SessionContext ctx, ByteBuf payloadData) throws AdaptorException { + public static JsonElement validateJsonPayload(SessionId sessionId, ByteBuf payloadData) throws AdaptorException { + String payload = validatePayload(sessionId, payloadData); + try { + return new JsonParser().parse(payload); + } catch (JsonSyntaxException ex) { + throw new AdaptorException(ex); + } + } + + public static String validatePayload(SessionId sessionId, ByteBuf payloadData) throws AdaptorException { try { String payload = payloadData.toString(UTF8); if (payload == null) { - log.warn("[{}] Payload is empty!", ctx.getSessionId()); + log.warn("[{}] Payload is empty!", sessionId); throw new AdaptorException(new IllegalArgumentException("Payload is empty!")); } return payload; diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java new file mode 100644 index 0000000000..02f29be016 --- /dev/null +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java @@ -0,0 +1,49 @@ +/** + * Copyright © 2016-2017 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.adaptors; + +import com.google.gson.Gson; +import io.netty.buffer.ByteBufAllocator; +import io.netty.buffer.UnpooledByteBufAllocator; +import io.netty.handler.codec.mqtt.MqttMessage; +import org.thingsboard.server.common.msg.session.AdaptorToSessionActorMsg; +import org.thingsboard.server.common.msg.session.MsgType; +import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg; +import org.thingsboard.server.common.transport.adaptor.AdaptorException; +import org.thingsboard.server.transport.mqtt.session.GatewaySessionCtx; + +import java.nio.charset.Charset; +import java.util.Optional; + +/** + * Created by ashvayka on 19.01.17. + */ +public class JsonMqttGatewayAdaptor implements MqttGatewayAdaptor { + + private static final Gson GSON = new Gson(); + private static final Charset UTF8 = Charset.forName("UTF-8"); + private static final ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false); + + @Override + public AdaptorToSessionActorMsg convertToActorMsg(GatewaySessionCtx ctx, MsgType type, MqttMessage inbound) throws AdaptorException { + return null; + } + + @Override + public Optional convertToAdaptorMsg(GatewaySessionCtx ctx, SessionActorToAdaptorMsg msg) throws AdaptorException { + return null; + } +} diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttGatewayAdaptor.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttGatewayAdaptor.java new file mode 100644 index 0000000000..5641af593d --- /dev/null +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttGatewayAdaptor.java @@ -0,0 +1,36 @@ +/** + * Copyright © 2016-2017 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.adaptors; + +import io.netty.handler.codec.mqtt.MqttMessage; +import org.thingsboard.server.common.msg.session.AdaptorToSessionActorMsg; +import org.thingsboard.server.common.msg.session.MsgType; +import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg; +import org.thingsboard.server.common.transport.adaptor.AdaptorException; +import org.thingsboard.server.transport.mqtt.session.GatewaySessionCtx; + +import java.util.Optional; + +/** + * Created by ashvayka on 19.01.17. + */ +public interface MqttGatewayAdaptor { + + AdaptorToSessionActorMsg convertToActorMsg(GatewaySessionCtx ctx, MsgType type, MqttMessage inbound) throws AdaptorException; + + Optional convertToAdaptorMsg(GatewaySessionCtx ctx, SessionActorToAdaptorMsg msg) throws AdaptorException; + +} diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java index 0de11e4d5f..7f8e1d7d69 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java @@ -17,10 +17,10 @@ package org.thingsboard.server.transport.mqtt.adaptors; import io.netty.handler.codec.mqtt.MqttMessage; import org.thingsboard.server.common.transport.TransportAdaptor; -import org.thingsboard.server.transport.mqtt.session.MqttSessionCtx; +import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; /** * @author Andrew Shvayka */ -public interface MqttTransportAdaptor extends TransportAdaptor { +public interface MqttTransportAdaptor extends TransportAdaptor { } diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttSessionCtx.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java similarity index 95% rename from transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttSessionCtx.java rename to transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index 7d7931e021..f7996facc1 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttSessionCtx.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -36,7 +36,7 @@ import java.util.concurrent.atomic.AtomicInteger; * @author Andrew Shvayka */ @Slf4j -public class MqttSessionCtx extends DeviceAwareSessionContext { +public class DeviceSessionCtx extends DeviceAwareSessionContext { private final MqttTransportAdaptor adaptor; private final MqttSessionId sessionId; @@ -44,7 +44,7 @@ public class MqttSessionCtx extends DeviceAwareSessionContext { private volatile boolean allowAttributeResponses; private AtomicInteger msgIdSeq = new AtomicInteger(0); - public MqttSessionCtx(SessionMsgProcessor processor, DeviceAuthService authService, MqttTransportAdaptor adaptor) { + public DeviceSessionCtx(SessionMsgProcessor processor, DeviceAuthService authService, MqttTransportAdaptor adaptor) { super(processor, authService); this.adaptor = adaptor; this.sessionId = new MqttSessionId(); diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java new file mode 100644 index 0000000000..fef155f62b --- /dev/null +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java @@ -0,0 +1,79 @@ +/** + * Copyright © 2016-2017 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.session; + +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.id.SessionId; +import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg; +import org.thingsboard.server.common.msg.session.SessionCtrlMsg; +import org.thingsboard.server.common.msg.session.SessionType; +import org.thingsboard.server.common.msg.session.ex.SessionException; +import org.thingsboard.server.common.transport.SessionMsgProcessor; +import org.thingsboard.server.common.transport.auth.DeviceAuthService; +import org.thingsboard.server.common.transport.session.DeviceAwareSessionContext; + +/** + * Created by ashvayka on 19.01.17. + */ +public class GatewayDeviceSessionCtx extends DeviceAwareSessionContext { + + private final MqttSessionId sessionId; + private volatile boolean closed; + + public GatewayDeviceSessionCtx(SessionMsgProcessor processor, DeviceAuthService authService, Device device) { + super(processor, authService, device); + this.sessionId = new MqttSessionId(); + } + + @Override + public SessionId getSessionId() { + return sessionId; + } + + @Override + public SessionType getSessionType() { + return SessionType.ASYNC; + } + + @Override + public void onMsg(SessionActorToAdaptorMsg msg) throws SessionException { + + } + + @Override + public void onMsg(SessionCtrlMsg msg) throws SessionException { + + } + + @Override + public void onError(SessionException e) { + + } + + @Override + public boolean isClosed() { + return closed; + } + + public void setClosed(boolean closed) { + this.closed = closed; + } + + @Override + public long getTimeout() { + return 0; + } +} diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java new file mode 100644 index 0000000000..54336d9982 --- /dev/null +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java @@ -0,0 +1,69 @@ +/** + * Copyright © 2016-2017 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.session; + +import java.util.HashMap; +import java.util.Map; +import java.util.Optional; + +import org.springframework.util.StringUtils; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.transport.SessionMsgProcessor; +import org.thingsboard.server.common.transport.auth.DeviceAuthService; +import org.thingsboard.server.dao.device.DeviceService; + +/** + * Created by ashvayka on 19.01.17. + */ +public class GatewaySessionCtx { + + private final Device gateway; + private final SessionMsgProcessor processor; + private final DeviceService deviceService; + private final DeviceAuthService authService; + private final Map devices; + + public GatewaySessionCtx(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, Device gateway) { + this.processor = processor; + this.deviceService = deviceService; + this.authService = authService; + this.gateway = gateway; + this.devices = new HashMap<>(); + } + + public void connect(String deviceName) { + checkDeviceName(deviceName); + Optional deviceOpt = deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), deviceName); + Device device = deviceOpt.orElseGet(() -> { + Device newDevice = new Device(); + newDevice.setTenantId(gateway.getTenantId()); + return deviceService.saveDevice(newDevice); + }); + devices.put(deviceName, new GatewayDeviceSessionCtx(processor, authService, device)); + } + + public void disconnect(String deviceName) { + checkDeviceName(deviceName); + devices.remove(deviceName); + } + + private void checkDeviceName(String deviceName) { + if (StringUtils.isEmpty(deviceName)) { + throw new RuntimeException(); + } + } + +} From 0da7724c09eee28228f4baca4bcaf29a43a1ab50 Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Fri, 20 Jan 2017 17:12:48 +0200 Subject: [PATCH 2/3] TB-33: Implementation --- .../src/main/resources/thingsboard.yml | 8 +- .../common/msg/session/SessionContext.java | 2 - .../msg/session/ctrl/SessionCloseMsg.java | 4 + .../transport/adaptor/JsonConverter.java | 4 +- .../coap/session/CoapSessionCtx.java | 11 -- .../http/session/HttpSessionCtx.java | 5 - .../transport/mqtt/MqttTransportHandler.java | 32 ++-- .../mqtt/adaptors/JsonMqttAdaptor.java | 2 +- .../mqtt/adaptors/JsonMqttGatewayAdaptor.java | 49 ------ .../mqtt/session/DeviceSessionCtx.java | 5 - .../mqtt/session/GatewayDeviceSessionCtx.java | 42 +++-- .../mqtt/session/GatewaySessionCtx.java | 150 ++++++++++++++++-- 12 files changed, 194 insertions(+), 120 deletions(-) delete mode 100644 transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 288914822d..f150574e66 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -76,10 +76,10 @@ mqtt: adaptor: "${MQTT_ADAPTOR_NAME:JsonMqttAdaptor}" timeout: "${MQTT_TIMEOUT:10000}" # Uncomment the following lines to enable ssl for MQTT - ssl: - key_store: keystore/mqttserver.jks - key_store_password: password - key_store_type: JKS +# ssl: +# key_store: keystore/mqttserver.jks +# key_store_password: password +# key_store_type: JKS # CoAP server parameters coap: 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 4a392253ab..0b138dfaac 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 @@ -27,8 +27,6 @@ public interface SessionContext extends SessionAwareMsg { void onMsg(SessionCtrlMsg msg) throws SessionException; - void onError(SessionException e); - boolean isClosed(); long getTimeout(); diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/session/ctrl/SessionCloseMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/session/ctrl/SessionCloseMsg.java index 6d957ceef6..c7baaafc70 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/session/ctrl/SessionCloseMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/session/ctrl/SessionCloseMsg.java @@ -24,6 +24,10 @@ public class SessionCloseMsg implements SessionCtrlMsg { private final boolean revoked; private final boolean timeout; + public static SessionCloseMsg onDisconnect(SessionId sessionId) { + return new SessionCloseMsg(sessionId, false, false); + } + public static SessionCloseMsg onError(SessionId sessionId) { return new SessionCloseMsg(sessionId, false, false); } diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java index 93e764ef47..640cce798b 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java @@ -74,7 +74,7 @@ public class JsonConverter { } } - private static void parseWithTs(BasicTelemetryUploadRequest request, JsonObject jo) { + public static void parseWithTs(BasicTelemetryUploadRequest request, JsonObject jo) { long ts = jo.get("ts").getAsLong(); JsonObject valuesObject = jo.get("values").getAsJsonObject(); for (KvEntry entry : parseValues(valuesObject)) { @@ -82,7 +82,7 @@ public class JsonConverter { } } - private static List parseValues(JsonObject valuesObject) { + public static List parseValues(JsonObject valuesObject) { List result = new ArrayList<>(); for (Entry valueEntry : valuesObject.entrySet()) { JsonElement element = valueEntry.getValue(); diff --git a/transport/coap/src/main/java/org/thingsboard/server/transport/coap/session/CoapSessionCtx.java b/transport/coap/src/main/java/org/thingsboard/server/transport/coap/session/CoapSessionCtx.java index e9b8e22404..cecc42d31a 100644 --- a/transport/coap/src/main/java/org/thingsboard/server/transport/coap/session/CoapSessionCtx.java +++ b/transport/coap/src/main/java/org/thingsboard/server/transport/coap/session/CoapSessionCtx.java @@ -95,17 +95,6 @@ public class CoapSessionCtx extends DeviceAwareSessionContext { } } - @Override - public void onError(SessionException e) { - if (e instanceof SessionAuthException) { - log.warn("[{}] onError: {}", sessionId, e.getMessage()); - exchange.respond(ResponseCode.UNAUTHORIZED); - } else { - log.warn("[{}] onError: {}", sessionId, e.getMessage(), e); - exchange.respond(ResponseCode.BAD_REQUEST); - } - } - @Override public SessionId getSessionId() { return sessionId; diff --git a/transport/http/src/main/java/org/thingsboard/server/transport/http/session/HttpSessionCtx.java b/transport/http/src/main/java/org/thingsboard/server/transport/http/session/HttpSessionCtx.java index efaa0cdd2b..4bee59533d 100644 --- a/transport/http/src/main/java/org/thingsboard/server/transport/http/session/HttpSessionCtx.java +++ b/transport/http/src/main/java/org/thingsboard/server/transport/http/session/HttpSessionCtx.java @@ -140,11 +140,6 @@ public class HttpSessionCtx extends DeviceAwareSessionContext { } - @Override - public void onError(SessionException e) { - - } - @Override public boolean isClosed() { return false; diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index e74495572f..fbc53ff843 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -16,7 +16,6 @@ package org.thingsboard.server.transport.mqtt; import com.fasterxml.jackson.databind.JsonNode; -import com.google.gson.JsonElement; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; import io.netty.handler.codec.mqtt.*; @@ -38,7 +37,6 @@ import org.thingsboard.server.common.transport.adaptor.AdaptorException; import org.thingsboard.server.common.transport.auth.DeviceAuthService; import org.thingsboard.server.dao.EncryptionUtil; import org.thingsboard.server.dao.device.DeviceService; -import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; import org.thingsboard.server.transport.mqtt.session.GatewaySessionCtx; import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; @@ -129,13 +127,17 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement log.trace("[{}] Processing publish msg [{}][{}]!", sessionId, topicName, msgId); if (topicName.startsWith(BASE_GATEWAY_API_TOPIC)) { - AdaptorToSessionActorMsg msg = null; if (gatewaySessionCtx != null) { + gatewaySessionCtx.setChannel(ctx); try { - if (topicName.equals(GATEWAY_CONNECT_TOPIC)) { - gatewaySessionCtx.connect(getDeviceName(mqttMsg)); + if (topicName.equals(GATEWAY_TELEMETRY_TOPIC)) { + gatewaySessionCtx.onDeviceTelemetry(mqttMsg); + } else if (topicName.equals(GATEWAY_ATTRIBUTES_TOPIC)) { + gatewaySessionCtx.onDeviceAttributes(mqttMsg); + } else if (topicName.equals(GATEWAY_CONNECT_TOPIC)) { + gatewaySessionCtx.onDeviceConnect(mqttMsg); } else if (topicName.equals(GATEWAY_DISCONNECT_TOPIC)) { - gatewaySessionCtx.disconnect(getDeviceName(mqttMsg)); + gatewaySessionCtx.onDeviceDisconnect(mqttMsg); } } catch (RuntimeException | AdaptorException e) { log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); @@ -146,11 +148,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private String getDeviceName(MqttPublishMessage mqttMsg) throws AdaptorException { - JsonElement json = JsonMqttAdaptor.validateJsonPayload(deviceSessionCtx.getSessionId(), mqttMsg.payload()); - return json.getAsJsonObject().get("device").getAsString(); - } - private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) { AdaptorToSessionActorMsg msg = null; try { @@ -309,6 +306,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void processDisconnect(ChannelHandlerContext ctx) { ctx.close(); + processor.process(SessionCloseMsg.onDisconnect(deviceSessionCtx.getSessionId())); + if (gatewaySessionCtx != null) { + gatewaySessionCtx.onGatewayDisconnect(); + } } private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode) { @@ -362,9 +363,12 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void checkGatewaySession() { Device device = deviceSessionCtx.getDevice(); - JsonNode gatewayNode = device.getAdditionalInfo().get("gateway"); - if (gatewayNode != null && gatewayNode.asBoolean()) { - gatewaySessionCtx = new GatewaySessionCtx(processor, deviceService, authService, device); + JsonNode infoNode = device.getAdditionalInfo(); + if (infoNode != null) { + JsonNode gatewayNode = infoNode.get("gateway"); + if (gatewayNode != null && gatewayNode.asBoolean()) { + gatewaySessionCtx = new GatewaySessionCtx(processor, deviceService, authService, deviceSessionCtx); + } } } diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index ae49cd589a..bf033dcae3 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -248,7 +248,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { String payload = payloadData.toString(UTF8); if (payload == null) { - log.warn("[{}] Payload is empty!", sessionId); + log.warn("[{}] Payload is empty!", sessionId.toUidStr()); throw new AdaptorException(new IllegalArgumentException("Payload is empty!")); } return payload; diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java deleted file mode 100644 index 02f29be016..0000000000 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttGatewayAdaptor.java +++ /dev/null @@ -1,49 +0,0 @@ -/** - * Copyright © 2016-2017 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.adaptors; - -import com.google.gson.Gson; -import io.netty.buffer.ByteBufAllocator; -import io.netty.buffer.UnpooledByteBufAllocator; -import io.netty.handler.codec.mqtt.MqttMessage; -import org.thingsboard.server.common.msg.session.AdaptorToSessionActorMsg; -import org.thingsboard.server.common.msg.session.MsgType; -import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg; -import org.thingsboard.server.common.transport.adaptor.AdaptorException; -import org.thingsboard.server.transport.mqtt.session.GatewaySessionCtx; - -import java.nio.charset.Charset; -import java.util.Optional; - -/** - * Created by ashvayka on 19.01.17. - */ -public class JsonMqttGatewayAdaptor implements MqttGatewayAdaptor { - - private static final Gson GSON = new Gson(); - private static final Charset UTF8 = Charset.forName("UTF-8"); - private static final ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false); - - @Override - public AdaptorToSessionActorMsg convertToActorMsg(GatewaySessionCtx ctx, MsgType type, MqttMessage inbound) throws AdaptorException { - return null; - } - - @Override - public Optional convertToAdaptorMsg(GatewaySessionCtx ctx, SessionActorToAdaptorMsg msg) throws AdaptorException { - return null; - } -} diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index f7996facc1..f458b86248 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -82,11 +82,6 @@ public class DeviceSessionCtx extends DeviceAwareSessionContext { } } - @Override - public void onError(SessionException e) { - - } - @Override public boolean isClosed() { return false; diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java index fef155f62b..9c4bacfd87 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java @@ -15,26 +15,29 @@ */ package org.thingsboard.server.transport.mqtt.session; +import io.netty.handler.codec.mqtt.MqttMessage; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.id.SessionId; -import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg; -import org.thingsboard.server.common.msg.session.SessionCtrlMsg; -import org.thingsboard.server.common.msg.session.SessionType; +import org.thingsboard.server.common.msg.core.ResponseMsg; +import org.thingsboard.server.common.msg.session.*; import org.thingsboard.server.common.msg.session.ex.SessionException; -import org.thingsboard.server.common.transport.SessionMsgProcessor; -import org.thingsboard.server.common.transport.auth.DeviceAuthService; import org.thingsboard.server.common.transport.session.DeviceAwareSessionContext; +import org.thingsboard.server.transport.mqtt.MqttTransportHandler; + +import java.util.Optional; /** * Created by ashvayka on 19.01.17. */ public class GatewayDeviceSessionCtx extends DeviceAwareSessionContext { + private GatewaySessionCtx parent; private final MqttSessionId sessionId; private volatile boolean closed; - public GatewayDeviceSessionCtx(SessionMsgProcessor processor, DeviceAuthService authService, Device device) { - super(processor, authService, device); + public GatewayDeviceSessionCtx(GatewaySessionCtx parent, Device device) { + super(parent.getProcessor(), parent.getAuthService(), device); + this.parent = parent; this.sessionId = new MqttSessionId(); } @@ -49,17 +52,30 @@ public class GatewayDeviceSessionCtx extends DeviceAwareSessionContext { } @Override - public void onMsg(SessionActorToAdaptorMsg msg) throws SessionException { - + public void onMsg(SessionActorToAdaptorMsg sessionMsg) throws SessionException { + Optional message = getToDeviceMsg(sessionMsg); + message.ifPresent(parent::writeAndFlush); } - @Override - public void onMsg(SessionCtrlMsg msg) throws SessionException { - + private Optional getToDeviceMsg(SessionActorToAdaptorMsg sessionMsg) { + ToDeviceMsg msg = sessionMsg.getMsg(); + switch (msg.getMsgType()) { + case STATUS_CODE_RESPONSE: + ResponseMsg responseMsg = (ResponseMsg) msg; + if (responseMsg.isSuccess()) { + MsgType requestMsgType = responseMsg.getRequestMsgType(); + Integer requestId = responseMsg.getRequestId(); + if (requestMsgType == MsgType.POST_ATTRIBUTES_REQUEST || requestMsgType == MsgType.POST_TELEMETRY_REQUEST) { + return Optional.of(MqttTransportHandler.createMqttPubAckMsg(requestId)); + } + } + break; + } + return Optional.empty(); } @Override - public void onError(SessionException e) { + public void onMsg(SessionCtrlMsg msg) throws SessionException { } diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java index 54336d9982..2badd3ab79 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java @@ -15,55 +15,177 @@ */ package org.thingsboard.server.transport.mqtt.session; -import java.util.HashMap; -import java.util.Map; -import java.util.Optional; - +import com.google.gson.*; +import io.netty.buffer.ByteBufAllocator; +import io.netty.buffer.UnpooledByteBufAllocator; +import io.netty.channel.ChannelHandlerContext; +import io.netty.handler.codec.mqtt.MqttMessage; +import io.netty.handler.codec.mqtt.MqttPublishMessage; +import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.id.SessionId; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.msg.core.BasicTelemetryUploadRequest; +import org.thingsboard.server.common.msg.core.BasicUpdateAttributesRequest; +import org.thingsboard.server.common.msg.core.TelemetryUploadRequest; +import org.thingsboard.server.common.msg.session.BasicAdaptorToSessionActorMsg; +import org.thingsboard.server.common.msg.session.BasicToDeviceActorSessionMsg; +import org.thingsboard.server.common.msg.session.ctrl.SessionCloseMsg; import org.thingsboard.server.common.transport.SessionMsgProcessor; +import org.thingsboard.server.common.transport.adaptor.AdaptorException; +import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.common.transport.auth.DeviceAuthService; import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.transport.mqtt.MqttTransportHandler; +import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor; + +import java.nio.charset.Charset; +import java.util.HashMap; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; + +import static org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor.validateJsonPayload; /** * Created by ashvayka on 19.01.17. */ +@Slf4j public class GatewaySessionCtx { + private static final Gson GSON = new Gson(); + private static final Charset UTF8 = Charset.forName("UTF-8"); + private static final ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false); + private final Device gateway; + private final SessionId gatewaySessionId; private final SessionMsgProcessor processor; private final DeviceService deviceService; private final DeviceAuthService authService; private final Map devices; + private ChannelHandlerContext channel; - public GatewaySessionCtx(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, Device gateway) { + public GatewaySessionCtx(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, DeviceSessionCtx gatewaySessionCtx) { this.processor = processor; this.deviceService = deviceService; this.authService = authService; - this.gateway = gateway; + this.gateway = gatewaySessionCtx.getDevice(); + this.gatewaySessionId = gatewaySessionCtx.getSessionId(); this.devices = new HashMap<>(); } - public void connect(String deviceName) { - checkDeviceName(deviceName); + public void onDeviceConnect(MqttPublishMessage msg) throws AdaptorException { + String deviceName = checkDeviceName(getDeviceName(msg)); Optional deviceOpt = deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), deviceName); Device device = deviceOpt.orElseGet(() -> { Device newDevice = new Device(); newDevice.setTenantId(gateway.getTenantId()); + newDevice.setName(deviceName); return deviceService.saveDevice(newDevice); }); - devices.put(deviceName, new GatewayDeviceSessionCtx(processor, authService, device)); + devices.put(deviceName, new GatewayDeviceSessionCtx(this, device)); + ack(msg); + } + + public void onDeviceDisconnect(MqttPublishMessage msg) throws AdaptorException { + String deviceName = checkDeviceName(getDeviceName(msg)); + GatewayDeviceSessionCtx deviceSessionCtx = devices.remove(deviceName); + deviceSessionCtx.setClosed(true); + ack(msg); + } + + public void onGatewayDisconnect() { + devices.forEach((k, v) -> { + processor.process(SessionCloseMsg.onDisconnect(v.getSessionId())); + }); } - public void disconnect(String deviceName) { - checkDeviceName(deviceName); - devices.remove(deviceName); + public void onDeviceTelemetry(MqttPublishMessage mqttMsg) throws AdaptorException { + JsonElement json = validateJsonPayload(gatewaySessionId, mqttMsg.payload()); + int requestId = mqttMsg.variableHeader().messageId(); + if (json.isJsonObject()) { + JsonObject jsonObj = json.getAsJsonObject(); + for (Map.Entry deviceEntry : jsonObj.entrySet()) { + String deviceName = checkDeviceConnected(deviceEntry.getKey()); + if (!deviceEntry.getValue().isJsonArray()) { + throw new JsonSyntaxException("Can't parse value: " + json); + } + BasicTelemetryUploadRequest request = new BasicTelemetryUploadRequest(requestId); + JsonArray deviceData = deviceEntry.getValue().getAsJsonArray(); + for (JsonElement element : deviceData) { + JsonConverter.parseWithTs(request, element.getAsJsonObject()); + } + GatewayDeviceSessionCtx deviceSessionCtx = devices.get(deviceName); + processor.process(new BasicToDeviceActorSessionMsg(deviceSessionCtx.getDevice(), + new BasicAdaptorToSessionActorMsg(deviceSessionCtx, request))); + } + } else { + throw new JsonSyntaxException("Can't parse value: " + json); + } } - private void checkDeviceName(String deviceName) { + public void onDeviceAttributes(MqttPublishMessage mqttMsg) throws AdaptorException { + JsonElement json = validateJsonPayload(gatewaySessionId, mqttMsg.payload()); + int requestId = mqttMsg.variableHeader().messageId(); + if (json.isJsonObject()) { + JsonObject jsonObj = json.getAsJsonObject(); + for (Map.Entry deviceEntry : jsonObj.entrySet()) { + String deviceName = checkDeviceConnected(deviceEntry.getKey()); + if (!deviceEntry.getValue().isJsonObject()) { + throw new JsonSyntaxException("Can't parse value: " + json); + } + long ts = System.currentTimeMillis(); + BasicUpdateAttributesRequest request = new BasicUpdateAttributesRequest(requestId); + JsonObject deviceData = deviceEntry.getValue().getAsJsonObject(); + request.add(JsonConverter.parseValues(deviceData).stream().map(kv -> new BaseAttributeKvEntry(kv, ts)).collect(Collectors.toList())); + GatewayDeviceSessionCtx deviceSessionCtx = devices.get(deviceName); + processor.process(new BasicToDeviceActorSessionMsg(deviceSessionCtx.getDevice(), + new BasicAdaptorToSessionActorMsg(deviceSessionCtx, request))); + } + } else { + throw new JsonSyntaxException("Can't parse value: " + json); + } + } + + private String checkDeviceConnected(String deviceName) { + if (!devices.containsKey(deviceName)) { + throw new RuntimeException("Device is not connected!"); + } else { + return deviceName; + } + } + + private String checkDeviceName(String deviceName) { if (StringUtils.isEmpty(deviceName)) { - throw new RuntimeException(); + throw new RuntimeException("Device name is empty!"); + } else { + return deviceName; } } + private String getDeviceName(MqttPublishMessage mqttMsg) throws AdaptorException { + JsonElement json = JsonMqttAdaptor.validateJsonPayload(gatewaySessionId, mqttMsg.payload()); + return json.getAsJsonObject().get("device").getAsString(); + } + + protected SessionMsgProcessor getProcessor() { + return processor; + } + + protected DeviceAuthService getAuthService() { + return authService; + } + + public void setChannel(ChannelHandlerContext channel) { + this.channel = channel; + } + + private void ack(MqttPublishMessage msg) { + writeAndFlush(MqttTransportHandler.createMqttPubAckMsg(msg.variableHeader().messageId())); + } + + protected void writeAndFlush(MqttMessage mqttMessage) { + channel.writeAndFlush(mqttMessage); + } } From bfb27e87bd9789b4e43b82c2be5fb194aa414fec Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Sun, 29 Jan 2017 03:08:08 +0200 Subject: [PATCH 3/3] TB-33: SSL tools improvements --- .../src/main/resources/thingsboard.yml | 2 +- .../one-way-ssl-mqtt-client.py} | 7 ++-- .../simple-mqtt-client.py} | 0 .../two-way-ssl-mqtt-client.py} | 2 +- ...emqttclient.keygen.sh => client.keygen.sh} | 6 +-- tools/src/main/shell/keygen.properties | 10 ++--- tools/src/main/shell/server.keygen.sh | 38 +++++++++---------- 7 files changed, 32 insertions(+), 33 deletions(-) rename tools/src/main/{shell/onewaysslmqttclient.py => python/one-way-ssl-mqtt-client.py} (88%) rename tools/src/main/{shell/simplemqttclient.py => python/simple-mqtt-client.py} (100%) rename tools/src/main/{shell/twowaysslmqttclient.py => python/two-way-ssl-mqtt-client.py} (97%) rename tools/src/main/shell/{securemqttclient.keygen.sh => client.keygen.sh} (95%) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index a3ca9d0c16..d2cac80101 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -81,7 +81,7 @@ mqtt: worker_group_thread_count: "${NETTY_WORKER_GROUP_THREADS:12}" # Uncomment the following lines to enable ssl for MQTT # ssl: -# key_store: keystore/mqttserver.jks +# key_store: mqttserver.jks # key_store_password: server_ks_password # key_password: server_key_password # key_store_type: JKS diff --git a/tools/src/main/shell/onewaysslmqttclient.py b/tools/src/main/python/one-way-ssl-mqtt-client.py similarity index 88% rename from tools/src/main/shell/onewaysslmqttclient.py rename to tools/src/main/python/one-way-ssl-mqtt-client.py index d06face509..9266fbfc0b 100644 --- a/tools/src/main/shell/onewaysslmqttclient.py +++ b/tools/src/main/python/one-way-ssl-mqtt-client.py @@ -1,3 +1,4 @@ +# -*- coding: utf-8 -*- # # Copyright © 2016-2017 The Thingsboard Authors # @@ -41,14 +42,12 @@ client.on_connect = on_connect client.on_message = on_message client.publish('v1/devices/me/attributes/request/1', "{\"clientKeys\":\"model\"}", 1) -#client.tls_set(ca_certs="client_truststore.pem", certfile="mqttclient.nopass.pem", keyfile=None, cert_reqs=ssl.CERT_REQUIRED, -# tls_version=ssl.PROTOCOL_TLSv1, ciphers=None); client.tls_set(ca_certs="mqttserver.pub.pem", certfile=None, keyfile=None, cert_reqs=ssl.CERT_REQUIRED, tls_version=ssl.PROTOCOL_TLSv1, ciphers=None); -client.username_pw_set("B1_TEST_TOKEN") +client.username_pw_set("TEST_TOKEN") client.tls_insecure_set(False) -client.connect(socket.gethostname(), 1883, 1) +client.connect(socket.gethostname(), 8883, 1) # Blocking call that processes network traffic, dispatches callbacks and diff --git a/tools/src/main/shell/simplemqttclient.py b/tools/src/main/python/simple-mqtt-client.py similarity index 100% rename from tools/src/main/shell/simplemqttclient.py rename to tools/src/main/python/simple-mqtt-client.py diff --git a/tools/src/main/shell/twowaysslmqttclient.py b/tools/src/main/python/two-way-ssl-mqtt-client.py similarity index 97% rename from tools/src/main/shell/twowaysslmqttclient.py rename to tools/src/main/python/two-way-ssl-mqtt-client.py index a2fa8b617e..d3b32423e1 100644 --- a/tools/src/main/shell/twowaysslmqttclient.py +++ b/tools/src/main/python/two-way-ssl-mqtt-client.py @@ -46,7 +46,7 @@ client.tls_set(ca_certs="mqttserver.pub.pem", certfile="mqttclient.nopass.pem", tls_version=ssl.PROTOCOL_TLSv1, ciphers=None); client.tls_insecure_set(False) -client.connect(socket.gethostname(), 1883, 1) +client.connect(socket.gethostname(), 8883, 1) # Blocking call that processes network traffic, dispatches callbacks and diff --git a/tools/src/main/shell/securemqttclient.keygen.sh b/tools/src/main/shell/client.keygen.sh similarity index 95% rename from tools/src/main/shell/securemqttclient.keygen.sh rename to tools/src/main/shell/client.keygen.sh index f69dd52ba5..500cd0eaa5 100755 --- a/tools/src/main/shell/securemqttclient.keygen.sh +++ b/tools/src/main/shell/client.keygen.sh @@ -18,7 +18,7 @@ usage() { echo "This script generates client public/private rey pair, extracts them to a no-password RSA pem file," echo "and imports server public key to client keystore" - echo "usage: ./securemqttclient.keygen.sh [-p file]" + echo "usage: ./client.keygen.sh [-p file]" echo " -p | --props | --properties file Properties file. default value is ./keygen.properties" echo " -h | --help | ? Show this message" } @@ -48,7 +48,7 @@ if [ -f $CLIENT_FILE_PREFIX.jks ] || [ -f $CLIENT_FILE_PREFIX.pub.pem ] || [ -f then while : do - read -p "Output files from previous server.keygen.sh script run found. Overwrite?[yes]" response + read -p "Output files from previous server.keygen.sh script run found. Overwrite? [Y/N]: " response case $response in [nN]|[nN][oO]) echo "Skipping" @@ -74,7 +74,7 @@ echo "Generating SSL Key Pair..." keytool -genkeypair -v \ -alias $CLIENT_KEY_ALIAS \ - -dname "CN=$DOMAIN_SUFFIX, OU=Thingsboard, O=Thingsboard, L=Piscataway, ST=NJ, C=US" \ + -dname "CN=$DOMAIN_SUFFIX, OU=Thingsboard, O=Thingsboard, L=San Francisco, ST=CA, C=US" \ -keystore $CLIENT_FILE_PREFIX.jks \ -keypass $CLIENT_KEY_PASSWORD \ -storepass $CLIENT_KEYSTORE_PASSWORD \ diff --git a/tools/src/main/shell/keygen.properties b/tools/src/main/shell/keygen.properties index 9435746faf..8dd11f2ea4 100644 --- a/tools/src/main/shell/keygen.properties +++ b/tools/src/main/shell/keygen.properties @@ -17,8 +17,8 @@ DOMAIN_SUFFIX="$(hostname)" ORGANIZATIONAL_UNIT=Thingsboard ORGANIZATION=Thingsboard -CITY=Piscataway -STATE_OR_PROVINCE=NJ +CITY=San Francisco +STATE_OR_PROVINCE=CA TWO_LETTER_COUNTRY_CODE=US SERVER_KEYSTORE_PASSWORD=server_ks_password @@ -26,10 +26,10 @@ SERVER_KEY_PASSWORD=server_key_password SERVER_KEY_ALIAS="serveralias" SERVER_FILE_PREFIX="mqttserver" -SERVER_KEYSTORE_DIR="../../../../application/src/main/resources/keystore/" +SERVER_KEYSTORE_DIR="/etc/thingsboard/conf" -CLIENT_KEYSTORE_PASSWORD=client_ks_password -CLIENT_KEY_PASSWORD=client_key_password +CLIENT_KEYSTORE_PASSWORD=password +CLIENT_KEY_PASSWORD=password CLIENT_KEY_ALIAS="clientalias" CLIENT_FILE_PREFIX="mqttclient" diff --git a/tools/src/main/shell/server.keygen.sh b/tools/src/main/shell/server.keygen.sh index cfeaa0c501..cfa46837e8 100755 --- a/tools/src/main/shell/server.keygen.sh +++ b/tools/src/main/shell/server.keygen.sh @@ -122,25 +122,25 @@ fi if [[ $COPY = true ]]; then if [[ -z "$COPY_DIR" ]]; then - read -p "Do you want to copy $SERVER_FILE_PREFIX.jks to server directory?[yes]" yn - while : - do - case $yn in - [nN]|[nN][oO]) - break - ;; - [yY]|[yY][eE]|[yY][eE]|[sS]|[yY]|"") - read -p "(Default: $SERVER_KEYSTORE_DIR): " dir - if [[ ! -z $dir ]]; then - DESTINATION=$dir; - else - DESTINATION=$SERVER_KEYSTORE_DIR - fi; - break;; - *) echo "Please reply 'yes' or 'no'" - ;; - esac - done + while : + do + read -p "Do you want to copy $SERVER_FILE_PREFIX.jks to server directory? [Y/N]: " yn + case $yn in + [nN]|[nN][oO]) + break + ;; + [yY]|[yY][eE]|[yY][eE]|[sS]|[yY]|"") + read -p "(Default: $SERVER_KEYSTORE_DIR): " dir + if [[ ! -z $dir ]]; then + DESTINATION=$dir; + else + DESTINATION=$SERVER_KEYSTORE_DIR + fi; + break;; + *) echo "Please reply 'yes' or 'no'" + ;; + esac + done else DESTINATION=$COPY_DIR fi