|
|
@ -20,8 +20,7 @@ import io.netty.channel.Channel; |
|
|
import io.netty.channel.EventLoopGroup; |
|
|
import io.netty.channel.EventLoopGroup; |
|
|
import io.netty.channel.nio.NioEventLoopGroup; |
|
|
import io.netty.channel.nio.NioEventLoopGroup; |
|
|
import io.netty.channel.socket.nio.NioServerSocketChannel; |
|
|
import io.netty.channel.socket.nio.NioServerSocketChannel; |
|
|
import io.netty.handler.logging.LogLevel; |
|
|
import io.netty.util.ResourceLeakDetector; |
|
|
import io.netty.handler.logging.LoggingHandler; |
|
|
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
import org.springframework.beans.factory.annotation.Autowired; |
|
|
import org.springframework.beans.factory.annotation.Autowired; |
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
@ -67,6 +66,14 @@ public class MqttTransportService { |
|
|
@Value("${mqtt.adaptor}") |
|
|
@Value("${mqtt.adaptor}") |
|
|
private String adaptorName; |
|
|
private String adaptorName; |
|
|
|
|
|
|
|
|
|
|
|
@Value("${mqtt.netty.leak_detector_level}") |
|
|
|
|
|
private String leakDetectorLevel; |
|
|
|
|
|
@Value("${mqtt.netty.boss_group_thread_count}") |
|
|
|
|
|
private Integer bossGroupThreadCount; |
|
|
|
|
|
@Value("${mqtt.netty.worker_group_thread_count}") |
|
|
|
|
|
private Integer workerGroupThreadCount; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private MqttTransportAdaptor adaptor; |
|
|
private MqttTransportAdaptor adaptor; |
|
|
|
|
|
|
|
|
private Channel serverChannel; |
|
|
private Channel serverChannel; |
|
|
@ -75,17 +82,19 @@ public class MqttTransportService { |
|
|
|
|
|
|
|
|
@PostConstruct |
|
|
@PostConstruct |
|
|
public void init() throws Exception { |
|
|
public void init() throws Exception { |
|
|
|
|
|
log.info("Setting resource leak detector level to {}", leakDetectorLevel); |
|
|
|
|
|
ResourceLeakDetector.setLevel(ResourceLeakDetector.Level.valueOf(leakDetectorLevel.toUpperCase())); |
|
|
|
|
|
|
|
|
log.info("Starting MQTT transport..."); |
|
|
log.info("Starting MQTT transport..."); |
|
|
log.info("Lookup MQTT transport adaptor {}", adaptorName); |
|
|
log.info("Lookup MQTT transport adaptor {}", adaptorName); |
|
|
this.adaptor = (MqttTransportAdaptor) appContext.getBean(adaptorName); |
|
|
this.adaptor = (MqttTransportAdaptor) appContext.getBean(adaptorName); |
|
|
|
|
|
|
|
|
log.info("Starting MQTT transport server"); |
|
|
log.info("Starting MQTT transport server"); |
|
|
bossGroup = new NioEventLoopGroup(1); |
|
|
bossGroup = new NioEventLoopGroup(bossGroupThreadCount); |
|
|
workerGroup = new NioEventLoopGroup(); |
|
|
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) |
|
|
.handler(new LoggingHandler(LogLevel.TRACE)) |
|
|
|
|
|
.childHandler(new MqttTransportServerInitializer(processor, deviceService, authService, adaptor, sslHandlerProvider)); |
|
|
.childHandler(new MqttTransportServerInitializer(processor, deviceService, authService, adaptor, sslHandlerProvider)); |
|
|
|
|
|
|
|
|
serverChannel = b.bind(host, port).sync().channel(); |
|
|
serverChannel = b.bind(host, port).sync().channel(); |
|
|
|