|
|
|
@ -46,6 +46,7 @@ import org.thingsboard.server.common.data.DeviceProfile; |
|
|
|
import org.thingsboard.server.common.data.DeviceTransportType; |
|
|
|
import org.thingsboard.server.common.data.TransportPayloadType; |
|
|
|
import org.thingsboard.server.common.data.device.profile.MqttTopics; |
|
|
|
import org.thingsboard.server.common.data.id.FirmwareId; |
|
|
|
import org.thingsboard.server.common.msg.EncryptionUtil; |
|
|
|
import org.thingsboard.server.common.msg.tools.TbRateLimitsException; |
|
|
|
import org.thingsboard.server.common.transport.SessionMsgListener; |
|
|
|
@ -69,10 +70,10 @@ import org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher; |
|
|
|
import org.thingsboard.server.transport.mqtt.util.SslUtil; |
|
|
|
|
|
|
|
import javax.net.ssl.SSLPeerUnverifiedException; |
|
|
|
import java.security.cert.Certificate; |
|
|
|
import java.security.cert.X509Certificate; |
|
|
|
import java.io.IOException; |
|
|
|
import java.net.InetSocketAddress; |
|
|
|
import java.security.cert.Certificate; |
|
|
|
import java.security.cert.X509Certificate; |
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.List; |
|
|
|
import java.util.Optional; |
|
|
|
@ -80,7 +81,10 @@ import java.util.UUID; |
|
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
|
import java.util.concurrent.ConcurrentMap; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.regex.Matcher; |
|
|
|
import java.util.regex.Pattern; |
|
|
|
|
|
|
|
import static com.amazonaws.util.StringUtils.UTF8; |
|
|
|
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.MqttMessageType.CONNACK; |
|
|
|
@ -99,6 +103,10 @@ import static io.netty.handler.codec.mqtt.MqttQoS.FAILURE; |
|
|
|
@Slf4j |
|
|
|
public class MqttTransportHandler extends ChannelInboundHandlerAdapter implements GenericFutureListener<Future<? super Void>>, SessionMsgListener { |
|
|
|
|
|
|
|
private static final Pattern FW_PATTERN = Pattern.compile("v2/fw/request/(?<requestId>\\d+)/chunk/(?<chunk>\\d+)"); |
|
|
|
|
|
|
|
private static final String PAYLOAD_TOO_LARGE = "PAYLOAD_TOO_LARGE"; |
|
|
|
|
|
|
|
private static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; |
|
|
|
|
|
|
|
private final UUID sessionId; |
|
|
|
@ -112,6 +120,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
|
private volatile InetSocketAddress address; |
|
|
|
private volatile GatewaySessionHandler gatewaySessionHandler; |
|
|
|
|
|
|
|
private final ConcurrentHashMap<String, String> fwSessions; |
|
|
|
private final ConcurrentHashMap<String, Integer> fwChunkSizes; |
|
|
|
|
|
|
|
MqttTransportHandler(MqttTransportContext context, SslHandler sslHandler) { |
|
|
|
this.sessionId = UUID.randomUUID(); |
|
|
|
this.context = context; |
|
|
|
@ -120,6 +131,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
|
this.sslHandler = sslHandler; |
|
|
|
this.mqttQoSMap = new ConcurrentHashMap<>(); |
|
|
|
this.deviceSessionCtx = new DeviceSessionCtx(sessionId, mqttQoSMap, context); |
|
|
|
this.fwSessions = new ConcurrentHashMap<>(); |
|
|
|
this.fwChunkSizes = new ConcurrentHashMap<>(); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -280,6 +293,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
|
|
|
|
|
private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) { |
|
|
|
try { |
|
|
|
Matcher fwMatcher; |
|
|
|
MqttTransportAdaptor payloadAdaptor = deviceSessionCtx.getPayloadAdaptor(); |
|
|
|
if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) { |
|
|
|
TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); |
|
|
|
@ -299,6 +313,38 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
|
} else if (topicName.equals(MqttTopics.DEVICE_CLAIM_TOPIC)) { |
|
|
|
TransportProtos.ClaimDeviceMsg claimDeviceMsg = payloadAdaptor.convertToClaimDevice(deviceSessionCtx, mqttMsg); |
|
|
|
transportService.process(deviceSessionCtx.getSessionInfo(), claimDeviceMsg, getPubAckCallback(ctx, msgId, claimDeviceMsg)); |
|
|
|
} else if ((fwMatcher = FW_PATTERN.matcher(topicName)).find()) { |
|
|
|
String payload = mqttMsg.content().toString(UTF8); |
|
|
|
int chunkSize = payload != null ? Integer.parseInt(payload) : 0; |
|
|
|
String requestId = fwMatcher.group("requestId"); |
|
|
|
int chunk = Integer.parseInt(fwMatcher.group("chunk")); |
|
|
|
|
|
|
|
if (chunkSize > 0) { |
|
|
|
this.fwChunkSizes.put(requestId, chunkSize); |
|
|
|
} else { |
|
|
|
chunkSize = fwChunkSizes.getOrDefault(requestId, 0); |
|
|
|
} |
|
|
|
|
|
|
|
if (chunkSize > context.getMaxPayloadSize()) { |
|
|
|
sendFirmwareError(ctx, PAYLOAD_TOO_LARGE); |
|
|
|
return; |
|
|
|
} |
|
|
|
|
|
|
|
String firmwareId = fwSessions.get(requestId); |
|
|
|
|
|
|
|
if (firmwareId != null) { |
|
|
|
sendFirmware(ctx, mqttMsg.variableHeader().packetId(), firmwareId, requestId, chunkSize, chunk); |
|
|
|
} else { |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = deviceSessionCtx.getSessionInfo(); |
|
|
|
TransportProtos.GetFirmwareRequestMsg getFirmwareRequestMsg = TransportProtos.GetFirmwareRequestMsg.newBuilder() |
|
|
|
.setDeviceIdMSB(sessionInfo.getDeviceIdMSB()) |
|
|
|
.setDeviceIdLSB(sessionInfo.getDeviceIdLSB()) |
|
|
|
.setTenantIdMSB(sessionInfo.getTenantIdMSB()) |
|
|
|
.setTenantIdLSB(sessionInfo.getTenantIdLSB()) |
|
|
|
.build(); |
|
|
|
transportService.process(deviceSessionCtx.getSessionInfo(), getFirmwareRequestMsg, |
|
|
|
new FirmwareCallback(ctx, mqttMsg.variableHeader().packetId(), getFirmwareRequestMsg, requestId, chunkSize, chunk)); |
|
|
|
} |
|
|
|
} else { |
|
|
|
transportService.reportActivity(deviceSessionCtx.getSessionInfo()); |
|
|
|
ack(ctx, msgId); |
|
|
|
@ -366,6 +412,65 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private class FirmwareCallback implements TransportServiceCallback<TransportProtos.GetFirmwareResponseMsg> { |
|
|
|
private final ChannelHandlerContext ctx; |
|
|
|
private final int msgId; |
|
|
|
private final TransportProtos.GetFirmwareRequestMsg msg; |
|
|
|
private final String requestId; |
|
|
|
private final int chunkSize; |
|
|
|
private final int chunk; |
|
|
|
|
|
|
|
FirmwareCallback(ChannelHandlerContext ctx, int msgId, TransportProtos.GetFirmwareRequestMsg msg, String requestId, int chunkSize, int chunk) { |
|
|
|
this.ctx = ctx; |
|
|
|
this.msgId = msgId; |
|
|
|
this.msg = msg; |
|
|
|
this.requestId = requestId; |
|
|
|
this.chunkSize = chunkSize; |
|
|
|
this.chunk = chunk; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onSuccess(TransportProtos.GetFirmwareResponseMsg response) { |
|
|
|
if (TransportProtos.ResponseStatus.SUCCESS.equals(response.getResponseStatus())) { |
|
|
|
FirmwareId firmwareId = new FirmwareId(new UUID(response.getFirmwareIdMSB(), response.getFirmwareIdLSB())); |
|
|
|
fwSessions.put(requestId, firmwareId.toString()); |
|
|
|
sendFirmware(ctx, msgId, firmwareId.toString(), requestId, chunkSize, chunk); |
|
|
|
} else { |
|
|
|
sendFirmwareError(ctx, response.getResponseStatus().toString()); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onError(Throwable e) { |
|
|
|
log.trace("[{}] Failed to get firmware: {}", sessionId, msg, e); |
|
|
|
processDisconnect(ctx); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void sendFirmware(ChannelHandlerContext ctx, int msgId, String firmwareId, String requestId, int chunkSize, int chunk) { |
|
|
|
log.trace("[{}] Send firmware [{}] to device!", sessionId, firmwareId); |
|
|
|
ack(ctx, msgId); |
|
|
|
try { |
|
|
|
byte[] firmwareChunk = context.getFirmwareCacheReader().get(firmwareId, chunkSize, chunk); |
|
|
|
deviceSessionCtx.getPayloadAdaptor() |
|
|
|
.convertToPublish(deviceSessionCtx, firmwareChunk, requestId, chunk) |
|
|
|
.ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); |
|
|
|
if (firmwareChunk != null && chunkSize != firmwareChunk.length) { |
|
|
|
scheduler.schedule(() -> processDisconnect(ctx), 60, TimeUnit.SECONDS); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("[{}] Failed to send firmware response!", sessionId, e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void sendFirmwareError(ChannelHandlerContext ctx, String error) { |
|
|
|
log.warn("[{}] {}", sessionId, error); |
|
|
|
deviceSessionCtx.getChannel().writeAndFlush(deviceSessionCtx |
|
|
|
.getPayloadAdaptor() |
|
|
|
.createMqttPublishMsg(deviceSessionCtx, MqttTopics.DEVICE_FIRMWARE_ERROR_TOPIC, error.getBytes())); |
|
|
|
processDisconnect(ctx); |
|
|
|
} |
|
|
|
|
|
|
|
private void processSubscribe(ChannelHandlerContext ctx, MqttSubscribeMessage mqttMsg) { |
|
|
|
if (!checkConnected(ctx, mqttMsg)) { |
|
|
|
return; |
|
|
|
@ -396,6 +501,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement |
|
|
|
case MqttTopics.GATEWAY_RPC_TOPIC: |
|
|
|
case MqttTopics.GATEWAY_ATTRIBUTES_RESPONSE_TOPIC: |
|
|
|
case MqttTopics.DEVICE_PROVISION_RESPONSE_TOPIC: |
|
|
|
case MqttTopics.DEVICE_FIRMWARE_RESPONSES_TOPIC: |
|
|
|
case MqttTopics.DEVICE_FIRMWARE_ERROR_TOPIC: |
|
|
|
registerSubQoS(topic, grantedQoSList, reqQoS); |
|
|
|
break; |
|
|
|
default: |
|
|
|
|