diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java index 8a60168154..10245c5a24 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java @@ -99,7 +99,7 @@ public class MqttTransportContext extends TransportContext { } public void channelRegistered(boolean isSSL) { - if (isSSL) { + if (isSSL) { connectionsActiveCounterMQTTS.incrementAndGet(); } else { connectionsActiveCounterMQTT.incrementAndGet(); @@ -107,7 +107,7 @@ public class MqttTransportContext extends TransportContext { } public void channelUnregistered(boolean isSSL) { - if (isSSL) { + if (isSSL) { connectionsActiveCounterMQTTS.decrementAndGet(); } else { connectionsActiveCounterMQTT.decrementAndGet(); 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 aadeff0dee..d427c9b78e 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 @@ -180,7 +180,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement this.rpcAwaitingAck = new ConcurrentHashMap<>(); } - boolean isSSL() { + private boolean isSSL() { return sslHandler != null; } @@ -258,6 +258,17 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } + private String getClientAddr(ChannelHandlerContext ctx) { + try { + InetSocketAddress remote = getAddress(ctx); + if (remote == null) return "unknown"; + String host = remote.getAddress() != null ? remote.getAddress().getHostAddress() : remote.getHostString(); + return host + ":" + remote.getPort(); + } catch (Exception ignored) { + return "unknown"; + } + } + InetSocketAddress getAddress(ChannelHandlerContext ctx) { var address = ctx.channel().attr(MqttTransportService.ADDRESS).get(); if (address == null) { @@ -1152,15 +1163,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { - String clientAddr = null; - try { - InetSocketAddress remote = getAddress(ctx); - clientAddr = remote != null ? (remote.getAddress() != null ? remote.getAddress().getHostAddress() : remote.getHostString()) + ":" + remote.getPort() : "IPunknown"; - } catch (Exception ignored) { - } - if (cause instanceof IOException) { if (log.isDebugEnabled()) { + String clientAddr = getClientAddr(ctx); log.debug("[{}][{}][{}][{}] {}: {}", sessionId, Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), @@ -1169,6 +1174,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement cause.getMessage(), cause); } else if (log.isInfoEnabled()) { + String clientAddr = getClientAddr(ctx); log.info("[{}][{}][{}][{}] {}: {}", sessionId, Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), @@ -1177,7 +1183,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement cause.getMessage()); } } else { - log.error("[{}][{}] Unexpected Exception", sessionId, clientAddr, cause); + log.error("[{}][{}] Unexpected Exception", sessionId, getClientAddr(ctx), cause); } closeCtx(ctx, MqttReasonCodes.Disconnect.SERVER_SHUTTING_DOWN); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 47e341547e..e8ade498db 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -1260,10 +1260,10 @@ public class DefaultTransportService extends TransportActivityManager implements public void createGaugeStats(String statsName, AtomicInteger number, String... tags) { String key = "thingsboard" + "." + StatsType.TRANSPORT.getName() + "." + statsName; statsFactory.createGauge(key, number, tags); - statsMap.put(statsName + TagsKey(tags), number); + statsMap.put(statsName + tagsKey(tags), number); } - String TagsKey(String... tags) { + private String tagsKey(String... tags) { if (tags == null || tags.length < 2) return ""; StringBuilder sb = new StringBuilder("["); for (int i = 0; i < tags.length; i += 2) {