From a9aeff3ae12b922037de6d45861cea1ab91636c7 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 25 Apr 2017 11:51:24 -0400 Subject: [PATCH] Fix HTTP Reactive support according latest SF --- .../outbound/ReactiveHttpRequestExecutingMessageHandler.java | 3 ++- .../springframework/integration/http/dsl/HttpDslTests.java | 2 +- .../ReactiveHttpRequestExecutingMessageHandlerTests.java | 4 ++-- 3 files changed, 5 insertions(+), 4 deletions(-) diff --git a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandler.java b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandler.java index 6d2733aea2..91f816b5c8 100644 --- a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandler.java +++ b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandler.java @@ -33,6 +33,7 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandler; import org.springframework.util.Assert; import org.springframework.web.reactive.function.BodyExtractors; +import org.springframework.web.reactive.function.BodyInserters; import org.springframework.web.reactive.function.client.ClientResponse; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.web.reactive.function.client.WebClientException; @@ -120,7 +121,7 @@ public class ReactiveHttpRequestExecutingMessageHandler extends AbstractHttpRequ .headers(httpRequest.getHeaders()); if (httpRequest.hasBody()) { - requestSpec.body(httpRequest.getBody()); + requestSpec.body(BodyInserters.fromObject(httpRequest.getBody())); } Mono responseMono = requestSpec.exchange() diff --git a/spring-integration-http/src/test/java/org/springframework/integration/http/dsl/HttpDslTests.java b/spring-integration-http/src/test/java/org/springframework/integration/http/dsl/HttpDslTests.java index 84b2603363..c43b371f3f 100644 --- a/spring-integration-http/src/test/java/org/springframework/integration/http/dsl/HttpDslTests.java +++ b/spring-integration-http/src/test/java/org/springframework/integration/http/dsl/HttpDslTests.java @@ -121,7 +121,7 @@ public class HttpDslTests { response.getHeaders().setContentType(MediaType.TEXT_PLAIN); return response.writeWith(Mono.just(response.bufferFactory().wrap("FOO".getBytes()))) - .then(response::setComplete); + .then(Mono.defer(response::setComplete)); }); WebClient webClient = WebClient.builder() diff --git a/spring-integration-http/src/test/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandlerTests.java b/spring-integration-http/src/test/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandlerTests.java index e1f2dc4c36..dde40bbf77 100644 --- a/spring-integration-http/src/test/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandlerTests.java +++ b/spring-integration-http/src/test/java/org/springframework/integration/http/outbound/ReactiveHttpRequestExecutingMessageHandlerTests.java @@ -52,7 +52,7 @@ public class ReactiveHttpRequestExecutingMessageHandlerTests { ClientHttpConnector httpConnector = new HttpHandlerConnector((request, response) -> { response.setStatusCode(HttpStatus.OK); return Mono.empty() - .then(response::setComplete); + .then(Mono.defer(response::setComplete)); }); WebClient webClient = WebClient.builder() @@ -79,7 +79,7 @@ public class ReactiveHttpRequestExecutingMessageHandlerTests { ClientHttpConnector httpConnector = new HttpHandlerConnector((request, response) -> { response.setStatusCode(HttpStatus.UNAUTHORIZED); return Mono.empty() - .then(response::setComplete); + .then(Mono.defer(response::setComplete)); }); WebClient webClient = WebClient.builder()