This is an automated email from the ASF dual-hosted git repository. reta pushed a commit to branch 3.5.x-fixes in repository https://gitbox.apache.org/repos/asf/cxf.git
commit e7a80e40e701a67a418e4b5da3eef9e3ea39bb27 Author: Andriy Redko <[email protected]> AuthorDate: Wed Feb 5 17:10:54 2025 -0500 CXF-8629: AsyncHTTPConduit (hc5) should support chunked request / response. Add test cases with auto-redirect (#2202) (cherry picked from commit 51ad92012fbcfbdd77b722214631303850315799) (cherry picked from commit d6dbb447a75bfc506f20012efc6e9138da36a293) # Conflicts: # systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/FileStore.java # systests/transports/src/test/java/org/apache/cxf/systest/jaxrs/FileStore.java (cherry picked from commit af434bf3e47ad3e77e7c5dd189fce027de9a3b84) # Conflicts: # rt/transports/http-hc5/src/main/java/org/apache/cxf/transport/http/asyncclient/hc5/URLConnectionAsyncHTTPConduit.java # systests/jaxrs/src/test/java/org/apache/cxf/systest/jaxrs/JAXRSAsyncClientChunkingTest.java # systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/FileStore.java # systests/transports/src/test/java/org/apache/cxf/systest/jaxrs/FileStore.java --- .../http/asyncclient/hc5/AsyncHTTPConduit.java | 5 +- .../http/netty/client/NettyHttpConduit.java | 5 +- .../apache/cxf/systest/hc5/jaxrs/FileStore.java | 59 +++++++- .../hc5/jaxrs/JAXRSAsyncClientChunkingTest.java | 151 +++++++++++++++++++-- 4 files changed, 206 insertions(+), 14 deletions(-) diff --git a/rt/transports/http-hc5/src/main/java/org/apache/cxf/transport/http/asyncclient/hc5/AsyncHTTPConduit.java b/rt/transports/http-hc5/src/main/java/org/apache/cxf/transport/http/asyncclient/hc5/AsyncHTTPConduit.java index cc6ad8967b..e911969149 100644 --- a/rt/transports/http-hc5/src/main/java/org/apache/cxf/transport/http/asyncclient/hc5/AsyncHTTPConduit.java +++ b/rt/transports/http-hc5/src/main/java/org/apache/cxf/transport/http/asyncclient/hc5/AsyncHTTPConduit.java @@ -709,7 +709,10 @@ public class AsyncHTTPConduit extends URLConnectionHTTPConduit { } protected void handleResponseAsync() throws IOException { - isAsync = true; + // The response hasn't been handled yet, should be handled asynchronously + if (httpResponse == null) { + isAsync = true; + } } protected void closeInputStream() throws IOException { diff --git a/rt/transports/http-netty/netty-client/src/main/java/org/apache/cxf/transport/http/netty/client/NettyHttpConduit.java b/rt/transports/http-netty/netty-client/src/main/java/org/apache/cxf/transport/http/netty/client/NettyHttpConduit.java index 32e1946a18..0679700e23 100644 --- a/rt/transports/http-netty/netty-client/src/main/java/org/apache/cxf/transport/http/netty/client/NettyHttpConduit.java +++ b/rt/transports/http-netty/netty-client/src/main/java/org/apache/cxf/transport/http/netty/client/NettyHttpConduit.java @@ -591,7 +591,10 @@ public class NettyHttpConduit extends URLConnectionHTTPConduit implements BusLif @Override protected void handleResponseAsync() throws IOException { - isAsync = true; + // The response hasn't been handled yet, should be handled asynchronously + if (httpResponse == null) { + isAsync = true; + } } @Override diff --git a/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/FileStore.java b/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/FileStore.java index 1490f9e21f..dc96f4ce27 100644 --- a/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/FileStore.java +++ b/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/FileStore.java @@ -29,6 +29,7 @@ import java.util.concurrent.ConcurrentMap; import javax.activation.DataHandler; import javax.ws.rs.Consumes; +import javax.ws.rs.GET; import javax.ws.rs.POST; import javax.ws.rs.Path; import javax.ws.rs.QueryParam; @@ -38,8 +39,10 @@ import javax.ws.rs.container.Suspended; import javax.ws.rs.core.Context; import javax.ws.rs.core.HttpHeaders; import javax.ws.rs.core.Response; +import javax.ws.rs.core.Response.ResponseBuilder; import javax.ws.rs.core.Response.Status; import javax.ws.rs.core.StreamingOutput; +import javax.ws.rs.core.UriBuilder; import javax.ws.rs.core.UriInfo; import org.apache.cxf.common.util.StringUtils; @@ -55,8 +58,9 @@ public class FileStore { @POST @Path("/stream") @Consumes("*/*") + @SuppressWarnings("PMD.UseTryWithResources") - public Response addBook(@QueryParam("chunked") boolean chunked, InputStream in) throws IOException { + public Response addFile(@QueryParam("chunked") boolean chunked, InputStream in) throws IOException { String transferEncoding = headers.getHeaderString("Transfer-Encoding"); if (chunked != Objects.equals("chunked", transferEncoding)) { @@ -81,11 +85,11 @@ public class FileStore { in.close(); } } - } + } @POST @Consumes("multipart/form-data") - public void addBook(@QueryParam("chunked") boolean chunked, + public void addFile(@QueryParam("chunked") boolean chunked, @Suspended final AsyncResponse response, @Context final UriInfo uri, final MultipartBody body) { String transferEncoding = headers.getHeaderString("Transfer-Encoding"); @@ -142,4 +146,53 @@ public class FileStore { } } } + + @GET + @Consumes("multipart/form-data") + public void getFile(@QueryParam("chunked") boolean chunked, @QueryParam("filename") String source, + @Suspended final AsyncResponse response) { + + if (StringUtils.isEmpty(source)) { + response.resume(Response.status(Status.BAD_REQUEST).build()); + return; + } + + try { + if (!store.containsKey(source)) { + response.resume(Response.status(Status.NOT_FOUND).build()); + return; + } + + final byte[] content = store.get(source); + if (response.isSuspended()) { + final StreamingOutput stream = new StreamingOutput() { + @Override + public void write(OutputStream os) throws IOException, WebApplicationException { + if (chunked) { + // Make sure we have enough data for chunking to kick in + for (int i = 0; i < 10; ++i) { + os.write(content); + } + } else { + os.write(content); + } + } + }; + response.resume(Response.ok().entity(stream).build()); + } + + } catch (final Exception ex) { + response.resume(Response.serverError().build()); + } + } + + @GET + @Path("/redirect") + public Response redirectFile(@Context UriInfo uriInfo) { + final UriBuilder builder = uriInfo.getBaseUriBuilder().path(getClass()); + uriInfo.getQueryParameters(true).forEach((p, v) -> builder.queryParam(p, v.get(0))); + + final ResponseBuilder response = Response.status(303).header("Location", builder.build()); + return response.build(); + } } diff --git a/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/JAXRSAsyncClientChunkingTest.java b/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/JAXRSAsyncClientChunkingTest.java index 8bac8c05f0..c7251a8465 100644 --- a/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/JAXRSAsyncClientChunkingTest.java +++ b/systests/transport-hc5/src/test/java/org/apache/cxf/systest/hc5/jaxrs/JAXRSAsyncClientChunkingTest.java @@ -25,6 +25,13 @@ import java.io.InputStream; import java.util.Arrays; import java.util.Collection; import java.util.Random; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.logging.Logger; import javax.ws.rs.client.Entity; import javax.ws.rs.core.MediaType; @@ -32,6 +39,7 @@ import javax.ws.rs.core.MultivaluedMap; import javax.ws.rs.core.Response; import org.apache.cxf.interceptor.LoggingInInterceptor; +import org.apache.cxf.interceptor.LoggingMessage; import org.apache.cxf.interceptor.LoggingOutInterceptor; import org.apache.cxf.jaxrs.client.ClientConfiguration; import org.apache.cxf.jaxrs.client.WebClient; @@ -40,6 +48,7 @@ import org.apache.cxf.jaxrs.ext.multipart.MultipartBody; import org.apache.cxf.jaxrs.impl.MetadataMap; import org.apache.cxf.jaxrs.model.AbstractResourceInfo; import org.apache.cxf.jaxrs.provider.MultipartProvider; +import org.apache.cxf.message.Message; import org.apache.cxf.testutil.common.AbstractBusClientServerTestBase; import org.apache.cxf.transport.http.asyncclient.hc5.AsyncHTTPConduit; @@ -50,6 +59,7 @@ import org.junit.runners.Parameterized.Parameters; import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.CoreMatchers.not; +import static org.hamcrest.CoreMatchers.startsWith; import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; @@ -57,9 +67,12 @@ import static org.junit.Assert.assertTrue; public class JAXRSAsyncClientChunkingTest extends AbstractBusClientServerTestBase { private static final String PORT = allocatePort(FileStoreServer.class); private final Boolean chunked; + private final Boolean autoRedirect; + private final ConcurrentMap<String, AtomicInteger> ids = new ConcurrentHashMap<>(); - public JAXRSAsyncClientChunkingTest(Boolean chunked) { + public JAXRSAsyncClientChunkingTest(Boolean chunked, Boolean autoRedirect) { this.chunked = chunked; + this.autoRedirect = autoRedirect; } @BeforeClass @@ -69,9 +82,14 @@ public class JAXRSAsyncClientChunkingTest extends AbstractBusClientServerTestBas createStaticBus(); } - @Parameters(name = "{0}") - public static Collection<Boolean> data() { - return Arrays.asList(new Boolean[] {Boolean.FALSE, Boolean.TRUE}); + @Parameters(name = "chunked {0}, auto-redirect {1}") + public static Collection<Boolean[]> data() { + return Arrays.asList(new Boolean[][] { + {Boolean.FALSE /* chunked */, Boolean.FALSE /* autoredirect */}, + {Boolean.FALSE /* chunked */, Boolean.TRUE /* autoredirect */}, + {Boolean.TRUE /* chunked */, Boolean.FALSE /* autoredirect */}, + {Boolean.TRUE /* chunked */, Boolean.TRUE /* autoredirect */}, + }); } @Test @@ -83,17 +101,18 @@ public class JAXRSAsyncClientChunkingTest extends AbstractBusClientServerTestBas final ClientConfiguration config = WebClient.getConfig(webClient); config.getBus().setProperty(AsyncHTTPConduit.USE_ASYNC, true); config.getHttpConduit().getClient().setAllowChunking(chunked); + config.getHttpConduit().getClient().setAutoRedirect(autoRedirect); configureLogging(config); + final String filename = "keymanagers.jks"; try { - final String filename = "keymanagers.jks"; final MultivaluedMap<String, String> headers = new MetadataMap<>(); headers.add("Content-ID", filename); headers.add("Content-Type", "application/binary"); - headers.add("Content-Disposition", "attachment; filename=" + chunked + "_" + filename); + headers.add("Content-Disposition", "attachment; filename=" + chunked + "_" + autoRedirect + "_" + filename); final Attachment att = new Attachment(getClass().getResourceAsStream("/" + filename), headers); final MultipartBody entity = new MultipartBody(att); - try (Response response = webClient.header("Content-Type", "multipart/form-data").post(entity)) { + try (Response response = webClient.header("Content-Type", MediaType.MULTIPART_FORM_DATA).post(entity)) { assertThat(response.getStatus(), equalTo(201)); assertThat(response.getHeaderString("Transfer-Encoding"), equalTo(chunked ? "chunked" : null)); assertThat(response.getEntity(), not(equalTo(null))); @@ -101,6 +120,43 @@ public class JAXRSAsyncClientChunkingTest extends AbstractBusClientServerTestBas } finally { webClient.close(); } + + assertRedirect(chunked + "_" + autoRedirect + "_" + filename); + } + + @Test + public void testMultipartChunkingAsync() throws InterruptedException, ExecutionException, TimeoutException { + final String url = "http://localhost:" + PORT + "/file-store"; + final WebClient webClient = WebClient.create(url, Arrays.asList(new MultipartProvider())) + .query("chunked", chunked); + + final ClientConfiguration config = WebClient.getConfig(webClient); + config.getBus().setProperty(AsyncHTTPConduit.USE_ASYNC, true); + config.getHttpConduit().getClient().setAllowChunking(chunked); + config.getHttpConduit().getClient().setAutoRedirect(autoRedirect); + configureLogging(config); + + final String filename = "keymanagers.jks"; + try { + final MultivaluedMap<String, String> headers = new MetadataMap<>(); + headers.add("Content-ID", filename); + headers.add("Content-Type", "application/binary"); + headers.add("Content-Disposition", "attachment; filename=" + chunked + + "_" + autoRedirect + "_async_" + filename); + final Attachment att = new Attachment(getClass().getResourceAsStream("/" + filename), headers); + final Entity<MultipartBody> entity = Entity.entity(new MultipartBody(att), + MediaType.MULTIPART_FORM_DATA_TYPE); + try (Response response = webClient.header("Content-Type", MediaType.MULTIPART_FORM_DATA).async() + .post(entity).get(10, TimeUnit.SECONDS)) { + assertThat(response.getStatus(), equalTo(201)); + assertThat(response.getHeaderString("Transfer-Encoding"), equalTo(chunked ? "chunked" : null)); + assertThat(response.getEntity(), not(equalTo(null))); + } + } finally { + webClient.close(); + } + + assertRedirect(chunked + "_" + autoRedirect + "_" + filename); } @Test @@ -111,6 +167,7 @@ public class JAXRSAsyncClientChunkingTest extends AbstractBusClientServerTestBas final ClientConfiguration config = WebClient.getConfig(webClient); config.getBus().setProperty(AsyncHTTPConduit.USE_ASYNC, true); config.getHttpConduit().getClient().setAllowChunking(chunked); + config.getHttpConduit().getClient().setAutoRedirect(autoRedirect); configureLogging(config); final byte[] bytes = new byte [32 * 1024]; @@ -127,13 +184,89 @@ public class JAXRSAsyncClientChunkingTest extends AbstractBusClientServerTestBas } finally { webClient.close(); } + + assertNoDuplicateLogging(); } - + + @Test + public void testStreamChunkingAsync() throws IOException, InterruptedException, + ExecutionException, TimeoutException { + final String url = "http://localhost:" + PORT + "/file-store/stream"; + final WebClient webClient = WebClient.create(url).query("chunked", chunked); + + final ClientConfiguration config = WebClient.getConfig(webClient); + config.getBus().setProperty(AsyncHTTPConduit.USE_ASYNC, true); + config.getHttpConduit().getClient().setAllowChunking(chunked); + config.getHttpConduit().getClient().setAutoRedirect(autoRedirect); + configureLogging(config); + + final byte[] bytes = new byte [32 * 1024]; + final Random random = new Random(); + random.nextBytes(bytes); + + try (InputStream in = new ByteArrayInputStream(bytes)) { + final Entity<InputStream> entity = Entity.entity(in, MediaType.APPLICATION_OCTET_STREAM); + try (Response response = webClient.async().post(entity).get(10, TimeUnit.SECONDS)) { + assertThat(response.getStatus(), equalTo(200)); + assertThat(response.getHeaderString("Transfer-Encoding"), equalTo(chunked ? "chunked" : null)); + assertThat(response.getEntity(), not(equalTo(null))); + } + } finally { + webClient.close(); + } + + assertNoDuplicateLogging(); + } + + private void assertRedirect(String filename) { + final String url = "http://localhost:" + PORT + "/file-store/redirect"; + + final WebClient webClient = WebClient.create(url, Arrays.asList(new MultipartProvider())) + .query("chunked", chunked) + .query("filename", filename); + + final ClientConfiguration config = WebClient.getConfig(webClient); + config.getBus().setProperty(AsyncHTTPConduit.USE_ASYNC, true); + config.getHttpConduit().getClient().setAllowChunking(chunked); + config.getHttpConduit().getClient().setAutoRedirect(autoRedirect); + configureLogging(config); + + try { + try (Response response = webClient.get()) { + if (autoRedirect) { + assertThat(response.getStatus(), equalTo(200)); + assertThat(response.getHeaderString("Transfer-Encoding"), equalTo(chunked ? "chunked" : null)); + assertThat(response.getEntity(), not(equalTo(null))); + } else { + assertThat(response.getStatus(), equalTo(303)); + assertThat(response.getHeaderString("Location"), + startsWith("http://localhost:" + PORT + "/file-store")); + } + } + } finally { + webClient.close(); + } + + assertNoDuplicateLogging(); + } + + private void assertNoDuplicateLogging() { + ids.forEach((id, counter) -> assertThat("Duplicate client logging for message " + id, + counter.get(), equalTo(1))); + } + private void configureLogging(final ClientConfiguration config) { final LoggingOutInterceptor out = new LoggingOutInterceptor(); out.setShowMultipartContent(false); - final LoggingInInterceptor in = new LoggingInInterceptor(); + final LoggingInInterceptor in = new LoggingInInterceptor() { + @Override + protected void logging(Logger logger, Message message) { + super.logging(logger, message); + final String id = (String) message.get(LoggingMessage.ID_KEY); + ids.computeIfAbsent(id, key -> new AtomicInteger()).incrementAndGet(); + } + }; in.setShowBinaryContent(false); config.getInInterceptors().add(in);
