Browse Source

Address code review: fix naming, visibility, and lazy address resolution

- Rename TagsKey → tagsKey and make it private (Java naming convention)
- Make isSSL() private in MqttTransportHandler (internal use only)
- Fix double space in if (isSSL) in MqttTransportContext
- Extract getClientAddr() helper and move clientAddr computation inside
  logging guards so address resolution is skipped when logging is disabled

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
pull/15112/head
Sergey Matvienko 7 months ago
committed by Sergii Matviienko
parent
commit
a643a340d8
  1. 4
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java
  2. 24
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  3. 4
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

4
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) { public void channelRegistered(boolean isSSL) {
if (isSSL) { if (isSSL) {
connectionsActiveCounterMQTTS.incrementAndGet(); connectionsActiveCounterMQTTS.incrementAndGet();
} else { } else {
connectionsActiveCounterMQTT.incrementAndGet(); connectionsActiveCounterMQTT.incrementAndGet();
@ -107,7 +107,7 @@ public class MqttTransportContext extends TransportContext {
} }
public void channelUnregistered(boolean isSSL) { public void channelUnregistered(boolean isSSL) {
if (isSSL) { if (isSSL) {
connectionsActiveCounterMQTTS.decrementAndGet(); connectionsActiveCounterMQTTS.decrementAndGet();
} else { } else {
connectionsActiveCounterMQTT.decrementAndGet(); connectionsActiveCounterMQTT.decrementAndGet();

24
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<>(); this.rpcAwaitingAck = new ConcurrentHashMap<>();
} }
boolean isSSL() { private boolean isSSL() {
return sslHandler != null; 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) { InetSocketAddress getAddress(ChannelHandlerContext ctx) {
var address = ctx.channel().attr(MqttTransportService.ADDRESS).get(); var address = ctx.channel().attr(MqttTransportService.ADDRESS).get();
if (address == null) { if (address == null) {
@ -1152,15 +1163,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
@Override @Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { 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 (cause instanceof IOException) {
if (log.isDebugEnabled()) { if (log.isDebugEnabled()) {
String clientAddr = getClientAddr(ctx);
log.debug("[{}][{}][{}][{}] {}: {}", sessionId, log.debug("[{}][{}][{}][{}] {}: {}", sessionId,
Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null),
Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""),
@ -1169,6 +1174,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
cause.getMessage(), cause.getMessage(),
cause); cause);
} else if (log.isInfoEnabled()) { } else if (log.isInfoEnabled()) {
String clientAddr = getClientAddr(ctx);
log.info("[{}][{}][{}][{}] {}: {}", sessionId, log.info("[{}][{}][{}][{}] {}: {}", sessionId,
Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null),
Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""),
@ -1177,7 +1183,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
cause.getMessage()); cause.getMessage());
} }
} else { } else {
log.error("[{}][{}] Unexpected Exception", sessionId, clientAddr, cause); log.error("[{}][{}] Unexpected Exception", sessionId, getClientAddr(ctx), cause);
} }
closeCtx(ctx, MqttReasonCodes.Disconnect.SERVER_SHUTTING_DOWN); closeCtx(ctx, MqttReasonCodes.Disconnect.SERVER_SHUTTING_DOWN);

4
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) { public void createGaugeStats(String statsName, AtomicInteger number, String... tags) {
String key = "thingsboard" + "." + StatsType.TRANSPORT.getName() + "." + statsName; String key = "thingsboard" + "." + StatsType.TRANSPORT.getName() + "." + statsName;
statsFactory.createGauge(key, number, tags); 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 ""; if (tags == null || tags.length < 2) return "";
StringBuilder sb = new StringBuilder("["); StringBuilder sb = new StringBuilder("[");
for (int i = 0; i < tags.length; i += 2) { for (int i = 0; i < tags.length; i += 2) {

Loading…
Cancel
Save