Copilot commented on code in PR #6913:
URL: https://github.com/apache/shenyu/pull/6913#discussion_r4043580028
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Publish.java:
##########
@@ -62,7 +64,7 @@ public void publish(final ChannelHandlerContext ctx, final
MqttPublishMessage ms
}
}
int packetId = msg.variableHeader().packetId();
- CompletableFuture.runAsync(() -> send(topic, payload, packetId));
+ CompletableFuture.runAsync(() -> send(topic, payload, mqttQoS));
Review Comment:
Every QoS 2 PUBLISH is fanned out immediately without recording its packet
ID. If PUBREC is lost, the publisher retransmits the same PUBLISH with DUP set,
and this line fans it out again, violating exactly-once delivery. Record
inbound QoS 2 state and suppress duplicate processing until the corresponding
PUBREL completes the exchange.
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/MqttFactory.java:
##########
@@ -65,6 +65,9 @@ public void connect() {
case PINGREQ:
messageType.pingReq(ctx);
break;
+ case PUBREL:
+ messageType.pubRel(ctx, msg);
+ break;
Review Comment:
The new broker-to-subscriber QoS exchanges have no acknowledgement path. QoS
1 subscribers return PUBACK, which still falls through to a no-op, while QoS 2
subscribers first return PUBREC, for which there is no case at all; therefore
the broker never completes or releases these outbound exchanges. Add outbound
in-flight state and dispatch PUBACK/PUBREC/PUBCOMP separately from this
inbound-publisher PUBREL path.
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/utils/MqttPacketIdGenerator.java:
##########
@@ -0,0 +1,60 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.protocol.mqtt.utils;
+
+import io.netty.channel.Channel;
+
+import java.util.Collections;
+import java.util.Map;
+import java.util.WeakHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
+
+/**
+ * Allocates packet identifiers for outbound messages from each channel's own
id space.
+ */
+public final class MqttPacketIdGenerator {
+
+ private static final int MIN_PACKET_ID = 1;
+
+ private static final int MAX_PACKET_ID = 0xFFFF;
+
+ // weak keys so channels closed without a DISCONNECT do not leak their id
space
+ private static final Map<Channel, AtomicInteger> CHANNEL_PACKET_ID_FACTORY
= Collections.synchronizedMap(new WeakHashMap<>());
+
+ private MqttPacketIdGenerator() {
+ }
+
+ /**
+ * get next packet id of the channel.
+ * @param channel channel
+ * @return next packet id
+ */
+ public static int next(final Channel channel) {
+ AtomicInteger packetId =
CHANNEL_PACKET_ID_FACTORY.computeIfAbsent(channel, key -> new AtomicInteger());
+ return packetId.updateAndGet(current -> current >= MAX_PACKET_ID ?
MIN_PACKET_ID : current + 1);
Review Comment:
Wrapping unconditionally can reissue packet ID 1 while the earlier packet
using ID 1 is still awaiting acknowledgement. After 65,535 outstanding QoS 1/2
publications, acknowledgements become ambiguous. Track leased IDs per channel,
release them on PUBACK/PUBCOMP, and apply backpressure or fail allocation when
all IDs are occupied.
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/PubRel.java:
##########
@@ -0,0 +1,41 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.protocol.mqtt;
+
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.handler.codec.mqtt.MqttFixedHeader;
+import io.netty.handler.codec.mqtt.MqttMessage;
+import io.netty.handler.codec.mqtt.MqttMessageIdVariableHeader;
+import io.netty.handler.codec.mqtt.MqttQoS;
+
+import static io.netty.handler.codec.mqtt.MqttMessageType.PUBCOMP;
+
+/**
+ * The PUBREL message is the third message of the QoS 2 protocol flow,
+ * the server responds with PUBCOMP to release the packet id.
+ */
+public class PubRel extends MessageType {
+
+ @Override
+ public void pubRel(final ChannelHandlerContext ctx, final MqttMessage msg)
{
+ MqttMessageIdVariableHeader variableHeader =
(MqttMessageIdVariableHeader) msg.variableHeader();
+ MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(PUBCOMP, false,
MqttQoS.AT_MOST_ONCE, false, 0);
+ MqttMessage mqttPubCompMessage = new MqttMessage(mqttFixedHeader,
variableHeader);
+ ctx.writeAndFlush(mqttPubCompMessage);
+ }
Review Comment:
Unlike the other post-CONNECT handlers, this new PUBREL path does not verify
connection state. A client can send PUBREL before CONNECT and receive PUBCOMP
instead of having the connection closed, violating the MQTT packet ordering
requirement.
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/repositories/SubscribeRepository.java:
##########
@@ -55,11 +53,11 @@ public void add(final List<String> topics, final
List<Channel> channels) {
* @param mqttTopicSubscription mqtt subscription info
*/
public void add(final Channel channel, final List<MqttTopicSubscription>
mqttTopicSubscription) {
- CompletableFuture.runAsync(() ->
mqttTopicSubscription.parallelStream().forEach(s -> {
- List<Channel> channels = get(s.topicName());
- channels.add(channel);
- TOPIC_CHANNEL_FACTORY.put(s.topicName(), channels);
- }));
+ CompletableFuture.runAsync(() -> mqttTopicSubscription.parallelStream()
+ .filter(s -> s.qualityOfService() != MqttQoS.FAILURE)
+ .forEach(s -> TOPIC_CHANNEL_FACTORY
+ .computeIfAbsent(s.topicName(), key -> new
ConcurrentHashMap<>())
+ .merge(channel, s.qualityOfService(),
SubscribeRepository::maxQoS)));
Review Comment:
The stored value is the requested QoS, not the granted QoS.
`Subscribe.sendSubAckMessage` still returns `AT_MOST_ONCE` for every topic
(`Subscribe.java:90-99`), so a client requesting QoS 1 is told that QoS 0 was
granted but is later sent QoS 1 publications. Compute the grant once and use
the same value both in this repository and in the SUBACK payload.
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Disconnect.java:
##########
@@ -45,5 +46,6 @@ public void disconnect(final ChannelHandlerContext ctx) {
private void cleanChannel(final Channel channel) {
//// todo ttl
Singleton.INST.get(ChannelRepository.class).remove(channel);
+ MqttPacketIdGenerator.remove(channel);
Review Comment:
This cleanup is never invoked for an incoming DISCONNECT packet because
`MqttFactory` still routes `DISCONNECT` to its no-op/default branch
(`MqttFactory.java:71-74`). Thus the stated cleanup-on-DISCONNECT behavior only
occurs later if the peer closes the socket. Dispatch DISCONNECT to
`MessageType.disconnect` so the server closes and cleans the channel
immediately.
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/repositories/SubscribeRepository.java:
##########
@@ -55,11 +53,11 @@ public void add(final List<String> topics, final
List<Channel> channels) {
* @param mqttTopicSubscription mqtt subscription info
*/
public void add(final Channel channel, final List<MqttTopicSubscription>
mqttTopicSubscription) {
- CompletableFuture.runAsync(() ->
mqttTopicSubscription.parallelStream().forEach(s -> {
- List<Channel> channels = get(s.topicName());
- channels.add(channel);
- TOPIC_CHANNEL_FACTORY.put(s.topicName(), channels);
- }));
+ CompletableFuture.runAsync(() -> mqttTopicSubscription.parallelStream()
+ .filter(s -> s.qualityOfService() != MqttQoS.FAILURE)
+ .forEach(s -> TOPIC_CHANNEL_FACTORY
+ .computeIfAbsent(s.topicName(), key -> new
ConcurrentHashMap<>())
+ .merge(channel, s.qualityOfService(),
SubscribeRepository::maxQoS)));
Review Comment:
Using `maxQoS` here prevents an existing subscription from being downgraded.
MQTT requires a new subscription with an identical topic filter to replace the
old subscription, so re-subscribing from QoS 2 to QoS 0 must store QoS 0 rather
than retaining QoS 2.
##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/repositories/SubscribeRepository.java:
##########
@@ -74,27 +72,41 @@ public void remove(final List<String> topics) {
*/
public void remove(final List<String> topics, final Channel channel) {
CompletableFuture.runAsync(() -> topics.parallelStream().forEach(topic
-> {
- List<Channel> channels = TOPIC_CHANNEL_FACTORY.get(topic);
- if (CollectionUtils.isNotEmpty(channels)) {
- channels.remove(channel);
+ Map<Channel, MqttQoS> subscribers =
TOPIC_CHANNEL_FACTORY.get(topic);
+ if (Objects.nonNull(subscribers)) {
+ subscribers.remove(channel);
}
}));
}
+ /**
+ * remove the channel from all topics it subscribed.
+ * @param channel channel
+ */
+ public void remove(final Channel channel) {
+ CompletableFuture.runAsync(() ->
TOPIC_CHANNEL_FACTORY.values().parallelStream()
+ .forEach(subscribers -> subscribers.remove(channel)));
Review Comment:
This asynchronous removal can race the asynchronous `add` operation
submitted for the same channel. If the close task removes the channel before an
earlier subscription task inserts it, that task subsequently leaves the closed
channel registered, so abrupt disconnects can still leak subscriptions. Make
these map mutations synchronous or serialize them per channel.
--
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]