|
|
@ -37,6 +37,7 @@ import org.thingsboard.server.common.transport.adaptor.AdaptorException; |
|
|
import org.thingsboard.server.common.transport.auth.DeviceAuthService; |
|
|
import org.thingsboard.server.common.transport.auth.DeviceAuthService; |
|
|
import org.thingsboard.server.common.transport.quota.QuotaService; |
|
|
import org.thingsboard.server.common.transport.quota.QuotaService; |
|
|
import org.thingsboard.server.dao.EncryptionUtil; |
|
|
import org.thingsboard.server.dao.EncryptionUtil; |
|
|
|
|
|
import org.thingsboard.server.dao.device.DeviceOfflineService; |
|
|
import org.thingsboard.server.dao.device.DeviceService; |
|
|
import org.thingsboard.server.dao.device.DeviceService; |
|
|
import org.thingsboard.server.dao.relation.RelationService; |
|
|
import org.thingsboard.server.dao.relation.RelationService; |
|
|
import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; |
|
|
import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; |
|
|
@ -72,13 +73,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
private final DeviceAuthService authService; |
|
|
private final DeviceAuthService authService; |
|
|
private final RelationService relationService; |
|
|
private final RelationService relationService; |
|
|
private final QuotaService quotaService; |
|
|
private final QuotaService quotaService; |
|
|
|
|
|
private final DeviceOfflineService offlineService; |
|
|
private final SslHandler sslHandler; |
|
|
private final SslHandler sslHandler; |
|
|
private volatile boolean connected; |
|
|
private volatile boolean connected; |
|
|
private volatile InetSocketAddress address; |
|
|
private volatile InetSocketAddress address; |
|
|
private volatile GatewaySessionCtx gatewaySessionCtx; |
|
|
private volatile GatewaySessionCtx gatewaySessionCtx; |
|
|
|
|
|
|
|
|
public MqttTransportHandler(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, RelationService relationService, |
|
|
public MqttTransportHandler(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, RelationService relationService, |
|
|
MqttTransportAdaptor adaptor, SslHandler sslHandler, QuotaService quotaService) { |
|
|
MqttTransportAdaptor adaptor, SslHandler sslHandler, QuotaService quotaService, DeviceOfflineService offlineService) { |
|
|
this.processor = processor; |
|
|
this.processor = processor; |
|
|
this.deviceService = deviceService; |
|
|
this.deviceService = deviceService; |
|
|
this.relationService = relationService; |
|
|
this.relationService = relationService; |
|
|
@ -88,6 +90,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
this.sessionId = deviceSessionCtx.getSessionId().toUidStr(); |
|
|
this.sessionId = deviceSessionCtx.getSessionId().toUidStr(); |
|
|
this.sslHandler = sslHandler; |
|
|
this.sslHandler = sslHandler; |
|
|
this.quotaService = quotaService; |
|
|
this.quotaService = quotaService; |
|
|
|
|
|
this.offlineService = offlineService; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
@ -129,11 +132,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
case PINGREQ: |
|
|
case PINGREQ: |
|
|
if (checkConnected(ctx)) { |
|
|
if (checkConnected(ctx)) { |
|
|
ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); |
|
|
ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), false); |
|
|
} |
|
|
} |
|
|
break; |
|
|
break; |
|
|
case DISCONNECT: |
|
|
case DISCONNECT: |
|
|
if (checkConnected(ctx)) { |
|
|
if (checkConnected(ctx)) { |
|
|
processDisconnect(ctx); |
|
|
processDisconnect(ctx); |
|
|
|
|
|
offlineService.offline(deviceSessionCtx.getDevice()); |
|
|
} |
|
|
} |
|
|
break; |
|
|
break; |
|
|
default: |
|
|
default: |
|
|
@ -185,23 +190,28 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
try { |
|
|
try { |
|
|
if (topicName.equals(DEVICE_TELEMETRY_TOPIC)) { |
|
|
if (topicName.equals(DEVICE_TELEMETRY_TOPIC)) { |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, POST_TELEMETRY_REQUEST, mqttMsg); |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, POST_TELEMETRY_REQUEST, mqttMsg); |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), true); |
|
|
} else if (topicName.equals(DEVICE_ATTRIBUTES_TOPIC)) { |
|
|
} else if (topicName.equals(DEVICE_ATTRIBUTES_TOPIC)) { |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, POST_ATTRIBUTES_REQUEST, mqttMsg); |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, POST_ATTRIBUTES_REQUEST, mqttMsg); |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), true); |
|
|
} else if (topicName.startsWith(DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { |
|
|
} else if (topicName.startsWith(DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, GET_ATTRIBUTES_REQUEST, mqttMsg); |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, GET_ATTRIBUTES_REQUEST, mqttMsg); |
|
|
if (msgId >= 0) { |
|
|
if (msgId >= 0) { |
|
|
ctx.writeAndFlush(createMqttPubAckMsg(msgId)); |
|
|
ctx.writeAndFlush(createMqttPubAckMsg(msgId)); |
|
|
} |
|
|
} |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), false); |
|
|
} else if (topicName.startsWith(DEVICE_RPC_RESPONSE_TOPIC)) { |
|
|
} else if (topicName.startsWith(DEVICE_RPC_RESPONSE_TOPIC)) { |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, TO_DEVICE_RPC_RESPONSE, mqttMsg); |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, TO_DEVICE_RPC_RESPONSE, mqttMsg); |
|
|
if (msgId >= 0) { |
|
|
if (msgId >= 0) { |
|
|
ctx.writeAndFlush(createMqttPubAckMsg(msgId)); |
|
|
ctx.writeAndFlush(createMqttPubAckMsg(msgId)); |
|
|
} |
|
|
} |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), true); |
|
|
} else if (topicName.startsWith(DEVICE_RPC_REQUESTS_TOPIC)) { |
|
|
} else if (topicName.startsWith(DEVICE_RPC_REQUESTS_TOPIC)) { |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, TO_SERVER_RPC_REQUEST, mqttMsg); |
|
|
msg = adaptor.convertToActorMsg(deviceSessionCtx, TO_SERVER_RPC_REQUEST, mqttMsg); |
|
|
if (msgId >= 0) { |
|
|
if (msgId >= 0) { |
|
|
ctx.writeAndFlush(createMqttPubAckMsg(msgId)); |
|
|
ctx.writeAndFlush(createMqttPubAckMsg(msgId)); |
|
|
} |
|
|
} |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), true); |
|
|
} |
|
|
} |
|
|
} catch (AdaptorException e) { |
|
|
} catch (AdaptorException e) { |
|
|
log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); |
|
|
log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e); |
|
|
@ -250,6 +260,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList)); |
|
|
ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList)); |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), false); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void processUnsubscribe(ChannelHandlerContext ctx, MqttUnsubscribeMessage mqttMsg) { |
|
|
private void processUnsubscribe(ChannelHandlerContext ctx, MqttUnsubscribeMessage mqttMsg) { |
|
|
@ -273,6 +284,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId())); |
|
|
ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId())); |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), false); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private MqttMessage createUnSubAckMessage(int msgId) { |
|
|
private MqttMessage createUnSubAckMessage(int msgId) { |
|
|
@ -304,6 +316,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); |
|
|
connected = true; |
|
|
connected = true; |
|
|
checkGatewaySession(); |
|
|
checkGatewaySession(); |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), false); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -315,6 +328,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); |
|
|
connected = true; |
|
|
connected = true; |
|
|
checkGatewaySession(); |
|
|
checkGatewaySession(); |
|
|
|
|
|
offlineService.online(deviceSessionCtx.getDevice(), false); |
|
|
} else { |
|
|
} else { |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); |
|
|
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); |
|
|
ctx.close(); |
|
|
ctx.close(); |
|
|
@ -365,6 +379,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { |
|
|
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { |
|
|
log.error("[{}] Unexpected Exception", sessionId, cause); |
|
|
log.error("[{}] Unexpected Exception", sessionId, cause); |
|
|
ctx.close(); |
|
|
ctx.close(); |
|
|
|
|
|
if(deviceSessionCtx.getDevice() != null) { |
|
|
|
|
|
offlineService.offline(deviceSessionCtx.getDevice()); |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private static MqttSubAckMessage createSubAckMessage(Integer msgId, List<Integer> grantedQoSList) { |
|
|
private static MqttSubAckMessage createSubAckMessage(Integer msgId, List<Integer> grantedQoSList) { |
|
|
@ -403,7 +420,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
if (infoNode != null) { |
|
|
if (infoNode != null) { |
|
|
JsonNode gatewayNode = infoNode.get("gateway"); |
|
|
JsonNode gatewayNode = infoNode.get("gateway"); |
|
|
if (gatewayNode != null && gatewayNode.asBoolean()) { |
|
|
if (gatewayNode != null && gatewayNode.asBoolean()) { |
|
|
gatewaySessionCtx = new GatewaySessionCtx(processor, deviceService, authService, relationService, deviceSessionCtx); |
|
|
gatewaySessionCtx = new GatewaySessionCtx(processor, deviceService, authService, |
|
|
|
|
|
relationService, deviceSessionCtx, offlineService); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
@ -411,5 +429,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
@Override |
|
|
@Override |
|
|
public void operationComplete(Future<? super Void> future) throws Exception { |
|
|
public void operationComplete(Future<? super Void> future) throws Exception { |
|
|
processor.process(SessionCloseMsg.onError(deviceSessionCtx.getSessionId())); |
|
|
processor.process(SessionCloseMsg.onError(deviceSessionCtx.getSessionId())); |
|
|
|
|
|
if(deviceSessionCtx.getDevice() != null) { |
|
|
|
|
|
offlineService.offline(deviceSessionCtx.getDevice()); |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|