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]

Reply via email to