Browse Source

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.
pull/15451/head
Oleksandra Matviienko 5 months ago
parent
commit
5b4d427008
  1. 23
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java

23
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())); ResourceLeakDetector.setLevel(ResourceLeakDetector.Level.valueOf(leakDetectorLevel.toUpperCase()));
log.info("Starting MQTT transport..."); log.info("Starting MQTT transport...");
bossGroup = new NioEventLoopGroup(bossGroupThreadCount);
workerGroup = new NioEventLoopGroup(workerGroupThreadCount);
try { try {
bossGroup = new NioEventLoopGroup(bossGroupThreadCount);
workerGroup = new NioEventLoopGroup(workerGroupThreadCount);
ServerBootstrap b = new ServerBootstrap(); ServerBootstrap b = new ServerBootstrap();
b.group(bossGroup, workerGroup) b.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class) .channel(NioServerSocketChannel.class)
@ -96,7 +96,7 @@ public class MqttTransportService implements TbTransportService {
.childOption(ChannelOption.SO_KEEPALIVE, keepAlive); .childOption(ChannelOption.SO_KEEPALIVE, keepAlive);
sslServerChannel = b.bind(sslHost, sslPort).sync().channel(); sslServerChannel = b.bind(sslHost, sslPort).sync().channel();
} }
} catch (Throwable e) { } catch (Exception e) {
log.error("Failed to start MQTT transport, releasing resources", e); log.error("Failed to start MQTT transport, releasing resources", e);
if (e instanceof InterruptedException) { if (e instanceof InterruptedException) {
Thread.currentThread().interrupt(); Thread.currentThread().interrupt();
@ -108,16 +108,17 @@ public class MqttTransportService implements TbTransportService {
if (sslServerChannel != null) { if (sslServerChannel != null) {
sslServerChannel.close().sync(); sslServerChannel.close().sync();
} }
} catch (InterruptedException ie) { } catch (Exception suppressed) {
Thread.currentThread().interrupt(); e.addSuppressed(suppressed);
} finally { } finally {
workerGroup.shutdownGracefully(); if (workerGroup != null) {
bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully();
} }
if (e instanceof Exception) { if (bossGroup != null) {
throw (Exception) e; bossGroup.shutdownGracefully();
}
} }
throw (Error) e; throw e;
} }
log.info("Mqtt transport started!"); log.info("Mqtt transport started!");
} }

Loading…
Cancel
Save