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


##########
shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttTransportHandlerTest.java:
##########
@@ -79,6 +93,24 @@ public void duplicateConnectCleansUpChannelRepository() {
         channel.finishAndReleaseAll();
     }
 
+    @Test
+    public void channelReadReleasesInboundMessage() throws Exception {
+        MqttFixedHeader fixedHeader = new 
MqttFixedHeader(MqttMessageType.PUBLISH, false, MqttQoS.AT_MOST_ONCE, false, 0);
+        MqttPublishVariableHeader variableHeader = new 
MqttPublishVariableHeader("test/topic", 1);
+        MqttPublishMessage msg = new MqttPublishMessage(fixedHeader, 
variableHeader, Unpooled.copiedBuffer("hello", CharsetUtil.UTF_8));
+        new 
MqttTransportHandler().channelRead(mock(ChannelHandlerContext.class), msg);

Review Comment:
   For a PUBLISH packet, `channelRead` dispatches to `Publish.publish`, which 
dereferences `ctx.channel()`; this unstubbed mock therefore throws 
`NullPointerException`, failing the test even though the `finally` block 
releases the message. Drive the handler through an `EmbeddedChannel` or stub a 
real connected channel on the context.
   
   This issue also appears on line 102 of the same file.



##########
shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/PublishTest.java:
##########
@@ -87,14 +96,14 @@ static void tearDown() {
 
     @Test
     public void retainedPublishStoresMessage() {
-        new Publish().publish(connectedContext(), 
publishMessage(RETAINED_TOPIC, "hello", true));
+        new Publish().publish(mock(ChannelHandlerContext.class), 
publishMessage(RETAINED_TOPIC, "hello", true));

Review Comment:
   `Publish.publish` immediately calls `ctx.channel()` and then 
`channel.attr(...)`; an unstubbed Mockito context returns `null`, so this test 
fails with a `NullPointerException` before storing the retained message. 
Restore a context backed by a connected `EmbeddedChannel` (the removed 
`connectedContext()` helper already provided this).
   
   This issue also appears in the following locations of the same file:
   - line 106
   - line 137
   - line 140
   - line 155



##########
shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/PublishTest.java:
##########
@@ -125,17 +134,33 @@ public void publishAfterConnectOnSameChannelIsAccepted() {
     @Test
     public void zeroByteRetainedPublishClearsRetainedMessage() {
         Publish publish = new Publish();
-        publish.publish(connectedContext(), publishMessage(CLEARED_TOPIC, 
"hello", true));
+        publish.publish(mock(ChannelHandlerContext.class), 
publishMessage(CLEARED_TOPIC, "hello", true));
         await().atMost(Duration.ofSeconds(5))
                 .until(() -> 
"hello".equals(topicRepository.get(CLEARED_TOPIC)));
-        publish.publish(connectedContext(), publishMessage(CLEARED_TOPIC, "", 
true));
+        publish.publish(mock(ChannelHandlerContext.class), 
publishMessage(CLEARED_TOPIC, "", true));
         assertNull(topicRepository.get(CLEARED_TOPIC));
     }
 
-    private ChannelHandlerContext connectedContext() {
-        EmbeddedChannel channel = new EmbeddedChannel(new 
ChannelInboundHandlerAdapter());
-        new MessageType().setConnected(channel, true);
-        return channel.pipeline().lastContext();
+    @Test
+    public void publishDeliversPayloadToEachSubscriber() {
+        Channel channel1 = mock(Channel.class);
+        when(channel1.isActive()).thenReturn(true);
+        Channel channel2 = mock(Channel.class);
+        when(channel2.isActive()).thenReturn(true);
+        SubscribeRepository subscribeRepository = 
Singleton.INST.get(SubscribeRepository.class);
+        subscribeRepository.add(Collections.singletonList("test/fanout"), 
Arrays.asList(channel1, channel2));
+        await().atMost(Duration.ofSeconds(5))
+                .until(() -> 
subscribeRepository.get("test/fanout").contains(channel1) && 
subscribeRepository.get("test/fanout").contains(channel2));
+        MqttPublishMessage msg = publishMessage("test/fanout", "hello", false);
+        new Publish().publish(mock(ChannelHandlerContext.class), msg);
+        // one reference held by the inbound message plus one per active 
subscriber after the send completes.
+        await().atMost(Duration.ofSeconds(5))
+                .until(() -> msg.payload().refCnt() == 3);
+        ArgumentCaptor<MqttPublishMessage> captor = 
ArgumentCaptor.forClass(MqttPublishMessage.class);
+        verify(channel1).writeAndFlush(captor.capture());
+        verify(channel2).writeAndFlush(captor.capture());
+        captor.getAllValues().forEach(message -> assertEquals("hello", 
message.payload().toString(CharsetUtil.UTF_8)));
+        captor.getAllValues().forEach(ReferenceCountUtil::safeRelease);

Review Comment:
   After releasing the two captured duplicates, the original inbound message 
still owns the reference described above (`refCnt() == 1`). Release `msg` as 
well so this reference-counting test does not itself leak its source buffer.



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