From c8906d89b5341ff7502aa98e74f7ab9973b53261 Mon Sep 17 00:00:00 2001 From: cyforkk <1139009771@qq.com> Date: Fri, 31 Jul 2026 19:30:06 +0800 Subject: [PATCH 1/2] [ISSUE #16399] Close late input after DefaultStreamingDecoder shutdown Mirror the #16391 fix: when decode() is called after the decoder has been closed (via close() or onStreamClosed()), close the input stream instead of ignoring it. This prevents ByteBuf leaks in the non-gRPC Triple HTTP/2 streaming path when late DATA frames arrive after the decoder has started closing. Co-Authored-By: Claude Fable 5 --- .../message/DefaultStreamingDecoder.java | 6 +- .../message/DefaultStreamingDecoderTest.java | 64 +++++++++++++++++++ 2 files changed, 69 insertions(+), 1 deletion(-) create mode 100644 dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java diff --git a/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoder.java b/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoder.java index 6d341783c0a5..f4696475779e 100644 --- a/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoder.java +++ b/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoder.java @@ -38,7 +38,11 @@ public void request(int numMessages) { @Override public void decode(InputStream inputStream) throws DecodeException { if (closed) { - // ignored + try { + inputStream.close(); + } catch (IOException e) { + throw new DecodeException(e); + } return; } accumulate.addInputStream(inputStream); diff --git a/dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java b/dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java new file mode 100644 index 000000000000..07d93fc4938e --- /dev/null +++ b/dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java @@ -0,0 +1,64 @@ +/* + * 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.dubbo.remoting.http12.message; + +import java.io.ByteArrayInputStream; +import java.io.IOException; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +class DefaultStreamingDecoderTest { + + @Test + void closesInputReceivedAfterStreamClosed() { + DefaultStreamingDecoder decoder = new DefaultStreamingDecoder(); + decoder.onStreamClosed(); + CloseTrackingInputStream inputStream = new CloseTrackingInputStream(new byte[] {0}); + + decoder.decode(inputStream); + + assertTrue(inputStream.closed); + } + + @Test + void closesInputReceivedAfterDecoderClosed() { + DefaultStreamingDecoder decoder = new DefaultStreamingDecoder(); + decoder.close(); + CloseTrackingInputStream inputStream = new CloseTrackingInputStream(new byte[] {0}); + + decoder.decode(inputStream); + + assertTrue(inputStream.closed); + } + + private static final class CloseTrackingInputStream extends ByteArrayInputStream { + + private boolean closed; + + private CloseTrackingInputStream(byte[] buf) { + super(buf); + } + + @Override + public void close() throws IOException { + this.closed = true; + super.close(); + } + } +} From 64aa67b3c06b715faecb2574b6aded1cbb8fe2b6 Mon Sep 17 00:00:00 2001 From: cyforkk <1139009771@qq.com> Date: Fri, 31 Jul 2026 20:06:04 +0800 Subject: [PATCH 2/2] Add test for close failure propagation in DefaultStreamingDecoder Cover the IOException -> DecodeException branch to improve patch coverage. Co-Authored-By: Claude Fable 5 --- .../message/DefaultStreamingDecoderTest.java | 23 +++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java b/dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java index 07d93fc4938e..5dd18382821e 100644 --- a/dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java +++ b/dubbo-remoting/dubbo-remoting-http12/src/test/java/org/apache/dubbo/remoting/http12/message/DefaultStreamingDecoderTest.java @@ -16,11 +16,15 @@ */ package org.apache.dubbo.remoting.http12.message; +import org.apache.dubbo.remoting.http12.exception.DecodeException; + import java.io.ByteArrayInputStream; import java.io.IOException; +import java.io.InputStream; import org.junit.jupiter.api.Test; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; class DefaultStreamingDecoderTest { @@ -47,6 +51,25 @@ void closesInputReceivedAfterDecoderClosed() { assertTrue(inputStream.closed); } + @Test + void propagatesCloseFailureAsDecodeException() { + DefaultStreamingDecoder decoder = new DefaultStreamingDecoder(); + decoder.onStreamClosed(); + InputStream failingStream = new InputStream() { + @Override + public int read() throws IOException { + return -1; + } + + @Override + public void close() throws IOException { + throw new IOException("close failed"); + } + }; + + assertThrows(DecodeException.class, () -> decoder.decode(failingStream)); + } + private static final class CloseTrackingInputStream extends ByteArrayInputStream { private boolean closed;