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]