Copilot commented on code in PR #6902:
URL: https://github.com/apache/shenyu/pull/6902#discussion_r3755186316


##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Publish.java:
##########
@@ -124,4 +126,23 @@ private void send(final String topic, final ByteBuf 
payload, final int packetId)
             }
         });
     }
+
+    /**
+     * Publish a Last Will message to all subscribers of the will topic.
+     *
+     * @param will the will entry containing topic, message, qos, and retain 
flag
+     */
+    static void publishWill(final WillRepository.WillEntry will) {
+        List<Channel> channels = 
Singleton.INST.get(SubscribeRepository.class).get(will.getTopic());
+        MqttQoS willQos = MqttQoS.valueOf(will.getQos());
+        channels.parallelStream().forEach(channel -> {
+            if (channel.isActive()) {
+                MqttFixedHeader mqttFixedHeader = new 
MqttFixedHeader(MqttMessageType.PUBLISH, false, willQos, will.isRetain(), 0);
+                MqttPublishVariableHeader mqttPublishVariableHeader = new 
MqttPublishVariableHeader(will.getTopic(), 0);
+                MqttPublishMessage mqttPublishMessage = new 
MqttPublishMessage(mqttFixedHeader, mqttPublishVariableHeader,
+                        Unpooled.wrappedBuffer(will.getMessage()));
+                channel.writeAndFlush(mqttPublishMessage);
+            }
+        });
+    }

Review Comment:
   publishWill always uses packetId=0 in the PUBLISH variable header. For QoS 
1/2, MQTT requires a non-zero Packet Identifier; sending 0 can cause clients to 
treat the packet as malformed or ignore it. Generate a packetId when QoS>0 (and 
also guard against null will/topic/message to avoid runtime exceptions during 
disconnect handling).



##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandler.java:
##########
@@ -38,6 +42,16 @@ public void channelRead(final ChannelHandlerContext ctx, 
final Object msg) throw
         }
     }
 
+    @Override
+    public void channelInactive(final ChannelHandlerContext ctx) throws 
Exception {
+        WillRepository.WillEntry will = 
Singleton.INST.get(WillRepository.class).get(ctx.channel());
+        if (Objects.nonNull(will)) {
+            Publish.publishWill(will);
+            Singleton.INST.get(WillRepository.class).remove(ctx.channel());
+        }
+        super.channelInactive(ctx);
+    }

Review Comment:
   channelInactive removes the will only if Publish.publishWill completes 
successfully. If publishing throws (e.g., due to invalid QoS/topic/message), 
the will entry will remain in the repository and the exception can disrupt the 
pipeline. Remove the will in a finally block and avoid repeated Singleton 
lookups.



##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/Connect.java:
##########
@@ -67,6 +68,17 @@ public void connect(final ChannelHandlerContext ctx, final 
MqttConnectMessage ms
 
         // record connect
         Singleton.INST.get(ChannelRepository.class).add(ctx.channel(), 
clientId);
+
+        // store will if present
+        if (msg.variableHeader().isWillFlag()) {
+            WillRepository.WillEntry will = new WillRepository.WillEntry(
+                    msg.payload().willTopic(),
+                    msg.payload().willMessageInBytes(),
+                    msg.variableHeader().willQos(),
+                    msg.variableHeader().isWillRetain());
+            Singleton.INST.get(WillRepository.class).add(ctx.channel(), will);
+        }

Review Comment:
   CONNECT stores the Will entry without validating willTopic / willMessage. If 
either is null/blank, later LWT publishing can throw (e.g., 
SubscribeRepository.getOrDefault(null, …) or Unpooled.wrappedBuffer(null)), 
which would break channelInactive handling. Add a guard and only store the will 
when required fields are present (or reject the CONNECT).



##########
shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/repositories/WillRepositoryTest.java:
##########
@@ -0,0 +1,112 @@
+/*
+ * 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.repositories;
+
+import io.netty.channel.embedded.EmbeddedChannel;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.notNullValue;
+import static org.hamcrest.Matchers.nullValue;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class WillRepositoryTest {
+
+    private WillRepository willRepository;
+
+    private EmbeddedChannel channel;
+
+    @BeforeEach
+    public void setUp() {
+        willRepository = new WillRepository();
+        channel = new EmbeddedChannel();
+    }
+
+    @AfterEach
+    public void tearDown() {
+        channel.close();
+    }

Review Comment:
   WillRepository uses a static backing map, but this test class only closes 
the EmbeddedChannel in tearDown. That leaves closed channels referenced in the 
static map when a test stores a will without removing it, causing avoidable 
cross-test state and memory retention. Remove the entry in tearDown before 
closing the channel.



##########
shenyu-protocol/shenyu-protocol-mqtt/src/main/java/org/apache/shenyu/protocol/mqtt/repositories/WillRepository.java:
##########
@@ -0,0 +1,85 @@
+/*
+ * 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.repositories;
+
+import io.netty.channel.Channel;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+
+/**
+ * Stores Last Will and Testament for connected clients.
+ * Will is set on CONNECT and cleared on graceful DISCONNECT.
+ * On ungraceful disconnect (channelInactive with will present), the will is 
published.
+ */
+public class WillRepository implements BaseRepository<Channel, 
WillRepository.WillEntry> {
+
+    private static final Map<Channel, WillEntry> WILL_FACTORY = new 
ConcurrentHashMap<>();
+
+    @Override
+    public void add(final Channel channel, final WillEntry willEntry) {
+        WILL_FACTORY.put(channel, willEntry);
+    }
+
+    @Override
+    public void remove(final Channel channel) {
+        WILL_FACTORY.remove(channel);
+    }
+
+    @Override
+    public WillEntry get(final Channel channel) {
+        return WILL_FACTORY.get(channel);
+    }
+
+    /**
+     * Holds the will message fields from a CONNECT payload.
+     */
+    public static class WillEntry {
+
+        private final String topic;
+
+        private final byte[] message;
+
+        private final int qos;
+
+        private final boolean retain;
+
+        public WillEntry(final String topic, final byte[] message, final int 
qos, final boolean retain) {
+            this.topic = topic;
+            this.message = message;
+            this.qos = qos;
+            this.retain = retain;
+        }

Review Comment:
   WillEntry stores the caller-provided byte[] directly. Because byte[] is 
mutable, code that still holds a reference can change the will payload after it 
has been stored (or after it has been retrieved), leading to publishing a 
different message than intended. Defensive-copy the array when storing it.
   
   This issue also appears on line 73 of the same file.



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