wy471x opened a new pull request, #6977: URL: https://github.com/apache/shenyu/pull/6977
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. <!-- Describe your PR here; e.g. Fixes #issueNo --> <!-- Thank you for proposing a pull request. This template will guide you through the essential steps necessary for a pull request. --> Make sure that: - [X] You have read the [contribution guidelines](https://shenyu.apache.org/community/contributor-guide). - [X] You submit test cases (unit or integration tests) that back your changes. - [X] Your local test passed `./mvnw clean install -Dmaven.javadoc.skip=true`. ## Summary ### Problems 1. Inbound leak: MqttTransportHandler extends ChannelInboundHandlerAdapter, so inbound messages are not auto-released, and the handler never called ReferenceCountUtil.release(msg). The MqttPublishMessage payload ByteBuf leaked on every PUBLISH. 2. Unsafe multi-subscriber fan-out: CompletableFuture.runAsync(() -> send(topic, payload, packetId)) wrapped the same payload with Unpooled.wrappedBuffer(payload) for every subscriber channel in parallel without a per-write retain — risking IllegalReferenceCountException or buffer corruption, and leaking after the first release. ### Changes - MqttTransportHandler.java — channelRead now releases the inbound message in a finally block via ReferenceCountUtil.release(msg). - Publish.java - payload.retain() before the asynchronous send, with ReferenceCountUtil.safeRelease(payload) in the task's finally, since the handler releases the inbound message once publish returns. - Replaced Unpooled.wrappedBuffer(payload) with payload.retainedDuplicate() so each subscriber write holds its own reference. ### Tests - MqttTransportHandlerTest (new) — verifies channelRead releases the inbound MqttPublishMessage (payload refCnt reaches 0) and closes the channel for non-MQTT messages. - PublishTest — added publishDeliversPayloadToEachSubscriber, verifying the payload is delivered to each active subscriber channel with correct reference counting. close [#6639](https://github.com/apache/shenyu/issues/6639) -- 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]
