diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java index 75826951ec..76598bbf57 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/EnricherSpec.java @@ -21,6 +21,7 @@ import java.util.Map; import java.util.function.Function; import org.springframework.expression.Expression; +import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.config.ConsumerEndpointFactoryBean; import org.springframework.integration.expression.FunctionExpression; import org.springframework.integration.expression.ValueExpression; @@ -40,6 +41,7 @@ import reactor.util.function.Tuple2; * * @author Artem Bilan * @author Tim Ysewyn + * @author Ian Bondoc * * @since 5.0 */ @@ -136,6 +138,23 @@ public class EnricherSpec extends ConsumerEndpointSpec receive = this.subFlowTestReplyChannel.receive(5000); + assertNotNull(receive); + assertEquals("Foo Bar (Reply Producing)", receive.getHeaders().get("foo")); + Object payload = receive.getPayload(); + assertThat(payload, instanceOf(TestPojo.class)); + TestPojo result = (TestPojo) payload; + assertThat(result.getName(), is("Foo Bar (Reply Producing)")); + + this.terminatingSubFlowEnricherInput.send(MessageBuilder.withPayload(new TestPojo("Bar")).build()); + receive = this.subFlowTestReplyChannel.receive(5000); + assertNotNull(receive); + assertEquals("Foo Bar (Terminating)", receive.getHeaders().get("foo")); + payload = receive.getPayload(); + assertThat(payload, instanceOf(TestPojo.class)); + result = (TestPojo) payload; + assertThat(result.getName(), is("Foo Bar (Terminating)")); + } + @Autowired @Qualifier("encodingFlow.input") private MessageChannel encodingFlowInput; @@ -287,6 +323,43 @@ public class TransformerTests { return idempotentReceiverInterceptor; } + @Bean + public PollableChannel subFlowTestReplyChannel() { + return new QueueChannel(); + } + + @Bean + public IntegrationFlow replyProducingSubFlowEnricher(SomeService someService) { + return f -> f + .enrich(e -> e.requestPayload(p -> p.getPayload().getName()) + .requestSubFlow(sf -> sf + .handle((p, h) -> someService.someServiceMethod(p))) + .headerFunction("foo", Message::getPayload) + .propertyFunction("name", Message::getPayload)) + .channel("subFlowTestReplyChannel"); + } + + @Bean + public MessageChannel enricherReplyChannel() { + return MessageChannels.direct().get(); + } + + @Bean + public IntegrationFlow terminatingSubFlowEnricher(SomeService someService) { + return f -> f + .enrich(e -> e.requestPayload(p -> p.getPayload().getName()) + .requestSubFlow(sf -> sf + .handle(someService::aTerminatingServiceMethod)) + .replyChannel("enricherReplyChannel") + .headerFunction("foo", Message::getPayload) + .propertyFunction("name", Message::getPayload)) + .channel("subFlowTestReplyChannel"); + } + + @Bean + public SomeService someService() { + return new SomeService(); + } } @@ -354,4 +427,21 @@ public class TransformerTests { } + public static class SomeService { + + @Autowired + @Qualifier("enricherReplyChannel") + public MessageChannel enricherReplyChannel; + + public String someServiceMethod(String value) { + return "Foo ".concat(value).concat(" (Reply Producing)"); + } + + public void aTerminatingServiceMethod(Message message) { + String payload = "Foo ".concat(message.getPayload().toString()).concat(" (Terminating)"); + enricherReplyChannel.send(MessageBuilder.withPayload(payload).copyHeaders(message.getHeaders()).build()); + } + + } + }