This is an automated email from the ASF dual-hosted git repository.

kwin pushed a commit to branch bugfix/consume-put-response
in repository https://gitbox.apache.org/repos/asf/maven-resolver.git

commit ea9d2e7fcb197e93aea47d1c7d9dd6eeba6ca9d1
Author: Konrad Windszus <[email protected]>
AuthorDate: Thu Jul 16 20:15:50 2026 +0200

    Make sure to always close input streams bound to responses
    
    Ensure all connections are closed when transporter is closed
    Add PUT IT where response contains body
    
    This closes #1964
---
 .../aether/internal/test/util/http/HttpServer.java | 31 ++++++++++++++++++++++
 .../test/util/http/HttpTransporterTest.java        | 18 ++++++++++++-
 .../aether/transport/jdk/JdkTransporter.java       | 16 +++++------
 .../aether/transport/jetty/JettyTransporter.java   |  3 +--
 4 files changed, 56 insertions(+), 12 deletions(-)

diff --git 
a/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpServer.java
 
b/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpServer.java
index f7aacb35f..0ab1b1eb6 100644
--- 
a/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpServer.java
+++ 
b/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpServer.java
@@ -27,6 +27,7 @@ import java.nio.file.Files;
 import java.nio.file.StandardOpenOption;
 import java.util.ArrayList;
 import java.util.Base64;
+import java.util.Collection;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
@@ -60,6 +61,7 @@ import 
org.eclipse.jetty.http3.server.HTTP3ServerConnectionFactory;
 import org.eclipse.jetty.http3.server.HTTP3ServerQuicConfiguration;
 import org.eclipse.jetty.io.ByteBufferPool;
 import org.eclipse.jetty.io.Content;
+import org.eclipse.jetty.io.EndPoint;
 import org.eclipse.jetty.quic.quiche.server.QuicheServerConnector;
 import org.eclipse.jetty.quic.quiche.server.QuicheServerQuicConfiguration;
 import org.eclipse.jetty.server.Handler;
