diff --git a/sdk/core/azure-core-http-jdk-httpclient/src/main/java/com/azure/core/http/jdk/httpclient/JdkHttpClient.java b/sdk/core/azure-core-http-jdk-httpclient/src/main/java/com/azure/core/http/jdk/httpclient/JdkHttpClient.java index 65b710343bea..45581015ae83 100644 --- a/sdk/core/azure-core-http-jdk-httpclient/src/main/java/com/azure/core/http/jdk/httpclient/JdkHttpClient.java +++ b/sdk/core/azure-core-http-jdk-httpclient/src/main/java/com/azure/core/http/jdk/httpclient/JdkHttpClient.java @@ -8,6 +8,7 @@ import com.azure.core.http.HttpHeaders; import com.azure.core.http.HttpRequest; import com.azure.core.http.HttpResponse; +import com.azure.core.http.jdk.httpclient.implementation.BodyIgnoringSubscriber; import com.azure.core.util.Context; import com.azure.core.util.Contexts; import com.azure.core.util.CoreUtils; @@ -20,7 +21,6 @@ import reactor.core.publisher.Mono; import java.io.IOException; -import java.io.InputStream; import java.io.UncheckedIOException; import java.net.URISyntaxException; import java.util.List; @@ -37,6 +37,9 @@ */ class JdkHttpClient implements HttpClient { private static final ClientLogger LOGGER = new ClientLogger(JdkHttpClient.class); + private static final String AZURE_EAGERLY_READ_RESPONSE = "azure-eagerly-read-response"; + private static final String AZURE_IGNORE_RESPONSE_BODY = "azure-ignore-response-body"; + private static final byte[] IGNORED_BODY = new byte[0]; private final java.net.http.HttpClient jdkHttpClient; @@ -61,11 +64,24 @@ public Mono send(HttpRequest request) { @Override public Mono send(HttpRequest request, Context context) { - boolean eagerlyReadResponse = (boolean) context.getData("azure-eagerly-read-response").orElse(false); + boolean eagerlyReadResponse = (boolean) context.getData(AZURE_EAGERLY_READ_RESPONSE).orElse(false); + boolean ignoreResponseBody = (boolean) context.getData(AZURE_IGNORE_RESPONSE_BODY).orElse(false); return Mono.fromCallable(() -> toJdkHttpRequest(request, context)) .flatMap(jdkRequest -> Mono.fromCompletionStage(jdkHttpClient.sendAsync(jdkRequest, ofPublisher())) .flatMap(jdKResponse -> { + // Ignoring the response body takes precedent over eagerly reading the response body. + // Both should never be true at the same time but this is acts as a safeguard. + if (ignoreResponseBody) { + HttpHeaders headers = fromJdkHttpHeaders(jdKResponse.headers()); + int statusCode = jdKResponse.statusCode(); + + return JdkFlowAdapter.flowPublisherToFlux(jdKResponse.body()) + .ignoreElements() + .then(Mono.fromSupplier(() -> + new JdkHttpResponseSync(request, statusCode, headers, IGNORED_BODY))); + } + if (eagerlyReadResponse) { HttpHeaders headers = fromJdkHttpHeaders(jdKResponse.headers()); int statusCode = jdKResponse.statusCode(); @@ -82,16 +98,24 @@ public Mono send(HttpRequest request, Context context) { @Override public HttpResponse sendSync(HttpRequest request, Context context) { - boolean eagerlyReadResponse = (boolean) context.getData("azure-eagerly-read-response").orElse(false); + boolean eagerlyReadResponse = (boolean) context.getData(AZURE_EAGERLY_READ_RESPONSE).orElse(false); + boolean ignoreResponseBody = (boolean) context.getData(AZURE_IGNORE_RESPONSE_BODY).orElse(false); java.net.http.HttpRequest jdkRequest = toJdkHttpRequest(request, context); try { - if (eagerlyReadResponse) { + // Ignoring the response body takes precedent over eagerly reading the response body. + // Both should never be true at the same time but this is acts as a safeguard. + if (ignoreResponseBody) { + java.net.http.HttpResponse jdKResponse = jdkHttpClient.send(jdkRequest, + responseInfo -> new BodyIgnoringSubscriber(LOGGER)); + return new JdkHttpResponseSync(request, jdKResponse.statusCode(), + fromJdkHttpHeaders(jdKResponse.headers()), IGNORED_BODY); + } else if (eagerlyReadResponse) { java.net.http.HttpResponse jdKResponse = jdkHttpClient.send(jdkRequest, ofByteArray()); - return new JdkHttpResponseSync(request, jdKResponse.statusCode(), fromJdkHttpHeaders(jdKResponse.headers()), jdKResponse.body()); + return new JdkHttpResponseSync(request, jdKResponse.statusCode(), + fromJdkHttpHeaders(jdKResponse.headers()), jdKResponse.body()); } else { - java.net.http.HttpResponse jdKResponse = jdkHttpClient.send(jdkRequest, ofInputStream()); - return new JdkHttpResponseSync(request, jdKResponse); + return new JdkHttpResponseSync(request, jdkHttpClient.send(jdkRequest, ofInputStream())); } } catch (IOException e) { throw LOGGER.logExceptionAsError(new UncheckedIOException(e)); diff --git a/sdk/core/azure-core-http-jdk-httpclient/src/main/java/com/azure/core/http/jdk/httpclient/implementation/BodyIgnoringSubscriber.java b/sdk/core/azure-core-http-jdk-httpclient/src/main/java/com/azure/core/http/jdk/httpclient/implementation/BodyIgnoringSubscriber.java new file mode 100644 index 000000000000..876e631d994a --- /dev/null +++ b/sdk/core/azure-core-http-jdk-httpclient/src/main/java/com/azure/core/http/jdk/httpclient/implementation/BodyIgnoringSubscriber.java @@ -0,0 +1,68 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +package com.azure.core.http.jdk.httpclient.implementation; + +import com.azure.core.http.HttpClient; +import com.azure.core.util.logging.ClientLogger; +import com.azure.core.util.logging.LogLevel; + +import java.net.http.HttpResponse; +import java.nio.ByteBuffer; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.Flow; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * Implementation of a {@link HttpResponse.BodySubscriber} that ignores the response body. + *

+ * This is used when the {@link HttpClient} is told to ignore the response body, used when the returned body value is + * {@code void} or {@code Void}. + *

+ * This will log a warning message if a response body is received to indicate that there was a bug either in determining + * that the response body should be ignored, the Swagger indicated no response body would be received but was, or that + * the server sent a response body when it shouldn't. + */ +public final class BodyIgnoringSubscriber implements HttpResponse.BodySubscriber { + private final CompletableFuture completableFuture; + private final ClientLogger logger; + private final AtomicBoolean subscribed = new AtomicBoolean(); + + public BodyIgnoringSubscriber(ClientLogger logger) { + this.completableFuture = new CompletableFuture<>(); + this.logger = logger; + } + + @Override + public CompletionStage getBody() { + return completableFuture; + } + + @Override + public void onSubscribe(Flow.Subscription subscription) { + if (!subscribed.compareAndSet(false, true)) { + // Only can have one subscription. + subscription.cancel(); + } else { + subscription.request(Long.MAX_VALUE); + } + } + + @Override + public void onNext(List item) { + logger.log(LogLevel.WARNING, () -> "Received HTTP response body when one wasn't expected. " + + "Response body will be ignored as directed."); + } + + @Override + public void onError(Throwable throwable) { + completableFuture.completeExceptionally(throwable); + } + + @Override + public void onComplete() { + completableFuture.complete(null); + } +} diff --git a/sdk/core/azure-core-http-netty/src/main/java/com/azure/core/http/netty/NettyAsyncHttpClient.java b/sdk/core/azure-core-http-netty/src/main/java/com/azure/core/http/netty/NettyAsyncHttpClient.java index 836b52f7ca75..70699134c66c 100644 --- a/sdk/core/azure-core-http-netty/src/main/java/com/azure/core/http/netty/NettyAsyncHttpClient.java +++ b/sdk/core/azure-core-http-netty/src/main/java/com/azure/core/http/netty/NettyAsyncHttpClient.java @@ -26,6 +26,7 @@ import com.azure.core.util.Contexts; import com.azure.core.util.ProgressReporter; import com.azure.core.util.logging.ClientLogger; +import com.azure.core.util.logging.LogLevel; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; import io.netty.channel.EventLoopGroup; @@ -51,6 +52,7 @@ import java.nio.file.StandardOpenOption; import java.time.Duration; import java.util.Objects; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.BiFunction; import static com.azure.core.http.netty.implementation.Utility.closeConnection; @@ -58,9 +60,9 @@ /** * This class provides a Netty-based implementation for the {@link HttpClient} interface. Creating an instance of this * class can be achieved by using the {@link NettyAsyncHttpClientBuilder} class, which offers Netty-specific API for - * features such as {@link NettyAsyncHttpClientBuilder#eventLoopGroup(EventLoopGroup) thread pooling}, {@link - * NettyAsyncHttpClientBuilder#wiretap(boolean) wiretapping}, {@link NettyAsyncHttpClientBuilder#proxy(ProxyOptions) - * setProxy configuration}, and much more. + * features such as {@link NettyAsyncHttpClientBuilder#eventLoopGroup(EventLoopGroup) thread pooling}, + * {@link NettyAsyncHttpClientBuilder#wiretap(boolean) wiretapping}, + * {@link NettyAsyncHttpClientBuilder#proxy(ProxyOptions) setProxy configuration}, and much more. * * @see HttpClient * @see NettyAsyncHttpClientBuilder @@ -70,6 +72,7 @@ class NettyAsyncHttpClient implements HttpClient { private static final byte[] EMPTY_BYTES = new byte[0]; private static final String AZURE_EAGERLY_READ_RESPONSE = "azure-eagerly-read-response"; + private static final String AZURE_IGNORE_RESPONSE_BODY = "azure-ignore-response-body"; private static final String AZURE_RESPONSE_TIMEOUT = "azure-response-timeout"; private static final String AZURE_EAGERLY_CONVERT_HEADERS = "azure-eagerly-convert-headers"; @@ -109,24 +112,24 @@ public Mono send(HttpRequest request, Context context) { Objects.requireNonNull(request.getUrl(), "'request.getUrl()' cannot be null."); Objects.requireNonNull(request.getUrl().getProtocol(), "'request.getUrl().getProtocol()' cannot be null."); - boolean effectiveEagerlyReadResponse = (boolean) context.getData(AZURE_EAGERLY_READ_RESPONSE).orElse(false); - long effectiveResponseTimeout = context.getData(AZURE_RESPONSE_TIMEOUT) + boolean eagerlyReadResponse = (boolean) context.getData(AZURE_EAGERLY_READ_RESPONSE).orElse(false); + boolean ignoreResponseBody = (boolean) context.getData(AZURE_IGNORE_RESPONSE_BODY).orElse(false); + boolean headersEagerlyConverted = (boolean) context.getData(AZURE_EAGERLY_CONVERT_HEADERS).orElse(false); + long responseTimeout = context.getData(AZURE_RESPONSE_TIMEOUT) .filter(timeoutDuration -> timeoutDuration instanceof Duration) .map(timeoutDuration -> ((Duration) timeoutDuration).toMillis()) .orElse(this.responseTimeout); - boolean effectiveHeadersEagerlyConverted = (boolean) context.getData(AZURE_EAGERLY_CONVERT_HEADERS) - .orElse(false); return nettyClient .doOnRequest((r, connection) -> addRequestHandlers(connection, context)) - .doAfterRequest((r, connection) -> doAfterRequest(connection, effectiveResponseTimeout)) + .doAfterRequest((r, connection) -> doAfterRequest(connection, responseTimeout)) .doOnResponse((response, connection) -> addReadTimeoutHandler(connection, readTimeout)) .doAfterResponseSuccess((response, connection) -> removeReadTimeoutHandler(connection)) .request(HttpMethod.valueOf(request.getHttpMethod().toString())) .uri(request.getUrl().toString()) .send(bodySendDelegate(request)) - .responseConnection(responseDelegate(request, disableBufferCopy, effectiveEagerlyReadResponse, - effectiveHeadersEagerlyConverted)) + .responseConnection(responseDelegate(request, disableBufferCopy, eagerlyReadResponse, ignoreResponseBody, + headersEagerlyConverted)) .single() .onErrorMap(throwable -> { // The exception was an SSLException that was caused by a failure to connect to a proxy. @@ -247,14 +250,32 @@ private static NettyOutbound sendInputStream(NettyOutbound reactorNettyOutbound, * @param restRequest the Rest request whose response this delegate handles * @param disableBufferCopy Flag indicating if the network response shouldn't be buffered. * @param eagerlyReadResponse Flag indicating if the network response should be eagerly read into memory. + * @param ignoreResponseBody Flag indicating if the network response should be ignored. * @param headersEagerlyConverted Flag indicating if the Netty HttpHeaders should be eagerly converted to Azure Core * HttpHeaders. * @return a delegate upon invocation setup Rest response object */ private static BiFunction> responseDelegate( - HttpRequest restRequest, boolean disableBufferCopy, boolean eagerlyReadResponse, + HttpRequest restRequest, boolean disableBufferCopy, boolean eagerlyReadResponse, boolean ignoreResponseBody, boolean headersEagerlyConverted) { return (reactorNettyResponse, reactorNettyConnection) -> { + // Ignoring the response body takes precedent over eagerly reading the response body. + // Both should never be true at the same time but this is acts as a safeguard. + if (ignoreResponseBody) { + AtomicBoolean firstNext = new AtomicBoolean(true); + return reactorNettyConnection.inbound().receive() + .doOnNext(ignored -> { + if (!firstNext.compareAndSet(true, false)) { + LOGGER.log(LogLevel.WARNING, () -> "Received HTTP response body when one wasn't expected. " + + "Response body will be ignored as directed."); + } + }) + .ignoreElements() + .doFinally(ignored -> closeConnection(reactorNettyConnection)) + .then(Mono.fromSupplier(() -> new NettyAsyncHttpBufferedResponse(reactorNettyResponse, restRequest, + EMPTY_BYTES, headersEagerlyConverted))); + } + /* * If the response is being eagerly read into memory the flag for buffer copying can be ignored as the * response MUST be deeply copied to ensure it can safely be used downstream. diff --git a/sdk/core/azure-core-http-okhttp/src/main/java/com/azure/core/http/okhttp/OkHttpAsyncHttpClient.java b/sdk/core/azure-core-http-okhttp/src/main/java/com/azure/core/http/okhttp/OkHttpAsyncHttpClient.java index 5cdccc181d4d..2a28c7454752 100644 --- a/sdk/core/azure-core-http-okhttp/src/main/java/com/azure/core/http/okhttp/OkHttpAsyncHttpClient.java +++ b/sdk/core/azure-core-http-okhttp/src/main/java/com/azure/core/http/okhttp/OkHttpAsyncHttpClient.java @@ -27,6 +27,7 @@ import com.azure.core.util.Contexts; import com.azure.core.util.ProgressReporter; import com.azure.core.util.logging.ClientLogger; +import com.azure.core.util.logging.LogLevel; import okhttp3.Call; import okhttp3.MediaType; import okhttp3.OkHttpClient; @@ -39,7 +40,6 @@ import java.io.IOException; import java.io.UncheckedIOException; -import java.util.Objects; /** * HttpClient implementation for OkHttp. @@ -47,9 +47,11 @@ class OkHttpAsyncHttpClient implements HttpClient { private static final ClientLogger LOGGER = new ClientLogger(OkHttpAsyncHttpClient.class); - private static final RequestBody EMPTY_REQUEST_BODY = RequestBody.create(new byte[0]); + private static final byte[] EMPTY_BODY = new byte[0]; + private static final RequestBody EMPTY_REQUEST_BODY = RequestBody.create(EMPTY_BODY); private static final String AZURE_EAGERLY_READ_RESPONSE = "azure-eagerly-read-response"; + private static final String AZURE_IGNORE_RESPONSE_BODY = "azure-ignore-response-body"; private static final String AZURE_EAGERLY_CONVERT_HEADERS = "azure-eagerly-convert-headers"; final OkHttpClient httpClient; @@ -66,6 +68,7 @@ public Mono send(HttpRequest request) { @Override public Mono send(HttpRequest request, Context context) { boolean eagerlyReadResponse = (boolean) context.getData(AZURE_EAGERLY_READ_RESPONSE).orElse(false); + boolean ignoreResponseBody = (boolean) context.getData(AZURE_IGNORE_RESPONSE_BODY).orElse(false); boolean eagerlyConvertHeaders = (boolean) context.getData(AZURE_EAGERLY_CONVERT_HEADERS).orElse(false); ProgressReporter progressReporter = Contexts.with(context).getHttpRequestProgressReporter(); @@ -87,7 +90,8 @@ public Mono send(HttpRequest request, Context context) { .subscribe(okHttpRequest -> { try { Call call = httpClient.newCall(okHttpRequest); - call.enqueue(new OkHttpCallback(sink, request, eagerlyReadResponse, eagerlyConvertHeaders)); + call.enqueue(new OkHttpCallback(sink, request, eagerlyReadResponse, ignoreResponseBody, + eagerlyConvertHeaders)); sink.onCancel(call::cancel); } catch (Exception ex) { sink.error(ex); @@ -99,6 +103,7 @@ public Mono send(HttpRequest request, Context context) { @Override public HttpResponse sendSync(HttpRequest request, Context context) { boolean eagerlyReadResponse = (boolean) context.getData(AZURE_EAGERLY_READ_RESPONSE).orElse(false); + boolean ignoreResponseBody = (boolean) context.getData(AZURE_IGNORE_RESPONSE_BODY).orElse(false); boolean eagerlyConvertHeaders = (boolean) context.getData(AZURE_EAGERLY_CONVERT_HEADERS).orElse(false); ProgressReporter progressReporter = Contexts.with(context).getHttpRequestProgressReporter(); @@ -106,7 +111,8 @@ public HttpResponse sendSync(HttpRequest request, Context context) { Request okHttpRequest = toOkHttpRequest(request, progressReporter); try { Response okHttpResponse = httpClient.newCall(okHttpRequest).execute(); - return toHttpResponse(request, okHttpResponse, eagerlyReadResponse, eagerlyConvertHeaders); + return toHttpResponse(request, okHttpResponse, eagerlyReadResponse, ignoreResponseBody, + eagerlyConvertHeaders); } catch (IOException e) { throw LOGGER.logExceptionAsError(new UncheckedIOException(e)); } @@ -198,20 +204,30 @@ private static long getRequestContentLength(BinaryDataContent content, HttpHeade } private static HttpResponse toHttpResponse(HttpRequest request, okhttp3.Response response, - boolean eagerlyReadResponse, boolean eagerlyConvertHeaders) throws IOException { + boolean eagerlyReadResponse, boolean ignoreResponseBody, boolean eagerlyConvertHeaders) throws IOException { + // Ignoring the response body takes precedent over eagerly reading the response body. + // Both should never be true at the same time but this is acts as a safeguard. + if (ignoreResponseBody) { + ResponseBody body = response.body(); + if (body != null) { + if (body.contentLength() > 0) { + LOGGER.log(LogLevel.WARNING, () -> "Received HTTP response body when one wasn't expected. " + + "Response body will be ignored as directed."); + } + body.close(); + } + + return new OkHttpAsyncBufferedResponse(response, request, EMPTY_BODY, eagerlyConvertHeaders); + } + /* * Use a buffered response when we are eagerly reading the response from the network and the body isn't * empty. */ if (eagerlyReadResponse) { try (ResponseBody body = response.body()) { - if (Objects.nonNull(body)) { - byte[] bytes = body.bytes(); - return new OkHttpAsyncBufferedResponse(response, request, bytes, eagerlyConvertHeaders); - } else { - // Body is null, use the non-buffering response. - return new OkHttpAsyncResponse(response, request, eagerlyConvertHeaders); - } + byte[] bytes = (body != null) ? body.bytes() : EMPTY_BODY; + return new OkHttpAsyncBufferedResponse(response, request, bytes, eagerlyConvertHeaders); } } else { return new OkHttpAsyncResponse(response, request, eagerlyConvertHeaders); @@ -222,13 +238,15 @@ private static class OkHttpCallback implements okhttp3.Callback { private final MonoSink sink; private final HttpRequest request; private final boolean eagerlyReadResponse; + private final boolean ignoreResponseBody; private final boolean eagerlyConvertHeaders; OkHttpCallback(MonoSink sink, HttpRequest request, boolean eagerlyReadResponse, - boolean eagerlyConvertHeaders) { + boolean ignoreResponseBody, boolean eagerlyConvertHeaders) { this.sink = sink; this.request = request; this.eagerlyReadResponse = eagerlyReadResponse; + this.ignoreResponseBody = ignoreResponseBody; this.eagerlyConvertHeaders = eagerlyConvertHeaders; } @@ -248,7 +266,8 @@ public void onFailure(okhttp3.Call call, IOException e) { @Override public void onResponse(okhttp3.Call call, okhttp3.Response response) { try { - sink.success(toHttpResponse(request, response, eagerlyReadResponse, eagerlyConvertHeaders)); + sink.success(toHttpResponse(request, response, eagerlyReadResponse, ignoreResponseBody, + eagerlyConvertHeaders)); } catch (IOException ex) { // Reading the body bytes may cause an IOException, if it happens propagate it. sink.error(ex); diff --git a/sdk/core/azure-core-test/src/main/java/com/azure/core/test/http/MockHttpClient.java b/sdk/core/azure-core-test/src/main/java/com/azure/core/test/http/MockHttpClient.java index 32b09ced5536..e2ce5b37b8ea 100644 --- a/sdk/core/azure-core-test/src/main/java/com/azure/core/test/http/MockHttpClient.java +++ b/sdk/core/azure-core-test/src/main/java/com/azure/core/test/http/MockHttpClient.java @@ -181,6 +181,8 @@ public Mono send(HttpRequest request) { final String statusCodeString = requestPathLower.substring("/status/".length()); final int statusCode = Integer.parseInt(statusCodeString); response = new MockHttpResponse(request, statusCode); + } else if (requestPathLower.startsWith("/voideagerreadoom")) { + response = new MockHttpResponse(request, 200); } } else if ("echo.org".equalsIgnoreCase(requestHost)) { return FluxUtil.collectBytesInByteBufferStream(request.getBody()) diff --git a/sdk/core/azure-core-test/src/main/java/com/azure/core/test/implementation/RestProxyTests.java b/sdk/core/azure-core-test/src/main/java/com/azure/core/test/implementation/RestProxyTests.java index 686c847b0ff3..97b1de142812 100644 --- a/sdk/core/azure-core-test/src/main/java/com/azure/core/test/implementation/RestProxyTests.java +++ b/sdk/core/azure-core-test/src/main/java/com/azure/core/test/implementation/RestProxyTests.java @@ -76,6 +76,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.function.Consumer; import java.util.stream.Stream; import static org.junit.jupiter.api.Assertions.assertArrayEquals; @@ -2035,6 +2036,61 @@ public void requestOptionsSetsAHeader() { assertEquals("randomValue2", response.getHeaderValue("randomHeader")); } + @Host("http://localhost") + @ServiceInterface(name = "Service28") + interface Service28 { + @Head("voideagerreadoom") + @ExpectedResponses({200}) + void headvoid(); + + @Head("voideagerreadoom") + @ExpectedResponses({200}) + Void headVoid(); + + @Head("voideagerreadoom") + @ExpectedResponses({200}) + Response headResponseVoid(); + + @Head("voideagerreadoom") + @ExpectedResponses({200}) + ResponseBase headResponseBaseVoid(); + + @Head("voideagerreadoom") + @ExpectedResponses({200}) + Mono headMonoVoid(); + + @Head("voideagerreadoom") + @ExpectedResponses({200}) + Mono> headMonoResponseVoid(); + + @Head("voideagerreadoom") + @ExpectedResponses({200}) + Mono> headMonoResponseBaseVoid(); + + @Head("voideagerreadoom") + @ExpectedResponses({200}) + Flux headFluxVoid(); + } + + @ParameterizedTest + @MethodSource("voidDoesNotEagerlyReadResponseSupplier") + public void voidDoesNotEagerlyReadResponse(Consumer executable) { + assertDoesNotThrow(() -> executable.accept(createService(Service28.class))); + } + + private static Stream> voidDoesNotEagerlyReadResponseSupplier() { + return Stream.of( + Service28::headvoid, + Service28::headVoid, + Service28::headResponseVoid, + Service28::headResponseBaseVoid, + Service28::headMonoVoid, + Service28::headMonoResponseVoid, + Service28::headMonoResponseBaseVoid, + Service28::headFluxVoid + ); + } + // Helpers protected T createService(Class serviceClass) { final HttpClient httpClient = createHttpClient(); diff --git a/sdk/core/azure-core-test/src/test/java/com/azure/core/test/RestProxyTestsWireMockServer.java b/sdk/core/azure-core-test/src/test/java/com/azure/core/test/RestProxyTestsWireMockServer.java index ae44b3ca223f..b3802238c7de 100644 --- a/sdk/core/azure-core-test/src/test/java/com/azure/core/test/RestProxyTestsWireMockServer.java +++ b/sdk/core/azure-core-test/src/test/java/com/azure/core/test/RestProxyTestsWireMockServer.java @@ -33,6 +33,7 @@ import java.util.Random; import java.util.stream.Collectors; +import static com.github.tomakehurst.wiremock.client.WireMock.aResponse; import static com.github.tomakehurst.wiremock.client.WireMock.delete; import static com.github.tomakehurst.wiremock.client.WireMock.get; import static com.github.tomakehurst.wiremock.client.WireMock.head; @@ -69,6 +70,16 @@ public static WireMockServer getRestProxyTestsServer() { server.stubFor(patch(urlPathMatching("/patch"))); server.stubFor(get("/get")); + // Validates a bug where a void, or Void, response type would previously attempt to eagerly read the response + // body. This resulted in OutOfMemoryErrors or high memory usage in APIs such as the getProperties on Blobs, + // Datalake, and Files where the size of the resource is the Content-Length header value. So, there could be + // an attempt to create a byte[] large enough to hold the response. + // + // This uses a size too large for a byte[], so if the incorrect handling is used an OutOfMemoryError will be + // thrown. + server.stubFor(head(urlPathMatching("/voideagerreadoom")).willReturn(aResponse() + .withHeader("Content-Length", "10737418240"))); + return server; } diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/RestProxyBase.java b/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/RestProxyBase.java index 3c5a173d857e..aeb963d5fb0a 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/RestProxyBase.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/RestProxyBase.java @@ -94,6 +94,10 @@ public final Object invoke(Object proxy, final Method method, RequestOptions opt context = context.addData("azure-eagerly-read-response", true); } + if (methodParser.isResponseBodyIgnored()) { + context = context.addData("azure-ignore-response-body", true); + } + if (methodParser.isHeadersEagerlyConverted()) { context = context.addData("azure-eagerly-convert-headers", true); } diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/SwaggerMethodParser.java b/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/SwaggerMethodParser.java index 4ee08e5672e2..9b5223235254 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/SwaggerMethodParser.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/implementation/http/rest/SwaggerMethodParser.java @@ -102,6 +102,7 @@ public class SwaggerMethodParser implements HttpResponseDecodeData { private final boolean isStreamResponse; private final boolean returnTypeDecodeable; private final boolean responseEagerlyRead; + private final boolean ignoreResponseBody; private final boolean headersEagerlyConverted; private final String spanName; @@ -278,6 +279,7 @@ public SwaggerMethodParser(Method swaggerMethod) { Type unwrappedReturnType = unwrapReturnType(returnType); this.returnTypeDecodeable = isReturnTypeDecodeable(unwrappedReturnType); this.responseEagerlyRead = isResponseEagerlyRead(unwrappedReturnType); + this.ignoreResponseBody = isResponseBodyIgnored(unwrappedReturnType); this.spanName = interfaceParser.getServiceName() + "." + swaggerMethod.getName(); } @@ -708,6 +710,11 @@ public boolean isResponseEagerlyRead() { return responseEagerlyRead; } + @Override + public boolean isResponseBodyIgnored() { + return ignoreResponseBody; + } + @Override public boolean isHeadersEagerlyConverted() { return headersEagerlyConverted; @@ -722,7 +729,7 @@ public String getSpanName() { return spanName; } - static boolean isReturnTypeDecodeable(Type unwrappedReturnType) { + public static boolean isReturnTypeDecodeable(Type unwrappedReturnType) { if (unwrappedReturnType == null) { return false; } @@ -735,17 +742,24 @@ static boolean isReturnTypeDecodeable(Type unwrappedReturnType) { && !TypeUtil.isTypeOrSubTypeOf(unwrappedReturnType, Void.class); } - static boolean isResponseEagerlyRead(Type unwrappedReturnType) { + public static boolean isResponseBodyIgnored(Type unwrappedReturnType) { if (unwrappedReturnType == null) { return false; } - return isReturnTypeDecodeable(unwrappedReturnType) - || TypeUtil.isTypeOrSubTypeOf(unwrappedReturnType, Void.TYPE) + return TypeUtil.isTypeOrSubTypeOf(unwrappedReturnType, Void.TYPE) || TypeUtil.isTypeOrSubTypeOf(unwrappedReturnType, Void.class); } - static Type unwrapReturnType(Type returnType) { + public static boolean isResponseEagerlyRead(Type unwrappedReturnType) { + if (unwrappedReturnType == null) { + return false; + } + + return isReturnTypeDecodeable(unwrappedReturnType); + } + + public static Type unwrapReturnType(Type returnType) { if (returnType == null) { return null; } diff --git a/sdk/core/azure-core/src/main/java/com/azure/core/implementation/serializer/HttpResponseDecodeData.java b/sdk/core/azure-core/src/main/java/com/azure/core/implementation/serializer/HttpResponseDecodeData.java index f18d4ca86d5c..76b4bfe1ce67 100644 --- a/sdk/core/azure-core/src/main/java/com/azure/core/implementation/serializer/HttpResponseDecodeData.java +++ b/sdk/core/azure-core/src/main/java/com/azure/core/implementation/serializer/HttpResponseDecodeData.java @@ -9,6 +9,7 @@ import com.azure.core.http.rest.ResponseBase; import com.azure.core.implementation.TypeUtil; import com.azure.core.implementation.http.UnexpectedExceptionInformation; +import com.azure.core.implementation.http.rest.SwaggerMethodParser; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -104,7 +105,9 @@ default UnexpectedExceptionInformation getUnexpectedException(int code) { * * @return Whether the return type is decode-able. */ - boolean isReturnTypeDecodeable(); + default boolean isReturnTypeDecodeable() { + return SwaggerMethodParser.isReturnTypeDecodeable(SwaggerMethodParser.unwrapReturnType(getReturnType())); + } /** * Whether the network response body should be eagerly read based on its {@link #getReturnType() returnType}. @@ -115,6 +118,8 @@ default UnexpectedExceptionInformation getUnexpectedException(int code) { *

  • byte[]
  • *
  • ByteBuffer
  • *
  • InputStream
  • + *
  • Void
  • + *
  • void
  • * * * Reactive, {@link Mono} and {@link Flux}, and Response, {@link Response} and {@link ResponseBase}, generics are @@ -122,7 +127,28 @@ default UnexpectedExceptionInformation getUnexpectedException(int code) { * * @return Whether the network response body should be eagerly read. */ - boolean isResponseEagerlyRead(); + default boolean isResponseEagerlyRead() { + return SwaggerMethodParser.isResponseEagerlyRead(SwaggerMethodParser.unwrapReturnType(getReturnType())); + } + + /** + * Whether the network response body will be ignored based on its {@link #getReturnType() returnType}. + *

    + * The following types, including subtypes, ignored the network response body: + *

      + *
    • Void
    • + *
    • void
    • + *
    + * + * Reactive, {@link Mono} and {@link Flux}, and Response, {@link Response} and {@link ResponseBase}, generics are + * cracked open and their generic types are inspected for being one of the types above. + * + * @return Whether the network response body will be ignored. + */ + default boolean isResponseBodyIgnored() { + return SwaggerMethodParser.isResponseBodyIgnored(SwaggerMethodParser.unwrapReturnType(getReturnType())); + + } /** * Whether the return type contains strongly-typed headers. diff --git a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/RestProxyTests.java b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/RestProxyTests.java index 9f201a6a1933..5994389ba98e 100644 --- a/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/RestProxyTests.java +++ b/sdk/core/azure-core/src/test/java/com/azure/core/http/rest/RestProxyTests.java @@ -227,7 +227,7 @@ public void voidReturningApiClosesResponse() { } @Test - public void voidReturningApiEagerlyReadsResponse() { + public void voidReturningApiIgnoresResponseBody() { LocalHttpClient client = new LocalHttpClient(); HttpPipeline pipeline = new HttpPipelineBuilder() .httpClient(client) @@ -237,12 +237,13 @@ public void voidReturningApiEagerlyReadsResponse() { testInterface.testVoidMethod(); - assertTrue(client.lastContext.getData("azure-eagerly-read-response").isPresent()); - assertTrue((Boolean) client.lastContext.getData("azure-eagerly-read-response").get()); + assertFalse(client.lastContext.getData("azure-eagerly-read-response").isPresent()); + assertTrue(client.lastContext.getData("azure-ignore-response-body").isPresent()); + assertTrue((boolean) client.lastContext.getData("azure-ignore-response-body").get()); } @Test - public void monoVoidReturningApiEagerlyReadsResponse() { + public void monoVoidReturningApiIgnoresResponseBody() { LocalHttpClient client = new LocalHttpClient(); HttpPipeline pipeline = new HttpPipelineBuilder() .httpClient(client) @@ -253,12 +254,13 @@ public void monoVoidReturningApiEagerlyReadsResponse() { testInterface.testMethodReturnsMonoVoid()) .verifyComplete(); - assertTrue(client.lastContext.getData("azure-eagerly-read-response").isPresent()); - assertTrue((Boolean) client.lastContext.getData("azure-eagerly-read-response").get()); + assertFalse(client.lastContext.getData("azure-eagerly-read-response").isPresent()); + assertTrue(client.lastContext.getData("azure-ignore-response-body").isPresent()); + assertTrue((boolean) client.lastContext.getData("azure-ignore-response-body").get()); } @Test - public void monoResponseVoidReturningApiEagerlyReadsResponse() { + public void monoResponseVoidReturningApiIgnoresResponseBody() { LocalHttpClient client = new LocalHttpClient(); HttpPipeline pipeline = new HttpPipelineBuilder() .httpClient(client) @@ -270,12 +272,13 @@ public void monoResponseVoidReturningApiEagerlyReadsResponse() { .expectNextCount(1) .verifyComplete(); - assertTrue(client.lastContext.getData("azure-eagerly-read-response").isPresent()); - assertTrue((Boolean) client.lastContext.getData("azure-eagerly-read-response").get()); + assertFalse(client.lastContext.getData("azure-eagerly-read-response").isPresent()); + assertTrue(client.lastContext.getData("azure-ignore-response-body").isPresent()); + assertTrue((boolean) client.lastContext.getData("azure-ignore-response-body").get()); } @Test - public void responseVoidReturningApiEagerlyReadsResponse() { + public void responseVoidReturningApiIgnoresResponseBody() { LocalHttpClient client = new LocalHttpClient(); HttpPipeline pipeline = new HttpPipelineBuilder() .httpClient(client) @@ -285,8 +288,9 @@ public void responseVoidReturningApiEagerlyReadsResponse() { TestInterface testInterface = RestProxy.create(TestInterface.class, pipeline); testInterface.testMethodReturnsResponseVoid(); - assertTrue(client.lastContext.getData("azure-eagerly-read-response").isPresent()); - assertTrue((Boolean) client.lastContext.getData("azure-eagerly-read-response").get()); + assertFalse(client.lastContext.getData("azure-eagerly-read-response").isPresent()); + assertTrue(client.lastContext.getData("azure-ignore-response-body").isPresent()); + assertTrue((boolean) client.lastContext.getData("azure-ignore-response-body").get()); } @Test diff --git a/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/ResponseConstructorsCacheBenchMarkTestData.java b/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/ResponseConstructorsCacheBenchMarkTestData.java index 8b996350a8bc..845161f3e633 100644 --- a/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/ResponseConstructorsCacheBenchMarkTestData.java +++ b/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/ResponseConstructorsCacheBenchMarkTestData.java @@ -267,8 +267,6 @@ private static byte[] asJsonByteArray(Object object) { class Input { private final Type returnType; - private final boolean returnTypeDecodeable; - private final boolean responseEagerlyRead; private final HttpResponseDecoder.HttpDecodedResponse decodedResponse; private final Object bodyAsObject; @@ -278,9 +276,6 @@ class Input { Mono httpResponse, Object bodyAsObject) { this.returnType = findMethod(serviceClass, methodName).getGenericReturnType(); - Type unwrappedReturnType = SwaggerMethodParser.unwrapReturnType(returnType); - this.returnTypeDecodeable = SwaggerMethodParser.isReturnTypeDecodeable(unwrappedReturnType); - this.responseEagerlyRead = SwaggerMethodParser.isResponseEagerlyRead(unwrappedReturnType); this.decodedResponse = decoder.decode(httpResponse, new HttpResponseDecodeData() { @Override public Type getReturnType() { @@ -292,16 +287,6 @@ public boolean isExpectedResponseStatusCode(int statusCode) { return false; } - @Override - public boolean isReturnTypeDecodeable() { - return returnTypeDecodeable; - } - - @Override - public boolean isResponseEagerlyRead() { - return responseEagerlyRead; - } - @Override public boolean isHeadersEagerlyConverted() { return false; diff --git a/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/SwaggerMethodParserTests.java b/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/SwaggerMethodParserTests.java index 5a810bd8856b..121ccf9e2c7e 100644 --- a/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/SwaggerMethodParserTests.java +++ b/sdk/core/azure-core/src/test/java/com/azure/core/implementation/http/rest/SwaggerMethodParserTests.java @@ -35,6 +35,7 @@ import com.azure.core.http.rest.Response; import com.azure.core.http.rest.ResponseBase; import com.azure.core.http.rest.SimpleResponse; +import com.azure.core.http.rest.StreamResponse; import com.azure.core.implementation.TypeUtil; import com.azure.core.models.JsonPatchDocument; import com.azure.core.util.Base64Url; @@ -692,7 +693,7 @@ public void isReturnTypeDecodable(Type returnType, boolean expected) { } private static Stream isReturnTypeDecodeableSupplier() { - return returnTypeSupplierForDecodeableAndEagerReading(false); + return returnTypeSupplierForDecodeableAndEagerReading(true, false, false); } @ParameterizedTest @@ -703,143 +704,156 @@ public void isResponseEagerlyRead(Type returnType, boolean expected) { } private static Stream isResponseEagerlyReadSupplier() { - return returnTypeSupplierForDecodeableAndEagerReading(true); + return returnTypeSupplierForDecodeableAndEagerReading(true, false, false); } - private static Stream returnTypeSupplierForDecodeableAndEagerReading(boolean voidTypeStatus) { + @ParameterizedTest + @MethodSource("isResponseBodyIgnoredSupplier") + public void isResponseBodyIgnored(Type returnType, boolean expected) { + Type unwrappedReturnType = SwaggerMethodParser.unwrapReturnType(returnType); + assertEquals(expected, SwaggerMethodParser.isResponseBodyIgnored(unwrappedReturnType)); + } + + private static Stream isResponseBodyIgnoredSupplier() { + return returnTypeSupplierForDecodeableAndEagerReading(false, false, true); + } + + private static Stream returnTypeSupplierForDecodeableAndEagerReading(boolean nonBinaryTypeStatus, + boolean binaryTypeStatus, boolean voidTypeStatus) { return Stream.of( // Unknown response type can't be determined to be decode-able. Arguments.of(null, false), // BinaryData, Byte arrays, ByteBuffers, InputStream, and voids aren't decode-able. - Arguments.of(BinaryData.class, false), + Arguments.of(BinaryData.class, binaryTypeStatus), - Arguments.of(byte[].class, false), + Arguments.of(byte[].class, binaryTypeStatus), // Both ByteBuffer and sub-types shouldn't be decode-able. - Arguments.of(ByteBuffer.class, false), - Arguments.of(MappedByteBuffer.class, false), + Arguments.of(ByteBuffer.class, binaryTypeStatus), + Arguments.of(MappedByteBuffer.class, binaryTypeStatus), // Both InputSteam and sub-types shouldn't be decode-able. - Arguments.of(InputStream.class, false), - Arguments.of(FileInputStream.class, false), + Arguments.of(InputStream.class, binaryTypeStatus), + Arguments.of(FileInputStream.class, binaryTypeStatus), Arguments.of(void.class, voidTypeStatus), Arguments.of(Void.class, voidTypeStatus), Arguments.of(Void.TYPE, voidTypeStatus), // Other POJO types are decode-able. - Arguments.of(JsonPatchDocument.class, true), + Arguments.of(JsonPatchDocument.class, nonBinaryTypeStatus), // In addition to the direct types, reactive and Response generic types should be handled. // Reactive generics. // Mono generics. - Arguments.of(createParameterizedMono(BinaryData.class), false), - Arguments.of(createParameterizedMono(byte[].class), false), - Arguments.of(createParameterizedMono(ByteBuffer.class), false), - Arguments.of(createParameterizedMono(MappedByteBuffer.class), false), - Arguments.of(createParameterizedMono(InputStream.class), false), - Arguments.of(createParameterizedMono(FileInputStream.class), false), + Arguments.of(createParameterizedMono(BinaryData.class), binaryTypeStatus), + Arguments.of(createParameterizedMono(byte[].class), binaryTypeStatus), + Arguments.of(createParameterizedMono(ByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedMono(MappedByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedMono(InputStream.class), binaryTypeStatus), + Arguments.of(createParameterizedMono(FileInputStream.class), binaryTypeStatus), Arguments.of(createParameterizedMono(void.class), voidTypeStatus), Arguments.of(createParameterizedMono(Void.class), voidTypeStatus), Arguments.of(createParameterizedMono(Void.TYPE), voidTypeStatus), - Arguments.of(createParameterizedMono(JsonPatchDocument.class), true), + Arguments.of(createParameterizedMono(JsonPatchDocument.class), nonBinaryTypeStatus), // Flux generics. - Arguments.of(createParameterizedFlux(BinaryData.class), false), - Arguments.of(createParameterizedFlux(byte[].class), false), - Arguments.of(createParameterizedFlux(ByteBuffer.class), false), - Arguments.of(createParameterizedFlux(MappedByteBuffer.class), false), - Arguments.of(createParameterizedFlux(InputStream.class), false), - Arguments.of(createParameterizedFlux(FileInputStream.class), false), + Arguments.of(createParameterizedFlux(BinaryData.class), binaryTypeStatus), + Arguments.of(createParameterizedFlux(byte[].class), binaryTypeStatus), + Arguments.of(createParameterizedFlux(ByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedFlux(MappedByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedFlux(InputStream.class), binaryTypeStatus), + Arguments.of(createParameterizedFlux(FileInputStream.class), binaryTypeStatus), Arguments.of(createParameterizedFlux(void.class), voidTypeStatus), Arguments.of(createParameterizedFlux(Void.class), voidTypeStatus), Arguments.of(createParameterizedFlux(Void.TYPE), voidTypeStatus), - Arguments.of(createParameterizedFlux(JsonPatchDocument.class), true), + Arguments.of(createParameterizedFlux(JsonPatchDocument.class), nonBinaryTypeStatus), // Response generics. // If the raw type is Response it should check the first, and only, generic type. - Arguments.of(createParameterizedResponse(BinaryData.class), false), - Arguments.of(createParameterizedResponse(byte[].class), false), - Arguments.of(createParameterizedResponse(ByteBuffer.class), false), - Arguments.of(createParameterizedResponse(MappedByteBuffer.class), false), - Arguments.of(createParameterizedResponse(InputStream.class), false), - Arguments.of(createParameterizedResponse(FileInputStream.class), false), + Arguments.of(createParameterizedResponse(BinaryData.class), binaryTypeStatus), + Arguments.of(createParameterizedResponse(byte[].class), binaryTypeStatus), + Arguments.of(createParameterizedResponse(ByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedResponse(MappedByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedResponse(InputStream.class), binaryTypeStatus), + Arguments.of(createParameterizedResponse(FileInputStream.class), binaryTypeStatus), Arguments.of(createParameterizedResponse(void.class), voidTypeStatus), Arguments.of(createParameterizedResponse(Void.class), voidTypeStatus), Arguments.of(createParameterizedResponse(Void.TYPE), voidTypeStatus), - Arguments.of(createParameterizedResponse(JsonPatchDocument.class), true), + Arguments.of(createParameterizedResponse(JsonPatchDocument.class), nonBinaryTypeStatus), // If the raw type is ResponseBase it should check the second generic type, the first is deserialized // headers. - Arguments.of(createParameterizedResponseBase(BinaryData.class), false), - Arguments.of(createParameterizedResponseBase(byte[].class), false), - Arguments.of(createParameterizedResponseBase(ByteBuffer.class), false), - Arguments.of(createParameterizedResponseBase(MappedByteBuffer.class), false), - Arguments.of(createParameterizedResponseBase(InputStream.class), false), - Arguments.of(createParameterizedResponseBase(FileInputStream.class), false), + Arguments.of(createParameterizedResponseBase(BinaryData.class), binaryTypeStatus), + Arguments.of(createParameterizedResponseBase(byte[].class), binaryTypeStatus), + Arguments.of(createParameterizedResponseBase(ByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedResponseBase(MappedByteBuffer.class), binaryTypeStatus), + Arguments.of(createParameterizedResponseBase(InputStream.class), binaryTypeStatus), + Arguments.of(createParameterizedResponseBase(FileInputStream.class), binaryTypeStatus), Arguments.of(createParameterizedResponseBase(void.class), voidTypeStatus), Arguments.of(createParameterizedResponseBase(Void.class), voidTypeStatus), Arguments.of(createParameterizedResponseBase(Void.TYPE), voidTypeStatus), - Arguments.of(createParameterizedResponseBase(JsonPatchDocument.class), true), + Arguments.of(createParameterizedResponseBase(JsonPatchDocument.class), nonBinaryTypeStatus), // Reactive generics containing response generics. // Mono of Response - Arguments.of(createParameterizedMono(createParameterizedResponse(BinaryData.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponse(byte[].class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponse(ByteBuffer.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponse(MappedByteBuffer.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponse(InputStream.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponse(FileInputStream.class)), false), + Arguments.of(createParameterizedMono(createParameterizedResponse(BinaryData.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponse(byte[].class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponse(ByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponse(MappedByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponse(InputStream.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponse(FileInputStream.class)), binaryTypeStatus), Arguments.of(createParameterizedMono(createParameterizedResponse(void.class)), voidTypeStatus), Arguments.of(createParameterizedMono(createParameterizedResponse(Void.class)), voidTypeStatus), Arguments.of(createParameterizedMono(createParameterizedResponse(Void.TYPE)), voidTypeStatus), - Arguments.of(createParameterizedMono(createParameterizedResponse(JsonPatchDocument.class)), true), + Arguments.of(createParameterizedMono(createParameterizedResponse(JsonPatchDocument.class)), nonBinaryTypeStatus), // Mono of ResponseBase - Arguments.of(createParameterizedMono(createParameterizedResponseBase(BinaryData.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponseBase(byte[].class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponseBase(ByteBuffer.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponseBase(MappedByteBuffer.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponseBase(InputStream.class)), false), - Arguments.of(createParameterizedMono(createParameterizedResponseBase(FileInputStream.class)), false), + Arguments.of(createParameterizedMono(createParameterizedResponseBase(BinaryData.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponseBase(byte[].class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponseBase(ByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponseBase(MappedByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponseBase(InputStream.class)), binaryTypeStatus), + Arguments.of(createParameterizedMono(createParameterizedResponseBase(FileInputStream.class)), binaryTypeStatus), Arguments.of(createParameterizedMono(createParameterizedResponseBase(void.class)), voidTypeStatus), Arguments.of(createParameterizedMono(createParameterizedResponseBase(Void.class)), voidTypeStatus), Arguments.of(createParameterizedMono(createParameterizedResponseBase(Void.TYPE)), voidTypeStatus), - Arguments.of(createParameterizedMono(createParameterizedResponseBase(JsonPatchDocument.class)), true), + Arguments.of(createParameterizedMono(createParameterizedResponseBase(JsonPatchDocument.class)), nonBinaryTypeStatus), // Flux of Response - Arguments.of(createParameterizedFlux(createParameterizedResponse(BinaryData.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponse(byte[].class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponse(ByteBuffer.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponse(MappedByteBuffer.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponse(InputStream.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponse(FileInputStream.class)), false), + Arguments.of(createParameterizedFlux(createParameterizedResponse(BinaryData.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponse(byte[].class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponse(ByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponse(MappedByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponse(InputStream.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponse(FileInputStream.class)), binaryTypeStatus), Arguments.of(createParameterizedFlux(createParameterizedResponse(void.class)), voidTypeStatus), Arguments.of(createParameterizedFlux(createParameterizedResponse(Void.class)), voidTypeStatus), Arguments.of(createParameterizedFlux(createParameterizedResponse(Void.TYPE)), voidTypeStatus), - Arguments.of(createParameterizedFlux(createParameterizedResponse(JsonPatchDocument.class)), true), + Arguments.of(createParameterizedFlux(createParameterizedResponse(JsonPatchDocument.class)), nonBinaryTypeStatus), // Flux of ResponseBase - Arguments.of(createParameterizedFlux(createParameterizedResponseBase(BinaryData.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponseBase(byte[].class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponseBase(ByteBuffer.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponseBase(MappedByteBuffer.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponseBase(InputStream.class)), false), - Arguments.of(createParameterizedFlux(createParameterizedResponseBase(FileInputStream.class)), false), + Arguments.of(createParameterizedFlux(createParameterizedResponseBase(BinaryData.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponseBase(byte[].class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponseBase(ByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponseBase(MappedByteBuffer.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponseBase(InputStream.class)), binaryTypeStatus), + Arguments.of(createParameterizedFlux(createParameterizedResponseBase(FileInputStream.class)), binaryTypeStatus), Arguments.of(createParameterizedFlux(createParameterizedResponseBase(void.class)), voidTypeStatus), Arguments.of(createParameterizedFlux(createParameterizedResponseBase(Void.class)), voidTypeStatus), Arguments.of(createParameterizedFlux(createParameterizedResponseBase(Void.TYPE)), voidTypeStatus), - Arguments.of(createParameterizedFlux(createParameterizedResponseBase(JsonPatchDocument.class)), true), + Arguments.of(createParameterizedFlux(createParameterizedResponseBase(JsonPatchDocument.class)), nonBinaryTypeStatus), // Custom implementations of Response and ResponseBase. Arguments.of(VoidResponse.class, voidTypeStatus), - Arguments.of(StringResponse.class, true), + Arguments.of(StringResponse.class, nonBinaryTypeStatus), + Arguments.of(StreamResponse.class, binaryTypeStatus), Arguments.of(VoidResponseWithDeserializedHeaders.class, voidTypeStatus), - Arguments.of(StringResponseWithDeserializedHeaders.class, true) + Arguments.of(StringResponseWithDeserializedHeaders.class, nonBinaryTypeStatus) ); }