|
|
@ -46,6 +46,14 @@ public class MqttTransportService implements TbTransportService { |
|
|
@Value("${transport.mqtt.bind_port}") |
|
|
@Value("${transport.mqtt.bind_port}") |
|
|
private Integer port; |
|
|
private Integer port; |
|
|
|
|
|
|
|
|
|
|
|
@Value("${transport.mqtt.ssl.enabled}") |
|
|
|
|
|
private boolean sslEnabled; |
|
|
|
|
|
|
|
|
|
|
|
@Value("${transport.mqtt.ssl.bind_address}") |
|
|
|
|
|
private String sslHost; |
|
|
|
|
|
@Value("${transport.mqtt.ssl.bind_port}") |
|
|
|
|
|
private Integer sslPort; |
|
|
|
|
|
|
|
|
@Value("${transport.mqtt.netty.leak_detector_level}") |
|
|
@Value("${transport.mqtt.netty.leak_detector_level}") |
|
|
private String leakDetectorLevel; |
|
|
private String leakDetectorLevel; |
|
|
@Value("${transport.mqtt.netty.boss_group_thread_count}") |
|
|
@Value("${transport.mqtt.netty.boss_group_thread_count}") |
|
|
@ -59,6 +67,7 @@ public class MqttTransportService implements TbTransportService { |
|
|
private MqttTransportContext context; |
|
|
private MqttTransportContext context; |
|
|
|
|
|
|
|
|
private Channel serverChannel; |
|
|
private Channel serverChannel; |
|
|
|
|
|
private Channel sslServerChannel; |
|
|
private EventLoopGroup bossGroup; |
|
|
private EventLoopGroup bossGroup; |
|
|
private EventLoopGroup workerGroup; |
|
|
private EventLoopGroup workerGroup; |
|
|
|
|
|
|
|
|
@ -73,10 +82,18 @@ public class MqttTransportService implements TbTransportService { |
|
|
ServerBootstrap b = new ServerBootstrap(); |
|
|
ServerBootstrap b = new ServerBootstrap(); |
|
|
b.group(bossGroup, workerGroup) |
|
|
b.group(bossGroup, workerGroup) |
|
|
.channel(NioServerSocketChannel.class) |
|
|
.channel(NioServerSocketChannel.class) |
|
|
.childHandler(new MqttTransportServerInitializer(context)) |
|
|
.childHandler(new MqttTransportServerInitializer(context, false)) |
|
|
.childOption(ChannelOption.SO_KEEPALIVE, keepAlive); |
|
|
.childOption(ChannelOption.SO_KEEPALIVE, keepAlive); |
|
|
|
|
|
|
|
|
serverChannel = b.bind(host, port).sync().channel(); |
|
|
serverChannel = b.bind(host, port).sync().channel(); |
|
|
|
|
|
if (sslEnabled) { |
|
|
|
|
|
b = new ServerBootstrap(); |
|
|
|
|
|
b.group(bossGroup, workerGroup) |
|
|
|
|
|
.channel(NioServerSocketChannel.class) |
|
|
|
|
|
.childHandler(new MqttTransportServerInitializer(context, true)) |
|
|
|
|
|
.childOption(ChannelOption.SO_KEEPALIVE, keepAlive); |
|
|
|
|
|
sslServerChannel = b.bind(sslHost, sslPort).sync().channel(); |
|
|
|
|
|
} |
|
|
log.info("Mqtt transport started!"); |
|
|
log.info("Mqtt transport started!"); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -85,6 +102,9 @@ public class MqttTransportService implements TbTransportService { |
|
|
log.info("Stopping MQTT transport!"); |
|
|
log.info("Stopping MQTT transport!"); |
|
|
try { |
|
|
try { |
|
|
serverChannel.close().sync(); |
|
|
serverChannel.close().sync(); |
|
|
|
|
|
if (sslEnabled) { |
|
|
|
|
|
sslServerChannel.close().sync(); |
|
|
|
|
|
} |
|
|
} finally { |
|
|
} finally { |
|
|
workerGroup.shutdownGracefully(); |
|
|
workerGroup.shutdownGracefully(); |
|
|
bossGroup.shutdownGracefully(); |
|
|
bossGroup.shutdownGracefully(); |
|
|
|