From 147875540d7e39329474a1d164f5b19902216132 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Fri, 29 May 2015 16:07:16 +0100 Subject: [PATCH] Regularize and test tap channel names --- roadmap.md | 6 ++++-- .../bus/runner/config/MessageBusAdapterConfiguration.java | 8 +++++--- .../bus/runner/config/MessageBusProperties.java | 6 +++++- .../config/MessageBusAdapterConfigurationTests.java | 3 ++- 4 files changed, 16 insertions(+), 7 deletions(-) diff --git a/roadmap.md b/roadmap.md index 08ebfa9f9..f3360a079 100644 --- a/roadmap.md +++ b/roadmap.md @@ -69,11 +69,13 @@ To be deployable as an XD module in a "traditional" way you need `/config/*.prop ## Backlog -- [ ] Support for multiple input and output channels +- [x] Support for multiple input and output channels - [ ] Partitioning -- [ ] Support for pubsub as "primary" input/output (in addition to the existing queue semantics) +- [x] Support for pubsub as "primary" input/output (in addition to the existing queue semantics) + +- [ ] Endpoint "/messages" for module configuration metadata ("/bus" is taken by Spring Cloud) - [ ] Support for more than one `MessageBus` (e.g. local and redis) in the same app diff --git a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java index b6b1ea132..45fd4fcd4 100644 --- a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java +++ b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java @@ -34,6 +34,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.ImportResource; import org.springframework.messaging.MessageChannel; +import org.springframework.util.StringUtils; import org.springframework.xd.dirt.integration.bus.MessageBus; import org.springframework.xd.dirt.integration.bus.MessageBusAwareRouterBeanPostProcessor; @@ -100,13 +101,14 @@ public class MessageBusAdapterConfiguration { prefix = channelName + "."; } } - channelName = prefix - + getPlainChannelName(module.getOutputChannelName()); + channelName = prefix + getPlainChannelName(module.getOutputChannelName()); channel = new OutputChannelSpec(channelName, beanFactory.getBean(name, MessageChannel.class)); } if (channel != null) { - String tapChannelName = prefix + module.getTapChannelName(); + String tapChannelName = StringUtils.hasText(prefix) ? module + .getTapChannelName(getPlainChannelName(channel.getName())) + : module.getTapChannelName(); channel.setTapChannelName(tapChannelName); channel.setTapped(true); // TODO: determine when this is the case channels.add(channel); diff --git a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusProperties.java b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusProperties.java index 3cd24267f..c4b39432c 100644 --- a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusProperties.java +++ b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusProperties.java @@ -89,9 +89,13 @@ public class MessageBusProperties { } public String getTapChannelName() { + return getTapChannelName(group); + } + + public String getTapChannelName(String prefix) { Assert.isTrue(!type.equals("job"), "Job module type not supported."); // for Stream return channel name with indexed elements - return String.format("%s.%s.%s", BusUtils.constructTapPrefix(group), name, index); + return String.format("%s.%s.%s", BusUtils.constructTapPrefix(prefix), name, index); } public void setOutputChannelName(String outputChannelName) { diff --git a/spring-bus-core/src/test/java/org/springframework/bus/runner/config/MessageBusAdapterConfigurationTests.java b/spring-bus-core/src/test/java/org/springframework/bus/runner/config/MessageBusAdapterConfigurationTests.java index 7311894c4..3bc75dfd2 100644 --- a/spring-bus-core/src/test/java/org/springframework/bus/runner/config/MessageBusAdapterConfigurationTests.java +++ b/spring-bus-core/src/test/java/org/springframework/bus/runner/config/MessageBusAdapterConfigurationTests.java @@ -105,7 +105,7 @@ public class MessageBusAdapterConfigurationTests { assertEquals(1, channels.size()); assertEquals("foo.bar", channels.iterator().next().getName()); // TODO: fix this. What should it be? - // assertEquals("tap:stream:group.module.0", channels.iterator().next().getTapChannelName()); + assertEquals("tap:stream:foo.bar.module.0", channels.iterator().next().getTapChannelName()); } @Test @@ -124,6 +124,7 @@ public class MessageBusAdapterConfigurationTests { Collection channels = configuration.getOutputChannels(); assertEquals(1, channels.size()); assertEquals("topic:foo.bar", channels.iterator().next().getName()); + assertEquals("tap:stream:foo.bar.module.0", channels.iterator().next().getTapChannelName()); } @Test