From f28819bbae1cd5bd67d1a050d03aa4ac4a09e45a Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 9 Aug 2021 09:49:51 -0400 Subject: [PATCH] Fix generics for customizeMonoReply() Related to: https://stackoverflow.com/questions/68637283/how-to-customize-response-in-spring-integration-using-webflux-when-a-specific-er The function provided for the `ConsumerEndpointSpec.customizeMonoReply()` may convert incoming value to something else. With wildcards it cannot be compiled without casting. * Add `` generic arg for the `customizeMonoReply()` to conform in and out types carrying. * Modify `WebFluxDslTests` to demonstrate the problem and confirm the fix **Cherry-pick to `5.4.x` & `5.3.x`** --- .../integration/dsl/ConsumerEndpointSpec.java | 8 ++++++-- .../integration/webflux/dsl/WebFluxDslTests.java | 9 +++++++-- 2 files changed, 13 insertions(+), 4 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java index bccdf98ec1..c04bfcd15e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/ConsumerEndpointSpec.java @@ -201,12 +201,16 @@ public abstract class ConsumerEndpointSpec, /** * Specify a {@link BiFunction} for customizing {@link Mono} replies via {@link ReactiveRequestHandlerAdvice}. * @param replyCustomizer the {@link BiFunction} to propagate into {@link ReactiveRequestHandlerAdvice}. + * @param inbound reply payload. + * @param outbound reply payload. * @return the spec. * @since 5.3 * @see ReactiveRequestHandlerAdvice */ - public S customizeMonoReply(BiFunction, Mono, Publisher> replyCustomizer) { - return advice(new ReactiveRequestHandlerAdvice(replyCustomizer)); + @SuppressWarnings({ "unchecked", "rawtypes"}) + public S customizeMonoReply(BiFunction, Mono, Publisher> replyCustomizer) { + return advice(new ReactiveRequestHandlerAdvice( + (BiFunction, Mono, Publisher>) (BiFunction) replyCustomizer)); } /** 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 0339ebb661..775866f48d 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 @@ -91,6 +91,7 @@ import org.springframework.web.context.WebApplicationContext; import org.springframework.web.reactive.config.EnableWebFlux; import org.springframework.web.reactive.config.WebFluxConfigurer; import org.springframework.web.reactive.function.client.WebClient; +import org.springframework.web.reactive.function.client.WebClientResponseException; import org.springframework.web.util.UriComponentsBuilder; import reactor.core.publisher.Flux; @@ -418,8 +419,12 @@ public class WebFluxDslTests { .id("webFluxWithReplyPayloadToFlux") .customizeMonoReply( (message, mono) -> - mono.timeout(Duration.ofMillis(100)) - .retry())); + mono. + timeout(Duration.ofMillis(100)) + .retry() + .onErrorResume( + WebClientResponseException.NotFound.class, + ex -> Mono.just("Not Found")))); } @Bean