|
17 | 17 | import java.util.ArrayList; |
18 | 18 | import java.util.List; |
19 | 19 | import java.util.Optional; |
| 20 | +import java.util.concurrent.CancellationException; |
| 21 | +import java.util.concurrent.CompletableFuture; |
| 22 | +import java.util.concurrent.CompletionException; |
20 | 23 | import java.util.concurrent.Flow; |
21 | 24 | import java.util.concurrent.Flow.Publisher; |
22 | 25 |
|
@@ -201,16 +204,45 @@ static <T> Flux<T> cancel(Publisher<List<ByteBuffer>> body) { |
201 | 204 | * Such a body must be subscribed to, or the connection it is read from is never |
202 | 205 | * released. Should the exchange be cancelled once the response has arrived but before |
203 | 206 | * its body could be subscribed to, the response is discarded, and its body cancelled. |
| 207 | + * |
| 208 | + * <p> |
| 209 | + * Cancelling the exchange before the response has arrived aborts the request. The |
| 210 | + * {@link HttpClient} then fails its future with a {@link CompletionException} |
| 211 | + * wrapping a {@link CancellationException}, which {@link Mono#fromFuture} does not |
| 212 | + * recognise as the outcome of its own cancellation and reports as a dropped error. |
| 213 | + * Only this method can cancel the future, so such a failure is always the expected |
| 214 | + * outcome of cancelling, and is ignored. |
204 | 215 | * @param httpClient the client to send the request with |
205 | 216 | * @param request the request to send |
206 | 217 | */ |
207 | 218 | static Mono<HttpResponse<Publisher<List<ByteBuffer>>>> sendAsync(HttpClient httpClient, HttpRequest request) { |
208 | | - return Mono.fromFuture(() -> httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher())) |
209 | | - .doOnDiscard(HttpResponse.class, response -> { |
210 | | - if (response.body() instanceof Publisher<?> body) { |
211 | | - cancelBody(body); |
| 219 | + return Mono.<HttpResponse<Publisher<List<ByteBuffer>>>>create(sink -> { |
| 220 | + CompletableFuture<HttpResponse<Publisher<List<ByteBuffer>>>> exchange = httpClient.sendAsync(request, |
| 221 | + HttpResponse.BodyHandlers.ofPublisher()); |
| 222 | + sink.onCancel(() -> exchange.cancel(true)); |
| 223 | + exchange.whenComplete((response, error) -> { |
| 224 | + if (error == null) { |
| 225 | + // Emit the response so the body can be consumed. |
| 226 | + // If the surrounding Mono was cancelled though and due to a race |
| 227 | + // the headers were already parsed, the below call will simply |
| 228 | + // discard the response. |
| 229 | + sink.success(response); |
| 230 | + return; |
| 231 | + } |
| 232 | + Throwable cause = error instanceof CompletionException && error.getCause() != null ? error.getCause() |
| 233 | + : error; |
| 234 | + if (cause instanceof CancellationException) { |
| 235 | + sink.success(); |
| 236 | + } |
| 237 | + else { |
| 238 | + sink.error(cause); |
212 | 239 | } |
213 | 240 | }); |
| 241 | + }).doOnDiscard(HttpResponse.class, response -> { |
| 242 | + if (response.body() instanceof Publisher<?> body) { |
| 243 | + cancelBody(body); |
| 244 | + } |
| 245 | + }); |
214 | 246 | } |
215 | 247 |
|
216 | 248 | private static void cancelBody(Publisher<?> body) { |
|
0 commit comments