From 9f77fcd763703aac24e7662d88c0c42e2958c618 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 6 Nov 2014 16:20:34 +0200 Subject: [PATCH] INT-3550: RoutingSlip Improvements JIRA: https://jira.spring.io/browse/INT-3550 * Make `RoutingSlipRouteStrategy#getNextPath` as `Object` return type. to allow to produce `MessageChannel` result, not only `beanName` * Make `RoutingSlipHeaderValueMessageProcessor` ctor to accept `Object... routingSlipPath` instead of just String. It is useful from JavaConfig, when we can use `RoutingSlipRouteStrategy` `@Bean` reference. * Add JavaConfig test case to demonstrate how `RoutingSlipRouteStrategy` can get deal with inline `FixedSubscriberChannel` and Lambdas together with Reactor Streams. INT-3550: Polishing according PR comments * Fix `AbstractMessageProducingHandler` to check if `nextPath` isn't empty String * Add `RoutingSlipHeaderValueMessageProcessor` ctor check for the `routingSlipPath` entries types * Add `RoutingSlipRouteStrategy` JavaDocs regarding the loop of strategy invocation * Add `Process Manager` doc Doc Polishing. --- .../AbstractMessageProducingHandler.java | 4 +- ...ionEvaluatingRoutingSlipRouteStrategy.java | 2 +- .../routingslip/RoutingSlipRouteStrategy.java | 5 +- ...outingSlipHeaderValueMessageProcessor.java | 44 ++++-- .../routingslip/RoutingSlipTests-context.xml | 2 + .../routingslip/RoutingSlipTests.java | 143 +++++++++++++++++- src/reference/docbook/router.xml | 41 +++++ 7 files changed, 222 insertions(+), 19 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java index f92b3f3153..577c485fd2 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java @@ -196,8 +196,8 @@ public abstract class AbstractMessageProducingHandler extends AbstractMessageHan return routingSlipPathValue; } else { - String nextPath = ((RoutingSlipRouteStrategy) routingSlipPathValue).getNextPath(requestMessage, reply); - if (StringUtils.hasText(nextPath)) { + Object nextPath = ((RoutingSlipRouteStrategy) routingSlipPathValue).getNextPath(requestMessage, reply); + if (nextPath != null && (!(nextPath instanceof String) || StringUtils.hasText((String) nextPath))) { return nextPath; } else { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/routingslip/ExpressionEvaluatingRoutingSlipRouteStrategy.java b/spring-integration-core/src/main/java/org/springframework/integration/routingslip/ExpressionEvaluatingRoutingSlipRouteStrategy.java index c743e40d53..05fb2baf65 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/routingslip/ExpressionEvaluatingRoutingSlipRouteStrategy.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/routingslip/ExpressionEvaluatingRoutingSlipRouteStrategy.java @@ -79,7 +79,7 @@ public class ExpressionEvaluatingRoutingSlipRouteStrategy } @Override - public String getNextPath(Message requestMessage, Object reply) { + public Object getNextPath(Message requestMessage, Object reply) { return this.expression.getValue(this.evaluationContext, new RequestAndReply(requestMessage, reply), String.class); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/routingslip/RoutingSlipRouteStrategy.java b/spring-integration-core/src/main/java/org/springframework/integration/routingslip/RoutingSlipRouteStrategy.java index 7995550be0..95d36e13c9 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/routingslip/RoutingSlipRouteStrategy.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/routingslip/RoutingSlipRouteStrategy.java @@ -20,12 +20,15 @@ import org.springframework.messaging.Message; /** * The {@code RoutingSlip} strategy to determine the next {@code replyChannel}. + *

