|
|
@ -17,7 +17,6 @@ package org.thingsboard.server.transport.mqtt; |
|
|
|
|
|
|
|
|
import com.fasterxml.jackson.databind.JsonNode; |
|
|
import com.fasterxml.jackson.databind.JsonNode; |
|
|
import com.google.gson.JsonParseException; |
|
|
import com.google.gson.JsonParseException; |
|
|
import com.google.gson.JsonSyntaxException; |
|
|
|
|
|
import io.netty.channel.ChannelHandlerContext; |
|
|
import io.netty.channel.ChannelHandlerContext; |
|
|
import io.netty.channel.ChannelInboundHandlerAdapter; |
|
|
import io.netty.channel.ChannelInboundHandlerAdapter; |
|
|
import io.netty.handler.codec.mqtt.MqttConnAckMessage; |
|
|
import io.netty.handler.codec.mqtt.MqttConnAckMessage; |
|
|
@ -70,11 +69,10 @@ import java.io.IOException; |
|
|
import java.net.InetSocketAddress; |
|
|
import java.net.InetSocketAddress; |
|
|
import java.util.ArrayList; |
|
|
import java.util.ArrayList; |
|
|
import java.util.List; |
|
|
import java.util.List; |
|
|
import java.util.Optional; |
|
|
|
|
|
import java.util.UUID; |
|
|
import java.util.UUID; |
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
import java.util.concurrent.ConcurrentMap; |
|
|
import java.util.concurrent.ConcurrentMap; |
|
|
import java.util.Date; |
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
|
|
|
|
|
import static io.netty.handler.codec.mqtt.MqttConnectReturnCode.CONNECTION_ACCEPTED; |
|
|
import static io.netty.handler.codec.mqtt.MqttConnectReturnCode.CONNECTION_ACCEPTED; |
|
|
import static io.netty.handler.codec.mqtt.MqttConnectReturnCode.CONNECTION_REFUSED_NOT_AUTHORIZED; |
|
|
import static io.netty.handler.codec.mqtt.MqttConnectReturnCode.CONNECTION_REFUSED_NOT_AUTHORIZED; |
|
|
@ -156,11 +154,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
if (topicName.equals(MqttTopics.DEVICE_PROVISION_REQUEST_TOPIC)) { |
|
|
if (topicName.equals(MqttTopics.DEVICE_PROVISION_REQUEST_TOPIC)) { |
|
|
try { |
|
|
try { |
|
|
TransportProtos.ProvisionDeviceRequestMsg provisionRequestMsg = deviceSessionCtx.getContext().getJsonMqttAdaptor().convertToProvisionRequestMsg(deviceSessionCtx, mqttMsg); |
|
|
TransportProtos.ProvisionDeviceRequestMsg provisionRequestMsg = deviceSessionCtx.getContext().getJsonMqttAdaptor().convertToProvisionRequestMsg(deviceSessionCtx, mqttMsg); |
|
|
|
|
|
validateProvisionMessage(provisionRequestMsg); |
|
|
transportService.process(provisionRequestMsg, new DeviceProvisionCallback(ctx, msgId, provisionRequestMsg)); |
|
|
transportService.process(provisionRequestMsg, new DeviceProvisionCallback(ctx, msgId, provisionRequestMsg)); |
|
|
log.trace("[{}][{}] Processing provision publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId); |
|
|
log.trace("[{}][{}] Processing provision publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId); |
|
|
} catch (Exception e) { |
|
|
} catch (Exception e) { |
|
|
if (e instanceof JsonParseException || (e.getCause() != null && e.getCause() instanceof JsonParseException)) { |
|
|
if (e instanceof JsonParseException || (e.getCause() != null && e.getCause() instanceof JsonParseException)) { |
|
|
TransportProtos.ProvisionDeviceRequestMsg provisionRequestMsg = deviceSessionCtx.getContext().getProtoMqttAdaptor().convertToProvisionRequestMsg(deviceSessionCtx, mqttMsg); |
|
|
TransportProtos.ProvisionDeviceRequestMsg provisionRequestMsg = deviceSessionCtx.getContext().getProtoMqttAdaptor().convertToProvisionRequestMsg(deviceSessionCtx, mqttMsg); |
|
|
|
|
|
validateProvisionMessage(provisionRequestMsg); |
|
|
transportService.process(provisionRequestMsg, new DeviceProvisionCallback(ctx, msgId, provisionRequestMsg)); |
|
|
transportService.process(provisionRequestMsg, new DeviceProvisionCallback(ctx, msgId, provisionRequestMsg)); |
|
|
deviceSessionCtx.setProvisionPayloadType(TransportPayloadType.PROTOBUF); |
|
|
deviceSessionCtx.setProvisionPayloadType(TransportPayloadType.PROTOBUF); |
|
|
log.trace("[{}][{}] Processing provision publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId); |
|
|
log.trace("[{}][{}] Processing provision publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId); |
|
|
@ -187,6 +187,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void validateProvisionMessage(TransportProtos.ProvisionDeviceRequestMsg provisionRequestMsg) { |
|
|
|
|
|
if (provisionRequestMsg.getProvisionDeviceCredentialsMsg().getProvisionDeviceKey() != null && |
|
|
|
|
|
provisionRequestMsg.getProvisionDeviceCredentialsMsg().getProvisionDeviceSecret() != null && |
|
|
|
|
|
provisionRequestMsg.getDeviceName() != null) |
|
|
|
|
|
throw new RuntimeException("Wrong credentials!"); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
private void processRegularSessionMsg(ChannelHandlerContext ctx, MqttMessage msg) { |
|
|
private void processRegularSessionMsg(ChannelHandlerContext ctx, MqttMessage msg) { |
|
|
switch (msg.fixedHeader().messageType()) { |
|
|
switch (msg.fixedHeader().messageType()) { |
|
|
case PUBLISH: |
|
|
case PUBLISH: |
|
|
@ -335,9 +342,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
} else { |
|
|
} else { |
|
|
deviceSessionCtx.getContext().getProtoMqttAdaptor().convertToPublish(deviceSessionCtx, provisionResponseMsg).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); |
|
|
deviceSessionCtx.getContext().getProtoMqttAdaptor().convertToPublish(deviceSessionCtx, provisionResponseMsg).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); |
|
|
} |
|
|
} |
|
|
|
|
|
transportService.getSchedulerExecutor().schedule(() -> processDisconnect(ctx), 60, TimeUnit.SECONDS); |
|
|
//TODO: close session with some delay.
|
|
|
|
|
|
//transportService.getScheduler().submit task with 60 seconds delay to close the session.
|
|
|
|
|
|
} catch (Exception e) { |
|
|
} catch (Exception e) { |
|
|
log.trace("[{}] Failed to convert device attributes response to MQTT msg", sessionId, e); |
|
|
log.trace("[{}] Failed to convert device attributes response to MQTT msg", sessionId, e); |
|
|
} |
|
|
} |
|
|
@ -379,7 +384,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
case MqttTopics.GATEWAY_ATTRIBUTES_TOPIC: |
|
|
case MqttTopics.GATEWAY_ATTRIBUTES_TOPIC: |
|
|
case MqttTopics.GATEWAY_RPC_TOPIC: |
|
|
case MqttTopics.GATEWAY_RPC_TOPIC: |
|
|
case MqttTopics.GATEWAY_ATTRIBUTES_RESPONSE_TOPIC: |
|
|
case MqttTopics.GATEWAY_ATTRIBUTES_RESPONSE_TOPIC: |
|
|
case MqttTopics.GATEWAY_PROVISION_RESPONSE_TOPIC: |
|
|
|
|
|
case MqttTopics.DEVICE_PROVISION_RESPONSE_TOPIC: |
|
|
case MqttTopics.DEVICE_PROVISION_RESPONSE_TOPIC: |
|
|
registerSubQoS(topic, grantedQoSList, reqQoS); |
|
|
registerSubQoS(topic, grantedQoSList, reqQoS); |
|
|
break; |
|
|
break; |
|
|
@ -447,7 +451,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
private void processConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) { |
|
|
private void processConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) { |
|
|
log.info("[{}] Processing connect msg for client: {}!", sessionId, msg.payload().clientIdentifier()); |
|
|
log.info("[{}] Processing connect msg for client: {}!", sessionId, msg.payload().clientIdentifier()); |
|
|
String userName = msg.payload().userName(); |
|
|
String userName = msg.payload().userName(); |
|
|
if (DataConstants.PROVISION.equals(userName)) { |
|
|
String clientId = msg.payload().clientIdentifier(); |
|
|
|
|
|
if (DataConstants.PROVISION.equals(userName) || DataConstants.PROVISION.equals(clientId)) { |
|
|
deviceSessionCtx.setProvisionOnly(true); |
|
|
deviceSessionCtx.setProvisionOnly(true); |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); |
|
|
} else { |
|
|
} else { |
|
|
|