From 5b4d4270084d8d4f8abe75c180bfd3a0c626f3ad Mon Sep 17 00:00:00 2001 From: Oleksandra Matviienko Date: Thu, 30 Apr 2026 13:36:15 +0200 Subject: [PATCH] Hardened MQTT transport init failure handling Moved NioEventLoopGroup allocations into the try block so that a constructor failure for the second group no longer leaks the first. Channel close failures during cleanup now attach via addSuppressed instead of replacing the original BindException. Narrowed the outer catch from Throwable to Exception, removing the brittle (Error) cast that would have masked any direct Throwable subclass. --- .../transport/mqtt/MqttTransportService.java | 23 ++++++++++--------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java index 5e52863791..4e41fa7eca 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java @@ -78,9 +78,9 @@ public class MqttTransportService implements TbTransportService { ResourceLeakDetector.setLevel(ResourceLeakDetector.Level.valueOf(leakDetectorLevel.toUpperCase())); log.info("Starting MQTT transport..."); - bossGroup = new NioEventLoopGroup(bossGroupThreadCount); - workerGroup = new NioEventLoopGroup(workerGroupThreadCount); try { + bossGroup = new NioEventLoopGroup(bossGroupThreadCount); + workerGroup = new NioEventLoopGroup(workerGroupThreadCount); ServerBootstrap b = new ServerBootstrap(); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) @@ -96,7 +96,7 @@ public class MqttTransportService implements TbTransportService { .childOption(ChannelOption.SO_KEEPALIVE, keepAlive); sslServerChannel = b.bind(sslHost, sslPort).sync().channel(); } - } catch (Throwable e) { + } catch (Exception e) { log.error("Failed to start MQTT transport, releasing resources", e); if (e instanceof InterruptedException) { Thread.currentThread().interrupt(); @@ -108,16 +108,17 @@ public class MqttTransportService implements TbTransportService { if (sslServerChannel != null) { sslServerChannel.close().sync(); } - } catch (InterruptedException ie) { - Thread.currentThread().interrupt(); + } catch (Exception suppressed) { + e.addSuppressed(suppressed); } finally { - workerGroup.shutdownGracefully(); - bossGroup.shutdownGracefully(); - } - if (e instanceof Exception) { - throw (Exception) e; + if (workerGroup != null) { + workerGroup.shutdownGracefully(); + } + if (bossGroup != null) { + bossGroup.shutdownGracefully(); + } } - throw (Error) e; + throw e; } log.info("Mqtt transport started!"); }