@@ -189,6 +191,8 @@ public class HttpServer {
 
     private final List<LogEntry> logEntries = Collections.synchronizedList(new 
ArrayList<>());
 
+    private String responseBodyForPut;
+
     public String getHost() {
         return "localhost";
     }
@@ -323,6 +327,11 @@ public class HttpServer {
         return this;
     }
 
+    public HttpServer setResponseBodyForPut(String body) {
+        this.responseBodyForPut = body;
+        return this;
+    }
+
     public HttpServer setWebDav(boolean webDav) {
         this.webDav = webDav;
         return this;
@@ -397,6 +406,23 @@ public class HttpServer {
         }
     }
 
+    public int getNumConnectedEndPoints() {
+        if (server.isStopped()) {
+            throw new IllegalStateException("Server is stopped");
+        }
+        Collection<EndPoint> connectedEndPoints = new ArrayList<>();
+        if (httpConnector != null) {
+            connectedEndPoints.addAll(httpConnector.getConnectedEndPoints());
+        }
+        if (httpsConnector != null) {
+            connectedEndPoints.addAll(httpsConnector.getConnectedEndPoints());
+        }
+        if (http3Connector != null) {
+            connectedEndPoints.addAll(http3Connector.getConnectedEndPoints());
+        }
+        return connectedEndPoints.size();
+    }
+
     private class CompressionEnforcingHandler extends CompressionHandler {
         // duplicate of CompressionHandler.pathConfigs which is private
         private final PathMappings<CompressionConfig> pathConfigs = new 
PathMappings<>();
@@ -662,6 +688,11 @@ public class HttpServer {
                         file.delete();
                         throw e;
                     }
+                    // optionally add some response body to test that the 
client can handle it, even though Maven
+                    // Repository Protocol doesn't mention it
+                    if (responseBodyForPut != null) {
+                        writeResponseBodyMessage(req, response, 
responseBodyForPut);
+                    }
                     response.setStatus(HttpServletResponse.SC_NO_CONTENT);
                 } else {
                     response.setStatus(HttpServletResponse.SC_FORBIDDEN);
diff --git 
a/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpTransporterTest.java
 
b/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpTransporterTest.java
index 546a3ac9a..50bfde3a2 100644
--- 
a/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpTransporterTest.java
+++ 
b/maven-resolver-test-http/src/main/java/org/eclipse/aether/internal/test/util/http/HttpTransporterTest.java
@@ -304,6 +304,8 @@ public abstract class HttpTransporterTest {
             closer.run();
             closer = null;
         }
+        // check for leaked connections (e.g., due to not closing response 
body streams)
+        assertEquals(0, httpServer.getConnectedEndPoints());
         if (httpServer != null) {
             httpServer.stop();
             httpServer = null;
@@ -918,7 +920,7 @@ public abstract class HttpTransporterTest {
         // issue 2 requests to ensure that the second request is HTTP/3 (the 
first one may be HTTP/2 if the TCP
         // connection is faster)
         transporter.get(task);
-        transporter.get(task);
+        //transporter.get(task);
         assertEquals("test", task.getDataString());
         assertEquals(0L, listener.getDataOffset());
         assertEquals(4L, listener.getDataLength());
@@ -1495,6 +1497,20 @@ public abstract class HttpTransporterTest {
         assertEquals(1, httpServer.getLogEntries().size()); // put w/ auth
     }
 
+    @Test
+    protected void testPut_WithResponseBody() throws Exception {
+        httpServer.setAuthentication("testuser", "testpass");
+        httpServer.setResponseBodyForPut("Some dummy response body");
+        auth = new AuthenticationBuilder()
+                .addUsername("testuser")
+                .addPassword("testpass")
+                .build();
+        newTransporter(httpServer.getHttpUrl());
+        PutTask task = new 
PutTask(URI.create("repo/file.txt")).setDataString("upload");
+        transporter.put(task);
+        // this leads to stuck threads in some transporters if the response 
body is not consumed
+    }
+
     @Test
     @Timeout(20)
     protected void testConcurrency() throws Exception {
diff --git 
a/maven-resolver-transport-jdk-parent/maven-resolver-transport-jdk11/src/main/java/org/eclipse/aether/transport/jdk/JdkTransporter.java
 
b/maven-resolver-transport-jdk-parent/maven-resolver-transport-jdk11/src/main/java/org/eclipse/aether/transport/jdk/JdkTransporter.java
index e6568459f..6d43cbee6 100644
--- 
a/maven-resolver-transport-jdk-parent/maven-resolver-transport-jdk11/src/main/java/org/eclipse/aether/transport/jdk/JdkTransporter.java
+++ 
b/maven-resolver-transport-jdk-parent/maven-resolver-transport-jdk11/src/main/java/org/eclipse/aether/transport/jdk/JdkTransporter.java
@@ -291,16 +291,11 @@ final class JdkTransporter extends AbstractTransporter 
implements HttpTransporte
                             resume = false;
                             continue;
                         }
-                        try {
-                            
JdkRFC9457Reporter.INSTANCE.generateException(response, (statusCode, 
reasonPhrase) -> {
-                                throw new HttpTransporterException(statusCode);
-                            });
-                        } finally {
-                            closeBody(response);
-                        }
+                        
JdkRFC9457Reporter.INSTANCE.generateException(response, (statusCode, 
reasonPhrase) -> {
+                            throw new HttpTransporterException(statusCode);
+                        });
                     }
                 } catch (ConnectException e) {
-                    closeBody(response);
                     throw enhance(e);
                 }
                 break;
@@ -437,8 +432,9 @@ final class JdkTransporter extends AbstractTransporter 
implements HttpTransporte
                     task.getDataLength()));
         }
         prepare(request);
+        HttpResponse<InputStream> response = null;
         try {
-            HttpResponse<InputStream> response = send(request.build(), 
HttpResponse.BodyHandlers.ofInputStream());
+            response = send(request.build(), 
HttpResponse.BodyHandlers.ofInputStream());
             
task.getListener().transportPropertiesAvailable(createTransportProperties(response));
             if (response.statusCode() >= MULTIPLE_CHOICES) {
                 try {
@@ -458,6 +454,8 @@ final class JdkTransporter extends AbstractTransporter 
implements HttpTransporte
                 throw (TransferCancelledException) rootCause;
             }
             throw e;
+        } finally {
+            closeBody(response);
         }
     }
 
diff --git 
a/maven-resolver-transport-jetty/src/main/java/org/eclipse/aether/transport/jetty/JettyTransporter.java
 
b/maven-resolver-transport-jetty/src/main/java/org/eclipse/aether/transport/jetty/JettyTransporter.java
index 9d57d5ea4..c8d32c1bb 100644
--- 
a/maven-resolver-transport-jetty/src/main/java/org/eclipse/aether/transport/jetty/JettyTransporter.java
+++ 
b/maven-resolver-transport-jetty/src/main/java/org/eclipse/aether/transport/jetty/JettyTransporter.java
@@ -373,8 +373,7 @@ final class JettyTransporter extends AbstractTransporter 
implements HttpTranspor
         });
         AtomicBoolean started = new AtomicBoolean(false);
         Response response;
-        InputStreamResponseListener listener = new 
InputStreamResponseListener();
-        try {
+        try (InputStreamResponseListener listener = new 
InputStreamResponseListener()) {
             request.onRequestCommit(r -> {
                         if (task.getDataLength() == 0) {
                             if (started.compareAndSet(false, true)) {

Reply via email to