INT-4541: Fix Reactive MessagingGateway Errors

JIRA: https://jira.spring.io/browse/INT-4541

Add test case to reproduce.

The gateway correctly sets up the `errorChannel` header so that a downstream
`MPEH` will send exceptions back to the gateway, but the `map()` function
did not check for an error message.

Check for an error message and throw the payload so that the `onErrorResume`
will route to the error channel.

# Conflicts:
#	spring-integration-webflux/src/test/java/org/springframework/integration/webflux/dsl/WebFluxDslTests.java
This commit is contained in:
Gary Russell
2018-10-08 16:20:08 -04:00
committed by Artem Bilan
parent 7faeee849d
commit 7403e36c1e
2 changed files with 47 additions and 6 deletions

View File

@@ -623,12 +623,23 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint
this.messageCount.incrementAndGet();
}
})
.<Message<?>>map(replyMessage ->
MessageBuilder.fromMessage(replyMessage)
.setHeader(MessageHeaders.REPLY_CHANNEL, originalReplyChannelHeader)
.setHeader(MessageHeaders.ERROR_CHANNEL, originalErrorChannelHeader)
.build())
.<Message<?>>map(replyMessage -> {
if (!error && replyMessage instanceof ErrorMessage) {
ErrorMessage em = (ErrorMessage) replyMessage;
if (em.getPayload() instanceof MessagingException) {
throw (MessagingException) em.getPayload();
}
else {
throw new MessagingException(requestMessage, em.getPayload());
}
}
else {
return MessageBuilder.fromMessage(replyMessage)
.setHeader(MessageHeaders.REPLY_CHANNEL, originalReplyChannelHeader)
.setHeader(MessageHeaders.ERROR_CHANNEL, originalErrorChannelHeader)
.build();
}
})
.onErrorResume(t -> error ? Mono.error(t) : handleSendError(requestMessage, t));
});
}

View File

@@ -51,6 +51,7 @@ 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.dsl.channel.MessageChannels;
import org.springframework.integration.http.HttpHeaders;
import org.springframework.integration.http.dsl.Http;
import org.springframework.integration.support.MessageBuilder;
@@ -231,6 +232,15 @@ public class WebFluxDslTests {
}
@Test
public void testHttpReactivePostWithError() {
this.webTestClient.post().uri("/reactivePostErrors")
.attributes(basicAuthenticationCredentials("guest", "guest"))
.body(Mono.just("foo\nbar\nbaz"), String.class)
.exchange()
.expectStatus().isEqualTo(HttpStatus.BAD_GATEWAY);
}
@Test
public void testSse() {
Flux<String> responseBody =
@@ -332,6 +342,26 @@ public class WebFluxDslTests {
.get();
}
@Bean
public IntegrationFlow httpReactiveInboundGatewayFlowWithErrors() {
return IntegrationFlows
.from(WebFlux.inboundGateway("/reactivePostErrors")
.requestMapping(m -> m.methods(HttpMethod.POST))
.requestPayloadType(ResolvableType.forClassWithGenerics(Flux.class, String.class))
.errorChannel(errorFlow().getInputChannel()))
.channel(MessageChannels.flux())
.handle((p, h) -> {
throw new RuntimeException("errorTest");
})
.get();
}
@Bean
public IntegrationFlow errorFlow() {
return f -> f
.enrichHeaders(h -> h.header(HttpHeaders.STATUS_CODE, HttpStatus.BAD_GATEWAY));
}
@Bean
public IntegrationFlow sseFlow() {
return IntegrationFlows