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

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


The following commit(s) were added to refs/heads/master by this push:
     new dbaf30986 Make sure to always close input streams bound to responses 
(#1970)
dbaf30986 is described below

commit dbaf30986168a4c66700193fba7b80b651a84445
Author: Konrad Windszus <[email protected]>
AuthorDate: Fri Jul 17 18:22:49 2026 +0200

    Make sure to always close input streams bound to responses (#1970)
    
    Ensure all connections are closed when transporter is closed.
    Add PUT IT where response contains body.
    
    This closes #1964
---
 maven-resolver-test-http/pom.xml                   |  4 +++
 .../aether/internal/test/util/http/HttpServer.java | 40 +++++++++++++++++++++-
 .../test/util/http/HttpTransporterTest.java        | 18 ++++++++++
 .../transport/apache/ApacheTransporterTest.java    | 13 +++++++
 .../aether/transport/jdk/JdkTransporter.java       | 26 ++++++--------
 .../aether/transport/jetty/JettyTransporter.java   | 14 ++++----
 pom.xml                                            |  6 ++++
 7 files changed, 96 insertions(+), 25 deletions(-)

diff --git a/maven-resolver-test-http/pom.xml b/maven-resolver-test-http/pom.xml
index 9521f6c74..cf43609e6 100644
--- a/maven-resolver-test-http/pom.xml
+++ b/maven-resolver-test-http/pom.xml
@@ -131,6 +131,10 @@
       <groupId>com.google.code.gson</groupId>
       <artifactId>gson</artifactId>
     </dependency>
+    <dependency>
+      <groupId>org.awaitility</groupId>
+      <artifactId>awaitility</artifactId>
+    </dependency>
   </dependencies>
 
 </project>
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..8d40d32e5 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,8 @@ 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.DatagramChannelEndPoint;
+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 +192,8 @@ public class HttpServer {
 
     private final List<LogEntry> logEntries = Collections.synchronizedList(new 
ArrayList<>());
 
+    private String responseBodyForPut;
+
     public String getHost() {
         return "localhost";
     }
@@ -323,6 +328,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 +407,27 @@ 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) {
+            // filter out the always present DatagramChannelEndPoint, which is 
not a real live connection
+            // (https://github.com/jetty/jetty.project/issues/15436)
+            http3Connector.getConnectedEndPoints().stream()
+                    .filter(endPoint -> !(endPoint instanceof 
DatagramChannelEndPoint))
+                    .forEach(connectedEndPoints::add);
+        }
+        return connectedEndPoints.size();
+    }
+
     private class CompressionEnforcingHandler extends CompressionHandler {
         // duplicate of CompressionHandler.pathConfigs which is private
         private final PathMappings<CompressionConfig> pathConfigs = new 
PathMappings<>();
@@ -662,7 +693,14 @@ public class HttpServer {
                         file.delete();
                         throw e;
                     }
-                    response.setStatus(HttpServletResponse.SC_NO_CONTENT);
+                    // 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_CREATED);
+                    } else {
+                        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 477b2aa19..46d9e6afb 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
@@ -39,10 +39,12 @@ import java.security.NoSuchAlgorithmException;
 import java.util.Enumeration;
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Supplier;
 import java.util.stream.Stream;
 
+import org.awaitility.Awaitility;
 import org.eclipse.aether.ConfigurationProperties;
 import org.eclipse.aether.DefaultRepositoryCache;
 import org.eclipse.aether.DefaultRepositorySystemSession;
@@ -305,6 +307,8 @@ public abstract class HttpTransporterTest {
             closer = null;
         }
         if (httpServer != null) {
+            // check for leaked connections (e.g., due to not closing response 
body streams)
+            Awaitility.await().atMost(5, TimeUnit.SECONDS).until(() -> 
httpServer.getNumConnectedEndPoints() == 0);
             httpServer.stop();
             httpServer = null;
         }
@@ -1499,6 +1503,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-apache/src/test/java/org/eclipse/aether/transport/apache/ApacheTransporterTest.java
 
b/maven-resolver-transport-apache/src/test/java/org/eclipse/aether/transport/apache/ApacheTransporterTest.java
index bbd28856e..ce3041c16 100644
--- 
a/maven-resolver-transport-apache/src/test/java/org/eclipse/aether/transport/apache/ApacheTransporterTest.java
+++ 
b/maven-resolver-transport-apache/src/test/java/org/eclipse/aether/transport/apache/ApacheTransporterTest.java
@@ -33,6 +33,7 @@ import 
org.eclipse.aether.internal.test.util.http.RecordingTransportListener;
 import org.eclipse.aether.spi.connector.transport.GetTask;
 import org.eclipse.aether.spi.connector.transport.PutTask;
 import org.eclipse.aether.spi.io.PathProcessorSupport;
+import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Test;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -68,6 +69,18 @@ class ApacheTransporterTest extends HttpTransporterTest {
         return false;
     }
 
+    @AfterEach
+    @Override
+    protected void tearDown() throws Exception {
+        // make sure to also release any connection in the global state 
(otherwise the check for connection leaks will
+        // fail)
+        GlobalState globalState = GlobalState.get(session);
+        if (globalState != null) {
+            globalState.close();
+        }
+        super.tearDown();
+    }
+
     @Test
     void testGet_WebDav() throws Exception {
         httpServer.setWebDav(true);
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..f6d52a368 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,17 +432,14 @@ 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 {
-                    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) {
             throw enhance(e);
@@ -458,6 +450,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..33801f29c 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)) {
@@ -404,6 +403,11 @@ final class JettyTransporter extends AbstractTransporter 
implements HttpTranspor
                     .send(listener);
             response = listener.get(requestTimeout, TimeUnit.MILLISECONDS);
             
task.getListener().transportPropertiesAvailable(createTransportProperties(request,
 rawResponseHeaders));
+            if (response.getStatus() >= MULTIPLE_CHOICES) {
+                JettyRFC9457Reporter.INSTANCE.generateException(listener, 
(statusCode, reasonPhrase) -> {
+                    throw new HttpTransporterException(statusCode);
+                });
+            }
         } catch (ExecutionException e) {
             Throwable t = e.getCause();
             if (t instanceof IOException ioex) {
@@ -418,12 +422,6 @@ final class JettyTransporter extends AbstractTransporter 
implements HttpTranspor
                 throw new RuntimeException(t);
             }
         }
-
-        if (response.getStatus() >= MULTIPLE_CHOICES) {
-            JettyRFC9457Reporter.INSTANCE.generateException(listener, 
(statusCode, reasonPhrase) -> {
-                throw new HttpTransporterException(statusCode);
-            });
-        }
     }
 
     @Override
diff --git a/pom.xml b/pom.xml
index 410feb06f..d65dc936f 100644
--- a/pom.xml
+++ b/pom.xml
@@ -338,6 +338,12 @@
         <artifactId>jmh-generator-annprocess</artifactId>
         <version>${jmhVersion}</version>
       </dependency>
+      <dependency>
+        <groupId>org.awaitility</groupId>
+        <artifactId>awaitility</artifactId>
+        <!-- release 4.3.1 was not successful 
(https://github.com/awaitility/awaitility/issues/306) -->
+        <version>4.3.0</version>
+      </dependency>
     </dependencies>
   </dependencyManagement>
 

Reply via email to