diff --git a/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/dsl/WebFluxMessageHandlerSpec.java b/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/dsl/WebFluxMessageHandlerSpec.java index 8a6755bc3b..cf4dcedf11 100644 --- a/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/dsl/WebFluxMessageHandlerSpec.java +++ b/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/dsl/WebFluxMessageHandlerSpec.java @@ -35,6 +35,7 @@ import reactor.core.publisher.Mono; * * @author Shiliang Li * @author Artem Bilan + * @author Abhijit Sarkar * * @since 5.0 * @@ -64,22 +65,26 @@ public class WebFluxMessageHandlerSpec * Defaults to {@code false} - simple value is pushed downstream. * Makes sense when {@code expectedResponseType} is configured. * @param replyPayloadToFlux represent reply payload as a {@link Flux} or as a value from the {@link Mono}. + * @return the spec * @since 5.0.1 * @see WebFluxRequestExecutingMessageHandler#setReplyPayloadToFlux(boolean) */ - public void replyPayloadToFlux(boolean replyPayloadToFlux) { + public WebFluxMessageHandlerSpec replyPayloadToFlux(boolean replyPayloadToFlux) { this.target.setReplyPayloadToFlux(replyPayloadToFlux); + return this; } /** * Specify a {@link BodyExtractor} as an alternative to the {@code expectedResponseType} * to allow to get low-level access to the received {@link ClientHttpResponse}. * @param bodyExtractor the {@link BodyExtractor} to use. + * @return the spec * @since 5.0.1 * @see WebFluxRequestExecutingMessageHandler#setBodyExtractor(BodyExtractor) */ - public void bodyExtractor(BodyExtractor bodyExtractor) { + public WebFluxMessageHandlerSpec bodyExtractor(BodyExtractor bodyExtractor) { this.target.setBodyExtractor(bodyExtractor); + return this; } @Override diff --git a/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/dsl/WebFluxDslTests.java b/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/dsl/WebFluxDslTests.java index f634e762a9..5d34dc5b68 100644 --- a/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/dsl/WebFluxDslTests.java +++ b/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/dsl/WebFluxDslTests.java @@ -26,6 +26,8 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers. import java.util.Collections; +import javax.annotation.Resource; + import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -33,19 +35,24 @@ import org.reactivestreams.Publisher; import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.ResolvableType; +import org.springframework.core.io.buffer.DataBufferFactory; import org.springframework.http.HttpMethod; import org.springframework.http.HttpStatus; import org.springframework.http.MediaType; import org.springframework.http.client.reactive.ClientHttpConnector; +import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.http.dsl.Http; +import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.webflux.outbound.WebFluxRequestExecutingMessageHandler; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; import org.springframework.messaging.PollableChannel; import org.springframework.security.access.AccessDecisionManager; import org.springframework.security.access.vote.AffirmativeBased; @@ -76,6 +83,7 @@ import reactor.test.StepVerifier; /** * @author Artem Bilan * @author Shiliang Li + * @author Abhijit Sarkar * * @since 5.0 */ @@ -88,7 +96,15 @@ public class WebFluxDslTests { private WebApplicationContext wac; @Autowired - private WebFluxRequestExecutingMessageHandler serviceInternalReactiveGatewayHandler; + @Qualifier("webFluxWithReplyPayloadToFlux.handler") + private WebFluxRequestExecutingMessageHandler webFluxWithReplyPayloadToFlux; + + @Resource(name = "org.springframework.integration.webflux.outbound.WebFluxRequestExecutingMessageHandler#1") + private WebFluxRequestExecutingMessageHandler httpReactiveProxyFlow; + + @Autowired + @Qualifier("webFluxFlowWithReplyPayloadToFlux.input") + private MessageChannel webFluxFlowWithReplyPayloadToFluxInput; private MockMvc mockMvc; @@ -106,6 +122,48 @@ public class WebFluxDslTests { .build(); } + @Test + public void testWebFluxFlowWithReplyPayloadToFlux() { + ClientHttpConnector httpConnector = new HttpHandlerConnector((request, response) -> { + response.setStatusCode(HttpStatus.OK); + response.getHeaders().setContentType(MediaType.TEXT_PLAIN); + + DataBufferFactory bufferFactory = response.bufferFactory(); + return response.writeWith( + Flux.just(bufferFactory.wrap("FOO".getBytes()), + bufferFactory.wrap("BAR".getBytes()))) + .then(Mono.defer(response::setComplete)); + }); + + WebClient webClient = WebClient.builder() + .clientConnector(httpConnector) + .build(); + + new DirectFieldAccessor(this.webFluxWithReplyPayloadToFlux) + .setPropertyValue("webClient", webClient); + + QueueChannel replyChannel = new QueueChannel(); + + Message testMessage = + MessageBuilder.withPayload("test") + .setReplyChannel(replyChannel) + .build(); + + this.webFluxFlowWithReplyPayloadToFluxInput.send(testMessage); + + Message receive = replyChannel.receive(10_000); + + assertNotNull(receive); + assertThat(receive.getPayload(), instanceOf(Flux.class)); + + @SuppressWarnings("unchecked") + Flux response = (Flux) receive.getPayload(); + + StepVerifier.create(response) + .expectNext("FOO", "BAR") + .verifyComplete(); + } + @Test public void testHttpReactiveProxyFlow() throws Exception { ClientHttpConnector httpConnector = new HttpHandlerConnector((request, response) -> { @@ -120,7 +178,7 @@ public class WebFluxDslTests { .clientConnector(httpConnector) .build(); - new DirectFieldAccessor(this.serviceInternalReactiveGatewayHandler) + new DirectFieldAccessor(this.httpReactiveProxyFlow) .setPropertyValue("webClient", webClient); this.mockMvc.perform( @@ -200,6 +258,16 @@ public class WebFluxDslTests { .anonymous().disable(); } + @Bean + public IntegrationFlow webFluxFlowWithReplyPayloadToFlux() { + return f -> f + .handle(WebFlux.outboundGateway("http://www.springsource.org/spring-integration") + .httpMethod(HttpMethod.GET) + .replyPayloadToFlux(true) + .expectedResponseType(String.class), + e -> e.id("webFluxWithReplyPayloadToFlux")); + } + @Bean public IntegrationFlow httpReactiveProxyFlow() { return IntegrationFlows