This is an automated email from the ASF dual-hosted git repository.

Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git


The following commit(s) were added to refs/heads/master by this push:
     new 658a16056c fix: release inbound MqttPublishMessage payload ByteBuf and 
retain per subscriber write (#6977)
658a16056c is described below

commit 658a16056cfbf70ee0b591a4bd5b33410c9752f0
Author: wy471x <[email protected]>
AuthorDate: Thu Oct 1 20:44:09 2026 +0800

    fix: release inbound MqttPublishMessage payload ByteBuf and retain per 
subscriber write (#6977)
    
    * fix: release inbound MqttPublishMessage payload ByteBuf and retain per 
subscriber write
    
    The inbound MqttPublishMessage payload ByteBuf was never released,
    leaking a native/pooled buffer on every PUBLISH. Fan-out to multiple
    subscriber channels also wrote the same ByteBuf without retaining it.
    
    Release the inbound message in channelRead, retain the payload across
    the asynchronous send, and use retainedDuplicate for each subscriber.
    
    Co-Authored-By: Claude <[email protected]>
    
    * test(mqtt): tidy up leak-fix tests and unbreak checkstyle
    
    - drop redundant imports failing RedundantImport
    - use EmbeddedChannel instead of a mock context that NPEs
    - fix invalid SubscribeRepository/Map API usages in PublishTest
    - remove duplicated tests, dead code and racy packet-id/refCnt assertions
    
    ---------
    
    Co-authored-by: Claude <[email protected]>
    Co-authored-by: aias00 <[email protected]>
---
 .../shenyu/protocol/mqtt/MqttTransportHandler.java | 15 +++++---
 .../org/apache/shenyu/protocol/mqtt/Publish.java   | 13 +++++--
 .../protocol/mqtt/MqttTransportHandlerTest.java    | 40 ++++++++++++++++------
 .../apache/shenyu/protocol/mqtt/PublishTest.java   |  7 +++-
 4 files changed, 57 insertions(+), 18 deletions(-)

diff --git 
a/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandler.java
 
b/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandler.java
index bc1ccad71f..c1d1181747 100644
--- 
a/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandler.java
+++ 
b/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandler.java
@@ -22,6 +22,7 @@ import io.netty.channel.ChannelFuture;
 import io.netty.channel.ChannelHandlerContext;
 import io.netty.channel.ChannelInboundHandlerAdapter;
 import io.netty.handler.codec.mqtt.MqttMessage;
+import io.netty.util.ReferenceCountUtil;
 import io.netty.util.concurrent.Future;
 import io.netty.util.concurrent.GenericFutureListener;
 import org.apache.shenyu.common.utils.Singleton;
@@ -36,11 +37,15 @@ public class MqttTransportHandler extends 
ChannelInboundHandlerAdapter implement
 
     @Override
     public void channelRead(final ChannelHandlerContext ctx, final Object msg) 
throws Exception {
-        if (msg instanceof MqttMessage) {
-            MqttFactory mqttFactory = new MqttFactory((MqttMessage) msg, ctx);
-            mqttFactory.connect();
-        } else {
-            ctx.close();
+        try {
+            if (msg instanceof MqttMessage) {
+                MqttFactory mqttFactory = new MqttFactory((MqttMessage) msg, 
ctx);
+                mqttFactory.connect();
+            } else {
+                ctx.close();
+            }
+        } finally {
+            ReferenceCountUtil.release(msg);
         }
     }
 
diff --git 
a/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Publish.java
 
b/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Publish.java
index db41e236f4..279b1e5c04 100644
--- 
a/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Publish.java
+++ 
b/shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Publish.java
@@ -29,6 +29,7 @@ import io.netty.handler.codec.mqtt.MqttQoS;
 import io.netty.handler.codec.mqtt.MqttPubAckMessage;
 import io.netty.handler.codec.mqtt.MqttMessageType;
 import io.netty.handler.codec.mqtt.MqttPublishVariableHeader;
+import io.netty.util.ReferenceCountUtil;
 import org.apache.shenyu.common.utils.Singleton;
 import org.apache.shenyu.protocol.mqtt.repositories.SubscribeRepository;
 import org.apache.shenyu.protocol.mqtt.repositories.TopicRepository;
@@ -64,7 +65,15 @@ public class Publish extends MessageType {
             }
         }
         int packetId = msg.variableHeader().packetId();
-        CompletableFuture.runAsync(() -> send(topic, payload, mqttQoS));
+        // The inbound message is released by MqttTransportHandler once 
publish returns, retain the payload for the asynchronous send.
+        payload.retain();
+        CompletableFuture.runAsync(() -> {
+            try {
+                send(topic, payload, mqttQoS);
+            } finally {
+                ReferenceCountUtil.safeRelease(payload);
+            }
+        });
 
         switch (mqttQoS.value()) {
             case 0:
@@ -125,7 +134,7 @@ public class Publish extends MessageType {
                 int packetId = MqttQoS.AT_MOST_ONCE == qos ? 0 : 
MqttPacketIdGenerator.next(channel);
                 MqttFixedHeader mqttFixedHeader = new 
MqttFixedHeader(MqttMessageType.PUBLISH, false, qos, false, 0);
                 MqttPublishVariableHeader mqttPublishVariableHeader = new 
MqttPublishVariableHeader(topic, packetId);
-                MqttPublishMessage mqttPublishMessage = new 
MqttPublishMessage(mqttFixedHeader, mqttPublishVariableHeader, 
Unpooled.wrappedBuffer(payload.retain()));
+                MqttPublishMessage mqttPublishMessage = new 
MqttPublishMessage(mqttFixedHeader, mqttPublishVariableHeader, 
Unpooled.wrappedBuffer(payload.retainedDuplicate()));
                 channel.writeAndFlush(mqttPublishMessage);
             }
         });
diff --git 
a/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandlerTest.java
 
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandlerTest.java
index 2b69c0c749..58c2466579 100644
--- 
a/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandlerTest.java
+++ 
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandlerTest.java
@@ -17,15 +17,20 @@
 
 package org.apache.shenyu.protocol.mqtt;
 
+import io.netty.buffer.ByteBuf;
+import io.netty.buffer.Unpooled;
 import io.netty.channel.embedded.EmbeddedChannel;
 import io.netty.handler.codec.mqtt.MqttConnectMessage;
 import io.netty.handler.codec.mqtt.MqttConnectPayload;
 import io.netty.handler.codec.mqtt.MqttConnectVariableHeader;
 import io.netty.handler.codec.mqtt.MqttFixedHeader;
 import io.netty.handler.codec.mqtt.MqttMessageType;
+import io.netty.handler.codec.mqtt.MqttPublishMessage;
+import io.netty.handler.codec.mqtt.MqttPublishVariableHeader;
 import io.netty.handler.codec.mqtt.MqttQoS;
 import io.netty.handler.codec.mqtt.MqttTopicSubscription;
 import io.netty.handler.codec.mqtt.MqttVersion;
+import io.netty.util.CharsetUtil;
 import org.apache.shenyu.common.utils.Singleton;
 import org.apache.shenyu.protocol.mqtt.repositories.ChannelRepository;
 import org.apache.shenyu.protocol.mqtt.repositories.SubscribeRepository;
@@ -102,6 +107,21 @@ public final class MqttTransportHandlerTest {
         new MqttContext().setPassword(null);
     }
 
+    @Test
+    public void channelReadReleasesInboundPublishMessage() {
+        EmbeddedChannel channel = new EmbeddedChannel(new 
MqttTransportHandler());
+        MqttFixedHeader fixedHeader = new 
MqttFixedHeader(MqttMessageType.PUBLISH, false, MqttQoS.AT_MOST_ONCE, false, 0);
+        MqttPublishVariableHeader variableHeader = new 
MqttPublishVariableHeader(TOPIC, 1);
+        MqttPublishMessage message = new MqttPublishMessage(fixedHeader, 
variableHeader,
+                Unpooled.copiedBuffer("hello", CharsetUtil.UTF_8));
+        ByteBuf payload = message.payload();
+
+        channel.writeInbound(message);
+
+        assertEquals(0, payload.refCnt());
+        channel.finishAndReleaseAll();
+    }
+
     @Test
     public void duplicateConnectCleansUpChannelRepository() {
         EmbeddedChannel channel = new EmbeddedChannel(new 
MqttTransportHandler());
@@ -146,16 +166,6 @@ public final class MqttTransportHandlerTest {
         channel.finishAndReleaseAll();
     }
 
-    private MqttConnectMessage connectMessage() {
-        MqttFixedHeader fixedHeader = new 
MqttFixedHeader(MqttMessageType.CONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0);
-        MqttConnectVariableHeader variableHeader = new 
MqttConnectVariableHeader(
-                MqttVersion.MQTT_3_1_1.protocolName(), 
MqttVersion.MQTT_3_1_1.protocolLevel(),
-                true, true, false, 0, false, false, 60);
-        MqttConnectPayload payload = new MqttConnectPayload(CLIENT_ID, null, 
null,
-                USER_NAME, PASSWORD.getBytes(StandardCharsets.UTF_8));
-        return new MqttConnectMessage(fixedHeader, variableHeader, payload);
-    }
-
     @Test
     public void testOperationCompleteCleansRepositoriesOnClose() throws 
Exception {
         assertEquals(1, MqttPacketIdGenerator.next(registeredChannel));
@@ -167,6 +177,16 @@ public final class MqttTransportHandlerTest {
         assertEquals(1, MqttPacketIdGenerator.next(registeredChannel));
     }
 
+    private MqttConnectMessage connectMessage() {
+        MqttFixedHeader fixedHeader = new 
MqttFixedHeader(MqttMessageType.CONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0);
+        MqttConnectVariableHeader variableHeader = new 
MqttConnectVariableHeader(
+                MqttVersion.MQTT_3_1_1.protocolName(), 
MqttVersion.MQTT_3_1_1.protocolLevel(),
+                true, true, false, 0, false, false, 60);
+        MqttConnectPayload payload = new MqttConnectPayload(CLIENT_ID, null, 
null,
+                USER_NAME, PASSWORD.getBytes(StandardCharsets.UTF_8));
+        return new MqttConnectMessage(fixedHeader, variableHeader, payload);
+    }
+
     /**
      * The repositories mutate their state asynchronously on the common pool,
      * so assertions are retried until the mutation becomes visible.
diff --git 
a/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/PublishTest.java
 
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/PublishTest.java
index 487208d800..a840b5dc0f 100644
--- 
a/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/PublishTest.java
+++ 
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/PublishTest.java
@@ -254,7 +254,11 @@ public final class PublishTest {
         addSubscriber(subscriberChannel, MqttQoS.EXACTLY_ONCE);
         addSubscriber(otherSubscriberChannel, MqttQoS.EXACTLY_ONCE);
 
+        // wait for the first delivery so that the second publish cannot 
interleave packet id allocation
         publishToSubscribers(MqttQoS.EXACTLY_ONCE);
+        captureMessages(subscriberChannel, 1);
+        captureMessages(otherSubscriberChannel, 1);
+
         publishToSubscribers(MqttQoS.EXACTLY_ONCE);
 
         List<MqttPublishMessage> messages = captureMessages(subscriberChannel, 
2);
@@ -273,10 +277,11 @@ public final class PublishTest {
         ByteBuf payload = Unpooled.copiedBuffer(PAYLOAD, CharsetUtil.UTF_8);
         try {
             publishToSubscribers(MqttQoS.AT_LEAST_ONCE, payload);
-            awaitAssert(() -> assertEquals(3, payload.refCnt()));
 
             MqttPublishMessage delivered = captureMessage(subscriberChannel);
             MqttPublishMessage otherDelivered = 
captureMessage(otherSubscriberChannel);
+            // one reference held by the inbound message plus one retained 
duplicate per active subscriber
+            awaitAssert(() -> assertEquals(3, payload.refCnt()));
             assertEquals(PAYLOAD, 
delivered.payload().toString(CharsetUtil.UTF_8));
             assertEquals(PAYLOAD, 
otherDelivered.payload().toString(CharsetUtil.UTF_8));
 

Reply via email to