From 17df347b1b38870ebebec3b0bb4a7396e88aa9e9 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Tue, 2 Jun 2015 12:39:24 +0100 Subject: [PATCH] Add localName property to channel specs (read only) /channels/taps responds to the local name as well as the "channel" name. Also includes a change to allow the default channel name to be output_topic: or input_topic: (for pubsub but no particular name). And also adds autoStartup to module properties. --- .../bus/runner/adapter/InputChannelSpec.java | 6 +++++ .../bus/runner/adapter/MessageBusAdapter.java | 10 +++++++ .../runner/config/LifecycleConfiguration.java | 5 +++- .../MessageBusAdapterConfiguration.java | 27 ++++++++++++++----- .../runner/config/MessageBusProperties.java | 10 +++++++ .../MessageBusAdapterConfigurationTests.java | 9 +++++++ 6 files changed, 60 insertions(+), 7 deletions(-) diff --git a/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/InputChannelSpec.java b/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/InputChannelSpec.java index 1c5503fc7..949a77129 100644 --- a/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/InputChannelSpec.java +++ b/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/InputChannelSpec.java @@ -16,6 +16,7 @@ package org.springframework.bus.runner.adapter; +import org.springframework.integration.support.context.NamedComponent; import org.springframework.messaging.MessageChannel; import com.fasterxml.jackson.annotation.JsonIgnore; @@ -38,6 +39,11 @@ public class InputChannelSpec { return name; } + public String getLocalName() { + return (channel instanceof NamedComponent) ? ((NamedComponent) channel) + .getComponentName() : channel.toString(); + } + @JsonIgnore public MessageChannel getMessageChannel() { return channel; diff --git a/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/MessageBusAdapter.java b/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/MessageBusAdapter.java index 5e77bcc99..ab6ec1664 100644 --- a/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/MessageBusAdapter.java +++ b/spring-bus-core/src/main/java/org/springframework/bus/runner/adapter/MessageBusAdapter.java @@ -129,6 +129,11 @@ public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { return spec; } } + for (OutputChannelSpec spec : outputChannels) { + if (name.equals(spec.getLocalName())) { + return spec; + } + } return null; } @@ -141,6 +146,11 @@ public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { return spec; } } + for (InputChannelSpec spec : inputChannels) { + if (name.equals(spec.getLocalName())) { + return spec; + } + } return null; } diff --git a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/LifecycleConfiguration.java b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/LifecycleConfiguration.java index b97b8734f..3248febeb 100644 --- a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/LifecycleConfiguration.java +++ b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/LifecycleConfiguration.java @@ -28,12 +28,15 @@ import org.springframework.context.annotation.Configuration; @Configuration public class LifecycleConfiguration implements CommandLineRunner { + @Autowired + private MessageBusProperties module; + @Autowired private MessageBusAdapter adapter; @Override public void run(String... args) throws Exception { - if (!adapter.isRunning()) { + if (!adapter.isRunning() && module.isAutoStartup()) { adapter.start(); } } 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 12ed8dd73..6c30a0852 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,8 @@ import org.springframework.messaging.MessageChannel; import org.springframework.xd.dirt.integration.bus.MessageBus; import org.springframework.xd.dirt.integration.bus.MessageBusAwareRouterBeanPostProcessor; +import reactor.util.StringUtils; + /** * @author Dave Syer * @@ -68,11 +70,12 @@ public class MessageBusAdapterConfiguration { String[] names = BeanFactoryUtils.beanNamesForTypeIncludingAncestors(beanFactory, MessageChannel.class); for (String name : names) { - String channelName = extractChannelName("output", name, module.getOutputChannelName()); + String channelName = extractChannelName("output", name, + module.getOutputChannelName()); if (channelName != null) { OutputChannelSpec channel = new OutputChannelSpec(channelName, beanFactory.getBean(name, MessageChannel.class)); - String tapChannelName = !channelName.equals(module.getOutputChannelName()) ? module + String tapChannelName = !isDefaultOuputChannel(channelName) ? module .getTapChannelName(getPlainChannelName(channel.getName())) : module.getTapChannelName(); channel.setTapChannelName(tapChannelName); @@ -83,7 +86,16 @@ public class MessageBusAdapterConfiguration { return channels; } - private String extractChannelName(String start, String name, String externalChannelName) { + private boolean isDefaultOuputChannel(String channelName) { + if (channelName.contains(":")) { + String[] tokens = channelName.split(":", 2); + channelName = tokens[1]; + } + return channelName.equals(module.getOutputChannelName()); + } + + private String extractChannelName(String start, String name, + String externalChannelName) { if (name.equals(start)) { return externalChannelName; } @@ -95,10 +107,12 @@ public class MessageBusAdapterConfiguration { String type = tokens[0]; if ("queue".equals(type)) { // omit the type for a queue - prefix = tokens[1] + "."; + if (StringUtils.hasText(tokens[1])) { + prefix = tokens[1] + "."; + } } else { - prefix = channelName + "."; + prefix = channelName + (channelName.endsWith(":") ? "" : "."); } } else { @@ -121,7 +135,8 @@ public class MessageBusAdapterConfiguration { String[] names = BeanFactoryUtils.beanNamesForTypeIncludingAncestors(beanFactory, MessageChannel.class); for (String name : names) { - String channelName = extractChannelName("input", name, module.getInputChannelName()); + String channelName = extractChannelName("input", name, + module.getInputChannelName()); if (channelName != null) { channels.add(new InputChannelSpec(channelName, beanFactory.getBean(name, MessageChannel.class))); 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 c4b39432c..204bda274 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 @@ -51,6 +51,8 @@ public class MessageBusProperties { private Tap tap; + private boolean autoStartup = true; + public String getName() { return name; } @@ -130,6 +132,14 @@ public class MessageBusProperties { this.producerProperties = producerProperties; } + public boolean isAutoStartup() { + return autoStartup; + } + + public void setAutoStartup(boolean autoStartup) { + this.autoStartup = autoStartup; + } + public Tap getTap() { return tap; } 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 ceba27562..2e41abcf2 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 @@ -66,6 +66,15 @@ public class MessageBusAdapterConfigurationTests { assertEquals("tap:stream:group.module.0", channels.iterator().next().getTapChannelName()); } + @Test + public void oneOutputTopic() throws Exception { + context.registerSingleton("output.topic:", new DirectChannel()); + Collection channels = configuration.getOutputChannels(); + assertEquals(1, channels.size()); + assertEquals("topic:group.0", channels.iterator().next().getName()); + assertEquals("tap:stream:group.module.0", channels.iterator().next().getTapChannelName()); + } + @Test public void twoOutputsWithTopic() throws Exception { context.registerSingleton("output", new DirectChannel());