|
|
|
@ -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<Future<? super Void>> { |
|
|
|
|
|
|
|
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<Integer> 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<Integer> 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<? super Void> future) throws Exception { |
|
|
|
processor.process(SessionCloseMsg.onError(sessionCtx.getSessionId())); |
|
|
|
processor.process(SessionCloseMsg.onError(deviceSessionCtx.getSessionId())); |
|
|
|
} |
|
|
|
} |
|
|
|
|