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));