diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index bf474cc266..fe39721f27 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -76,6 +76,7 @@ import org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher; import javax.net.ssl.SSLPeerUnverifiedException; import java.io.IOException; import java.net.InetSocketAddress; +import java.nio.ByteBuffer; import java.security.cert.Certificate; import java.security.cert.X509Certificate; import java.util.ArrayList; @@ -111,7 +112,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private static final Pattern FW_REQUEST_PATTERN = Pattern.compile(MqttTopics.DEVICE_FIRMWARE_REQUEST_TOPIC_PATTERN); private static final Pattern SW_REQUEST_PATTERN = Pattern.compile(MqttTopics.DEVICE_SOFTWARE_REQUEST_TOPIC_PATTERN); - private static final String PAYLOAD_TOO_LARGE = "PAYLOAD_TOO_LARGE"; private static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; @@ -124,7 +124,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private final ConcurrentMap mqttQoSMap; final DeviceSessionCtx deviceSessionCtx; - volatile InetSocketAddress address; + volatile int ip = 0; + volatile int port = 0; volatile GatewaySessionHandler gatewaySessionHandler; private final ConcurrentHashMap otaPackSessions; @@ -183,9 +184,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } void processMqttMsg(ChannelHandlerContext ctx, MqttMessage msg) { - address = getAddress(ctx); + if (port == 0) { + InetSocketAddress address = getAddress(ctx); + ip = getIpv4(address); //ipv6 will not appear in logs + port = address.getPort(); + } + if (msg.fixedHeader() == null) { - log.info("[{}:{}] Invalid message received", address.getHostName(), address.getPort()); + log.info("[{}:{}] Invalid message received", ip, port); ctx.close(); return; } @@ -199,6 +205,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } + int getIpv4(InetSocketAddress address) { + byte[] ipBytes = address.getAddress().getAddress(); + return ipBytes.length == 4 ? ByteBuffer.wrap(ipBytes).getInt() : -1; + } + InetSocketAddress getAddress(ChannelHandlerContext ctx) { return (InetSocketAddress) ctx.channel().remoteAddress(); } @@ -791,7 +802,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onError(Throwable e) { - log.trace("[{}] Failed to process credentials: {}", address, userName, e); + log.trace("[{}] Failed to process credentials: {}", ip, userName, e); ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage)); ctx.close(); } @@ -814,14 +825,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onError(Throwable e) { - log.trace("[{}] Failed to process credentials: {}", address, sha3Hash, e); + log.trace("[{}] Failed to process credentials: {}", ip, sha3Hash, e); ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage)); ctx.close(); } }); } catch (Exception e) { ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, connectMessage)); - log.trace("[{}] X509 auth failure: {}", sessionId, address, e); + log.trace("[{}] X509 auth failure: {}", sessionId, ip, e); ctx.close(); } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index 9f0c89714a..76e2911793 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -238,7 +238,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { public void release() { if (!msgQueue.isEmpty()) { - log.warn("doDisconnect for device {} but unprocessed messages {} left in the msg queue", getDeviceId(), msgQueue.size()); + log.warn("doDisconnect for device {} but unprocessed messages {} left in the msg queue. cleared", getDeviceId(), getMsgQueueSize()); msgQueue.forEach(ReferenceCountUtil::safeRelease); msgQueue.clear(); }