+ * This strategy is called repeatedly until null or an empty String is returned. * * @author Artem Bilan * @since 4.1 + * @see org.springframework.integration.handler.AbstractMessageProducingHandler */ public interface RoutingSlipRouteStrategy { - String getNextPath(Message requestMessage, Object reply); + Object getNextPath(Message requestMessage, Object reply); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/transformer/support/RoutingSlipHeaderValueMessageProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/transformer/support/RoutingSlipHeaderValueMessageProcessor.java index aace9f599f..dbf3a4f090 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/transformer/support/RoutingSlipHeaderValueMessageProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/transformer/support/RoutingSlipHeaderValueMessageProcessor.java @@ -48,7 +48,7 @@ public class RoutingSlipHeaderValueMessageProcessor extends AbstractHeaderValueMessageProcessor, Integer>> implements BeanFactoryAware { - private final List routingSlipPath; + private final List routingSlipPath; private EvaluationContext evaluationContext; @@ -56,9 +56,18 @@ public class RoutingSlipHeaderValueMessageProcessor private BeanFactory beanFactory; - public RoutingSlipHeaderValueMessageProcessor(String... routingSlipPath) { + public RoutingSlipHeaderValueMessageProcessor(Object... routingSlipPath) { Assert.notNull(routingSlipPath); Assert.noNullElements(routingSlipPath); + for (Object entry : routingSlipPath) { + if (!(entry instanceof String + || entry instanceof MessageChannel + || entry instanceof RoutingSlipRouteStrategy)) { + throw new IllegalArgumentException("The RoutingSlip can contain " + + "only bean names of MessageChannel or RoutingSlipRouteStrategy, " + + "or MessageChannel and RoutingSlipRouteStrategy instances: " + entry); + } + } this.routingSlipPath = Arrays.asList(routingSlipPath); } @@ -75,20 +84,29 @@ public class RoutingSlipHeaderValueMessageProcessor synchronized (this) { if (this.routingSlip == null) { List routingSlipValues = new ArrayList(this.routingSlipPath.size()); - for (String path : this.routingSlipPath) { - if (this.beanFactory.containsBean(path)) { - Object bean = this.beanFactory.getBean(path); - Assert.state(bean instanceof MessageChannel || bean instanceof RoutingSlipRouteStrategy, - "The RoutingSlip can contain only bean names of MessageChannel or " + - "RoutingSlipRouteStrategy: " + bean); - routingSlipValues.add(path); + for (Object path : this.routingSlipPath) { + if (path instanceof String) { + String entry = (String) path; + if (this.beanFactory.containsBean(entry)) { + Object bean = this.beanFactory.getBean(entry); + if (!(bean instanceof MessageChannel + || bean instanceof RoutingSlipRouteStrategy)) { + throw new IllegalArgumentException("The RoutingSlip can contain " + + "only bean names of MessageChannel or RoutingSlipRouteStrategy: " + bean); + } + routingSlipValues.add(entry); + } + else { + ExpressionEvaluatingRoutingSlipRouteStrategy strategy = new + ExpressionEvaluatingRoutingSlipRouteStrategy(entry); + strategy.setIntegrationEvaluationContext(this.evaluationContext); + routingSlipValues.add(strategy); + } } else { - ExpressionEvaluatingRoutingSlipRouteStrategy strategy = new - ExpressionEvaluatingRoutingSlipRouteStrategy(path); - strategy.setIntegrationEvaluationContext(this.evaluationContext); - routingSlipValues.add(strategy); + routingSlipValues.add(path); } + } this.routingSlip = Collections.singletonMap(Collections.unmodifiableList(routingSlipValues), 0); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests-context.xml index f6f82c87b3..8a6b0c0f84 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests-context.xml @@ -61,4 +61,6 @@ + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java b/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java index e33b2da099..18712a2978 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java @@ -16,11 +16,17 @@ package org.springframework.integration.routingslip; +import static org.hamcrest.Matchers.containsString; +import static org.hamcrest.Matchers.instanceOf; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; import java.util.Arrays; +import java.util.Collections; +import java.util.Date; import java.util.List; import java.util.Map; import java.util.Properties; @@ -30,27 +36,55 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; import org.springframework.integration.IntegrationMessageHeaderAccessor; +import org.springframework.integration.annotation.BridgeTo; +import org.springframework.integration.annotation.Transformer; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.channel.FixedSubscriberChannel; import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.core.MessagingTemplate; import org.springframework.integration.history.MessageHistory; import org.springframework.integration.support.MessageBuilder; +import org.springframework.integration.transformer.HeaderEnricher; +import org.springframework.integration.transformer.support.RoutingSlipHeaderValueMessageProcessor; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessagingException; import org.springframework.messaging.PollableChannel; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.support.GenericXmlContextLoader; + +import reactor.core.Environment; +import reactor.core.composable.spec.Streams; +import reactor.spring.context.config.EnableReactor; /** * @author Artem Bilan * @since 4.1 */ -@ContextConfiguration +@ContextConfiguration(loader = GenericXmlContextLoader.class) @RunWith(SpringJUnit4ClassRunner.class) +@DirtiesContext public class RoutingSlipTests { @Autowired private MessageChannel input; + @Autowired + private MessageChannel routingSlipHeaderChannel; + + @Autowired + private PollableChannel resultsChannel; + + @Autowired + private MessageChannel invalidRoutingSlipChannel; + @Test @SuppressWarnings("unchecked") public void testRoutingSlip() { @@ -76,6 +110,42 @@ public class RoutingSlipTests { } } + @Test + public void testDynamicRoutingSlipRoutStrategy() { + this.routingSlipHeaderChannel.send(new GenericMessage<>("foo")); + Message result = this.resultsChannel.receive(10000); + assertNotNull(result); + assertEquals("FOO", result.getPayload()); + + this.routingSlipHeaderChannel.send(new GenericMessage<>(2)); + result = this.resultsChannel.receive(10000); + assertNotNull(result); + assertEquals(4, result.getPayload()); + } + + @Test + public void testInvalidRoutingSlipRoutStrategy() { + try { + new RoutingSlipHeaderValueMessageProcessor(new Date()); + fail("IllegalArgumentException expected"); + } + catch (Exception e) { + assertThat(e, instanceOf(IllegalArgumentException.class)); + assertThat(e.getMessage(), + containsString("The RoutingSlip can contain " + + "only bean names of MessageChannel or RoutingSlipRouteStrategy, " + + "or MessageChannel and RoutingSlipRouteStrategy instances")); + } + try { + this.invalidRoutingSlipChannel.send(new GenericMessage<>("foo")); + fail("MessagingException expected"); + } + catch (Exception e) { + assertThat(e, instanceOf(MessagingException.class)); + assertThat(e.getMessage(), containsString("replyChannel must be a MessageChannel or String")); + } + } + public static class TestRoutingSlipRoutePojo { final String[] channels = {"channel2", "channel3"}; @@ -98,10 +168,79 @@ public class RoutingSlipTests { private AtomicBoolean invoked = new AtomicBoolean(); @Override - public String getNextPath(Message requestMessage, Object reply) { + public Object getNextPath(Message requestMessage, Object reply) { return !invoked.getAndSet(true) ? "channel4" : null; } } + @Configuration + @EnableReactor + @EnableIntegration + public static class RoutingSlipConfiguration { + + @Autowired + private Environment reactorEnv; + + @Bean + public MessagingTemplate messagingTemplate() { + return new MessagingTemplate(); + } + + @Bean + public PollableChannel resultsChannel() { + return new QueueChannel(); + } + + @Bean + public RoutingSlipRouteStrategy routeStrategy() { + return (requestMessage, reply) -> requestMessage.getPayload() instanceof String + ? new FixedSubscriberChannel(m -> + Streams.defer((String) m.getPayload()) + .env(this.reactorEnv) + .get() + .map(String::toUpperCase) + .consume(v -> messagingTemplate().convertAndSend(resultsChannel(), v)) + .flush()) + : new FixedSubscriberChannel(m -> + Streams.defer((Integer) m.getPayload()) + .env(this.reactorEnv) + .get() + .map(v -> v * 2) + .consume(v -> messagingTemplate().convertAndSend(resultsChannel(), v)) + .flush()); + } + + @Bean + public MessageChannel routingSlipHeaderChannel() { + return new DirectChannel(); + } + + @Bean + @BridgeTo + public MessageChannel processChannel() { + return new DirectChannel(); + } + + @Bean + @Transformer(inputChannel = "routingSlipHeaderChannel", outputChannel = "processChannel") + public HeaderEnricher headerEnricher() { + return new HeaderEnricher(Collections.singletonMap(IntegrationMessageHeaderAccessor.ROUTING_SLIP, + new RoutingSlipHeaderValueMessageProcessor(routeStrategy()))); + } + + @Bean + public MessageChannel invalidRoutingSlipChannel() { + return new DirectChannel(); + } + + @Bean + @Transformer(inputChannel = "invalidRoutingSlipChannel", outputChannel = "processChannel") + public HeaderEnricher headerEnricher2() { + return new HeaderEnricher(Collections.singletonMap(IntegrationMessageHeaderAccessor.ROUTING_SLIP, + new RoutingSlipHeaderValueMessageProcessor((RoutingSlipRouteStrategy) (message, r) -> new Date()))); + } + + } + } diff --git a/src/reference/docbook/router.xml b/src/reference/docbook/router.xml index 42ef5da66f..2ea6df48d5 100644 --- a/src/reference/docbook/router.xml +++ b/src/reference/docbook/router.xml @@ -1196,5 +1196,46 @@ public HeaderEnricher headerEnricher() { +
+ Process Manager Enterprise Integration Pattern + + The EIP also defines the + Process Manager pattern. + This pattern can now easily be implemented using custom Process Manager logic + encapsulated in a + RoutingSlipRouteStrategy within the routing slip. + In addition to a bean name, the RoutingSlipRouteStrategy can return any + MessageChannel object; and there is no requirement that this + MessageChannel instance is a bean in the application context. + This way, we can provide powerful dynamic routing logic, when there is no prediction which + channel should be used; a MessageChannel can be created + within the RoutingSlipRouteStrategy and returned. A + FixedSubscriberChannel with an associated MessageHandler + implementation is good combination for such cases. For example we can route to a + Reactor Stream: + + requestMessage.getPayload() instanceof String + ? new FixedSubscriberChannel(m -> + Streams.defer((String) m.getPayload()) + .env(this.reactorEnv) + .get() + .map(String::toUpperCase) + .consume(v -> messagingTemplate().convertAndSend(resultsChannel(), v)) + .flush()) + : new FixedSubscriberChannel(m -> + Streams.defer((Integer) m.getPayload()) + .env(this.reactorEnv) + .get() + .map(v -> v * 2) + .consume(v -> messagingTemplate().convertAndSend(resultsChannel(), v)) + .flush()); +}]]> +