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>