wy471x commented on PR #6913:
URL: https://github.com/apache/shenyu/pull/6913#issuecomment-5294886502

   > This is a solid fix that brings the MQTT fan-out into line with the spec: 
it now delivers at `min(published QoS, granted QoS)` and allocates a fresh 
packet identifier from each subscriber channel's own id space, instead of 
hard-coding `AT_MOST_ONCE` and reusing the publisher's packet id. The change is 
well-scoped and backed by thorough unit tests (`PublishTest`, 
`SubscribeRepositoryTest`, `MqttPacketIdGeneratorTest`) covering granted-QoS 
fan-out, QoS0 → packetId 0, per-subscriber id spaces, max-QoS merge, FAILURE 
filtering, and id wrap-around at 0xFFFF.
   > 
   > A few non-blocking notes:
   > 
   > 1. **Dead code:** `int packetId = msg.variableHeader().packetId();` in 
`Publish.publish` is now unused (it used to be passed to `send`). Please remove 
it, along with the now-resolved `//// todo qos` comment.
   > 2. **Id-space cleanup:** `MqttPacketIdGenerator.CHANNEL_PACKET_ID_FACTORY` 
is a static `Map<Channel, AtomicInteger>` with strong references, cleaned only 
via `Disconnect.cleanChannel`. If a channel is torn down without going through 
`cleanChannel` (e.g. abnormal disconnect), its id-space entry leaks for the 
process lifetime. Consider a `WeakHashMap` keyed by `Channel`, or guaranteeing 
cleanup on all close paths.
   > 3. **Buffer sharing (pre-existing):** `send` wraps the shared incoming 
`payload` with `Unpooled.wrappedBuffer(payload)` for every subscriber. Since 
`WrappedByteBuf` shares the underlying buffer's reference count, fan-out to 
multiple active subscribers releases the same buffer N times. This is 
pre-existing behavior, but worth double-checking the reference-count handling 
(e.g. `payload.retain()` per subscriber) so it doesn't over-release under 
multi-subscriber fan-out.
   > 
   > I also confirmed `SubscribeRepository`'s `BaseRepository` type change 
(`List<Channel>` → `Map<Channel, MqttQoS>`) doesn't break other callers — the 
only callers (`Subscribe.add(channel, subscriptions)`, 
`Unsubscribe.remove(topics, channel)`, and `Publish` itself) use unchanged 
signatures.
   > 
   > Approving.
   
   fix: prevent mqtt packet id space leak and payload over-release on fan-out — 
addresses two review findings:                                                  
   
                                                                                
                                                                                
                        
     1. Packet id space leak (MqttPacketIdGenerator)                            
                                                                                
                        
     - Replaced the static ConcurrentHashMap<Channel, AtomicInteger> with 
Collections.synchronizedMap(new WeakHashMap<>())
     - Id-space entries were previously only removed via 
Disconnect.cleanChannel; channels torn down without a DISCONNECT (abnormal 
disconnect) leaked their entry for the process      
     lifetime                                                                   
                                                                                
                        
     - Weak keys now let the entry die with the channel, while explicit 
remove() remains for clean disconnects                                          
                                
                                                                                
                                                                                
                        
     2. Payload over-release under fan-out (Publish.send)                       
                                                                                
                        
     - Changed Unpooled.wrappedBuffer(payload) to 
Unpooled.wrappedBuffer(payload.retain())                                        
                                                      
     - The wrapped buffer shares the inbound payload's ref-count, so N 
subscribers releasing their outbound messages would release the same buffer N 
times                              
     - Each subscriber now owns one reference, keeping the ref-count balanced   
                                                                                
                        
                                                                                
                                                                                
                        
     3. Regression test (PublishTest)                                           
                                                                                
                        
     - Added testPublishFanOutRetainsPayloadPerSubscriber: publishes to 2 
subscribers, asserts refCnt goes 1 → 3 after fan-out and back to 1 after the 
delivered messages are released  
                                                                                
                                                                                
                        
     The "dead code" note from the review was verified against the code: the 
//// todo qos comment had already been removed in the prior commit, and 
packetId is still used by the      
     qos1/qos2 PUBACK path, so it was intentionally left intact. 


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