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