Aias00 commented on code in PR #6939:
URL: https://github.com/apache/shenyu/pull/6939#discussion_r4072306389


##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/MqttBootstrapServer.java:
##########
@@ -23,52 +23,64 @@
 import io.netty.channel.nio.NioEventLoopGroup;
 import io.netty.channel.socket.nio.NioServerSocketChannel;
 import io.netty.util.ResourceLeakDetector;
-import java.util.Locale;
-import java.util.Objects;
 import org.apache.shenyu.common.utils.Singleton;
 import org.apache.shenyu.protocol.mqtt.repositories.BaseRepository;
 import org.reflections.Reflections;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.Locale;
+import java.util.Objects;
+import java.util.concurrent.CompletableFuture;
 
 /**
  * mqtt server.
  */
 public class MqttBootstrapServer implements BootstrapServer {
 
+    private static final Logger LOG = 
LoggerFactory.getLogger(MqttBootstrapServer.class);
+
     private static final String REPOSITORY_PACKAGE_NAME = 
"org.apache.shenyu.protocol.mqtt.repositories";
 
     private static final MqttContext ENV = new MqttContext();
 
-    private EventLoopGroup bossGroup;
+    private volatile EventLoopGroup bossGroup;
 
-    private EventLoopGroup workerGroup;
+    private volatile EventLoopGroup workerGroup;
 
-    private ChannelFuture future;
+    private volatile ChannelFuture future;
 
     @Override
     public void init() {
         try {
             initRepositories();
         } catch (Exception e) {
-            //// todo log
+            LOG.error("MQTT server init repositories failed", e);
         }
     }
 
     @Override
     public void start() {
-        //// todo thread start mqtt server
         
ResourceLeakDetector.setLevel(ResourceLeakDetector.Level.valueOf(ENV.getLeakDetectorLevel().toUpperCase(Locale.ROOT)));
-        bossGroup = new NioEventLoopGroup(ENV.getBossGroupThreadCount());
-        workerGroup = new NioEventLoopGroup(ENV.getWorkerGroupThreadCount());
-        ServerBootstrap bootstrap = new ServerBootstrap();
-        bootstrap.group(bossGroup, workerGroup)
-                .channel(NioServerSocketChannel.class)
-                .childHandler(new 
MqttTransportServerInitializer(ENV.getMaxPayloadSize()));
-        try {
-            future = bootstrap.bind(ENV.getPort()).sync();
-            //// todo log
-        } catch (InterruptedException e) {
-            //// todo log
-        }
+        CompletableFuture.runAsync(() -> {

Review Comment:
   🔴 **Blocking — this makes `start()` fire-and-forget and breaks its callers.**
   
   `MqttPluginDataHandler` relies on `start()` blocking until the port is bound:
   ```java
   if (pluginData.getEnabled()) {
       server.start();
   } else {
       server.shutdown();
   }
   ```
   Now `start()` returns immediately, so:
   - `MqttPluginDataHandlerTest.testEnableConfiguration` fails on this PR's CI 
(`expected: <true> but was: <false>` — the port isn't bound yet);
   - worse in production: an enable → disable toggle can call `shutdown()` 
while the bind is still in flight. `shutdown()` sees all three fields as 
`null`, returns silently, and the server then comes up unowned — leaked 
`NioEventLoopGroup`s and a listening port that can never be stopped.
   
   Also, `runAsync` runs on `ForkJoinPool.commonPool()` while the task blocks 
on `bind().sync()` — blocking common-pool threads is best avoided.
   
   Please restore the synchronous bind, or keep the `CompletableFuture` in a 
field and expose it so callers can await completion, with `shutdown()` 
cancelling any pending start.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to