diff --git a/spring-cloud-streams/pom.xml b/spring-cloud-streams/pom.xml index 029910288..6522e62f7 100644 --- a/spring-cloud-streams/pom.xml +++ b/spring-cloud-streams/pom.xml @@ -27,11 +27,6 @@ org.springframework.boot spring-boot-starter-actuator - - org.springframework.cloud - spring-cloud-starter - true - org.springframework.boot spring-boot-starter-web diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableMessageBus.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableMessageBus.java index aa669eec7..47ec0d788 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableMessageBus.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/EnableMessageBus.java @@ -24,7 +24,7 @@ import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; import org.springframework.cloud.streams.config.LifecycleConfiguration; -import org.springframework.cloud.streams.config.MessageBusAdapterConfiguration; +import org.springframework.cloud.streams.config.ChannelBindingAdapterConfiguration; import org.springframework.cloud.streams.config.RabbitServiceConfiguration; import org.springframework.cloud.streams.config.RedisServiceConfiguration; import org.springframework.context.annotation.Configuration; @@ -40,7 +40,7 @@ import org.springframework.context.annotation.Import; @Inherited @Configuration @Import({ RedisServiceConfiguration.class, RabbitServiceConfiguration.class, - MessageBusAdapterConfiguration.class, LifecycleConfiguration.class }) + ChannelBindingAdapterConfiguration.class, LifecycleConfiguration.class }) public @interface EnableMessageBus { } diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/MessageBusAdapter.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelBindingAdapter.java similarity index 91% rename from spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/MessageBusAdapter.java rename to spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelBindingAdapter.java index b57a9b3e6..65e47633a 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/MessageBusAdapter.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelBindingAdapter.java @@ -29,7 +29,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.BeansException; -import org.springframework.cloud.streams.config.MessageBusProperties; +import org.springframework.cloud.streams.config.ChannelBindingProperties; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.context.ConfigurableApplicationContext; @@ -57,9 +57,9 @@ import org.springframework.xd.dirt.integration.bus.XdHeaders; * @author Dave Syer */ @ManagedResource -public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { +public class ChannelBindingAdapter implements Lifecycle, ApplicationContextAware { - private static Logger logger = LoggerFactory.getLogger(MessageBusAdapter.class); + private static Logger logger = LoggerFactory.getLogger(ChannelBindingAdapter.class); private MessageBus messageBus; private MessageBuilderFactory messageBuilderFactory = new DefaultMessageBuilderFactory(); @@ -73,7 +73,7 @@ public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { private boolean trackHistory = false; - private MessageBusProperties module; + private ChannelBindingProperties module; private ConfigurableApplicationContext applicationContext; @@ -85,7 +85,7 @@ public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { private Map bindings = new HashMap(); - public MessageBusAdapter(MessageBusProperties module, MessageBus messageBus) { + public ChannelBindingAdapter(ChannelBindingProperties module, MessageBus messageBus) { this.module = module; this.messageBus = messageBus; this.inputChannelLocator = new DefaultChannelLocator(module); @@ -271,7 +271,7 @@ public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { MessageChannel outputChannel = this.channelResolver.resolveDestination(binding.getLocalName()); bindMessageProducer(outputChannel, name, this.module.getProducerProperties()); if (binding.isTapped()) { - String tapChannelName = getTapChannelName(name); + String tapChannelName = this.outputChannelLocator.tap(name); binding.setTapChannelName(tapChannelName); // tappableChannels.put(tapChannelName, outputChannel); // if (isTapActive(tapChannelName)) { @@ -319,29 +319,6 @@ public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { return located; } - // TODO: move this to ChannelLocator? - private String getTapChannelName(String name) { - return !isDefaultOuputChannel(name) ? this.module.getTapChannelName(getPlainChannelName(name)) - : this.module.getTapChannelName(); - } - - // TODO: move this to ChannelLocator? - private String getPlainChannelName(String name) { - if (name.contains(":")) { - name = name.substring(name.indexOf(":") + 1); - } - return name; - } - - // TODO: move this to ChannelLocator? - private boolean isDefaultOuputChannel(String channelName) { - if (channelName.contains(":")) { - String[] tokens = channelName.split(":", 2); - channelName = tokens[1]; - } - return channelName.equals(this.module.getOutputChannelName()); - } - /* * Following methods copied from parent to support the bindChannels() method above */ @@ -422,7 +399,7 @@ public class MessageBusAdapter implements Lifecycle, ApplicationContextAware { map.putAll(historyProps); map.put("thread", Thread.currentThread().getName()); history.add(map); - Message out = MessageBusAdapter.this.messageBuilderFactory.fromMessage(message) + Message out = ChannelBindingAdapter.this.messageBuilderFactory.fromMessage(message) .setHeader(XdHeaders.XD_HISTORY, history).build(); return out; } diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelLocator.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelLocator.java index 085de8b2e..3a7de7c17 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelLocator.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelLocator.java @@ -20,9 +20,10 @@ package org.springframework.cloud.streams.adapter; * @author Dave Syer * */ -// TODO: Use DestinationResolver? public interface ChannelLocator { - + String locate(String name); + String tap(String name); + } diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelsMetadata.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelsMetadata.java index 0053a3c04..70ca981c4 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelsMetadata.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/ChannelsMetadata.java @@ -19,7 +19,7 @@ package org.springframework.cloud.streams.adapter; import java.util.Collection; import java.util.Collections; -import org.springframework.cloud.streams.config.MessageBusProperties; +import org.springframework.cloud.streams.config.ChannelBindingProperties; /** * @author Dave Syer @@ -28,13 +28,13 @@ public class ChannelsMetadata { private Collection outputChannels = Collections.emptySet(); private Collection inputChannels = Collections.emptySet(); - private MessageBusProperties module; + private ChannelBindingProperties module; - public MessageBusProperties getModule() { + public ChannelBindingProperties getModule() { return this.module; } - public void setModule(MessageBusProperties module) { + public void setModule(ChannelBindingProperties module) { this.module = module; } diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/DefaultChannelLocator.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/DefaultChannelLocator.java index 90c3eda93..243e677b7 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/DefaultChannelLocator.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/adapter/DefaultChannelLocator.java @@ -16,7 +16,7 @@ package org.springframework.cloud.streams.adapter; -import org.springframework.cloud.streams.config.MessageBusProperties; +import org.springframework.cloud.streams.config.ChannelBindingProperties; import org.springframework.util.StringUtils; /** @@ -24,9 +24,9 @@ import org.springframework.util.StringUtils; */ public class DefaultChannelLocator implements ChannelLocator { - private MessageBusProperties module; + private ChannelBindingProperties module; - public DefaultChannelLocator(MessageBusProperties module) { + public DefaultChannelLocator(ChannelBindingProperties module) { this.module = module; } @@ -43,6 +43,19 @@ public class DefaultChannelLocator implements ChannelLocator { return null; } + @Override + public String tap(String name) { + return !isDefaultOuputChannel(name) ? this.module.getTapChannelName(getPlainChannelName(name)) + : this.module.getTapChannelName(); + } + + private boolean isDefaultOuputChannel(String channelName) { + if (channelName.contains(":")) { + String[] tokens = channelName.split(":", 2); + channelName = tokens[1]; + } + return channelName.equals(this.module.getOutputChannelName()); + } private String extractChannelName(String start, String name, String externalChannelName) { if (name.equals(start)) { diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/MessageBusAdapterConfiguration.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfiguration.java similarity index 83% rename from spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/MessageBusAdapterConfiguration.java rename to spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfiguration.java index 2f347f34a..1f34d235d 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/MessageBusAdapterConfiguration.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfiguration.java @@ -28,11 +28,11 @@ import org.springframework.aop.target.LazyInitTargetSource; import org.springframework.beans.factory.BeanFactoryUtils; import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.cloud.streams.adapter.ChannelBindingAdapter; import org.springframework.cloud.streams.adapter.ChannelLocator; import org.springframework.cloud.streams.adapter.Input; import org.springframework.cloud.streams.adapter.InputChannelBinding; -import org.springframework.cloud.streams.adapter.MessageBusAdapter; import org.springframework.cloud.streams.adapter.Output; import org.springframework.cloud.streams.adapter.OutputChannelBinding; import org.springframework.cloud.streams.endpoint.ChannelsEndpoint; @@ -50,11 +50,10 @@ import org.springframework.xd.dirt.integration.bus.MessageBusAwareRouterBeanPost */ @Configuration @ImportResource("classpath*:/META-INF/spring-xd/bus/codec.xml") -@EnableConfigurationProperties(MessageBusProperties.class) -public class MessageBusAdapterConfiguration { +public class ChannelBindingAdapterConfiguration { @Autowired - private MessageBusProperties module; + private ChannelBindingProperties module; @Autowired private ListableBeanFactory beanFactory; @@ -67,10 +66,12 @@ public class MessageBusAdapterConfiguration { @Output private ChannelLocator outputChannelLocator; + @Autowired + private MessageBus messageBus; + @Bean - public MessageBusAdapter messageBusAdapter(MessageBusProperties module, - MessageBus messageBus) { - MessageBusAdapter adapter = new MessageBusAdapter(module, messageBus); + public ChannelBindingAdapter messageBusAdapter() { + ChannelBindingAdapter adapter = new ChannelBindingAdapter(this.module, this.messageBus); adapter.setOutputChannels(getOutputChannels()); adapter.setInputChannels(getInputChannels()); if (this.inputChannelLocator!=null) { @@ -83,10 +84,16 @@ public class MessageBusAdapterConfiguration { } @Bean - public ChannelsEndpoint channelsEndpoint(MessageBusAdapter adapter) { + public ChannelsEndpoint channelsEndpoint(ChannelBindingAdapter adapter) { return new ChannelsEndpoint(adapter); } + public void refresh() { + ChannelBindingAdapter adapter = messageBusAdapter(); + adapter.setOutputChannels(getOutputChannels()); + adapter.setInputChannels(getInputChannels()); + } + protected Collection getOutputChannels() { Set channels = new LinkedHashSet(); String[] names = BeanFactoryUtils.beanNamesForTypeIncludingAncestors( @@ -158,4 +165,12 @@ public class MessageBusAdapterConfiguration { } + @Configuration + @ConditionalOnMissingBean(ChannelBindingProperties.class) + protected static class ModulePropertiesConfiguration { + @Bean(name="spring.cloud.channels.CONFIGURATION_PROPERTIES") + public ChannelBindingProperties moduleProperties() { + return new ChannelBindingProperties(); + } + } } diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingProperties.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingProperties.java new file mode 100644 index 000000000..c54d0b3ae --- /dev/null +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/ChannelBindingProperties.java @@ -0,0 +1,92 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.streams.config; + +import java.util.Properties; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonInclude.Include; + +/** + * @author Dave Syer + * + */ +@ConfigurationProperties("spring.cloud.channels") +@JsonInclude(Include.NON_DEFAULT) +public class ChannelBindingProperties { + + private String outputChannelName = "group.0"; + + private String inputChannelName = "group.0"; + + private Properties consumerProperties = new Properties(); + + private Properties producerProperties = new Properties(); + + private boolean autoStartup = true; + + public String getOutputChannelName() { + return this.outputChannelName; + } + + public String getInputChannelName() { + return this.inputChannelName; + } + + public void setOutputChannelName(String outputChannelName) { + this.outputChannelName = outputChannelName; + } + + public void setInputChannelName(String inputChannelName) { + this.inputChannelName = inputChannelName; + } + + public Properties getConsumerProperties() { + return this.consumerProperties; + } + + public void setConsumerProperties(Properties consumerProperties) { + this.consumerProperties = consumerProperties; + } + + public Properties getProducerProperties() { + return this.producerProperties; + } + + public void setProducerProperties(Properties producerProperties) { + this.producerProperties = producerProperties; + } + + public boolean isAutoStartup() { + return this.autoStartup; + } + + public void setAutoStartup(boolean autoStartup) { + this.autoStartup = autoStartup; + } + + public String getTapChannelName() { + return getTapChannelName(getOutputChannelName()); + } + + public String getTapChannelName(String prefix) { + return "tap:" + prefix; + } + +} diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/LifecycleConfiguration.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/LifecycleConfiguration.java index b4ac9966a..ab5ba9cd5 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/LifecycleConfiguration.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/LifecycleConfiguration.java @@ -18,7 +18,7 @@ package org.springframework.cloud.streams.config; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.CommandLineRunner; -import org.springframework.cloud.streams.adapter.MessageBusAdapter; +import org.springframework.cloud.streams.adapter.ChannelBindingAdapter; import org.springframework.context.annotation.Configuration; /** @@ -29,10 +29,10 @@ import org.springframework.context.annotation.Configuration; public class LifecycleConfiguration implements CommandLineRunner { @Autowired - private MessageBusProperties module; + private ChannelBindingProperties module; @Autowired - private MessageBusAdapter adapter; + private ChannelBindingAdapter adapter; @Override public void run(String... args) throws Exception { diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/MessageBusProperties.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/MessageBusProperties.java deleted file mode 100644 index b85714426..000000000 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/config/MessageBusProperties.java +++ /dev/null @@ -1,236 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.streams.config; - -import java.util.Properties; - -import org.springframework.boot.context.properties.ConfigurationProperties; -import org.springframework.util.Assert; -import org.springframework.xd.dirt.integration.bus.BusUtils; - -import com.fasterxml.jackson.annotation.JsonInclude; -import com.fasterxml.jackson.annotation.JsonInclude.Include; - -/** - * @author Dave Syer - * - */ -@ConfigurationProperties("spring.bus") -@JsonInclude(Include.NON_DEFAULT) -public class MessageBusProperties { - - private String name = "module"; - - private String group = "group"; - - private int index = 0; - - private String outputChannelName; - - private String inputChannelName; - - private String type = "processor"; - - private Properties consumerProperties = new Properties(); - - private Properties producerProperties = new Properties(); - - private Tap tap; - - private Discovery discovery = new Discovery(); - - private boolean autoStartup = true; - - public String getName() { - return name; - } - - public void setName(String name) { - this.name = name; - } - - public String getGroup() { - return group; - } - - public void setGroup(String group) { - this.group = group; - } - - public int getIndex() { - return index; - } - - public void setIndex(int index) { - this.index = index; - } - - public String getInputChannelName() { - if (isTap()) { - return String.format("%s.%s.%s", BusUtils.constructTapPrefix(tap.getGroup()), - tap.getName(), tap.getIndex()); - } - return (inputChannelName != null) ? inputChannelName : BusUtils - .constructPipeName(group, index > 0 ? index - 1 : index); - } - - public String getOutputChannelName() { - return (outputChannelName != null) ? outputChannelName : BusUtils - .constructPipeName(group, index); - } - - 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(prefix), name, index); - } - - public void setOutputChannelName(String outputChannelName) { - this.outputChannelName = outputChannelName; - } - - public void setInputChannelName(String inputChannelName) { - this.inputChannelName = inputChannelName; - } - - public String getType() { - return type; - } - - public void setType(String type) { - this.type = type; - } - - public Properties getConsumerProperties() { - return consumerProperties; - } - - public void setConsumerProperties(Properties consumerProperties) { - this.consumerProperties = consumerProperties; - } - - public Properties getProducerProperties() { - return producerProperties; - } - - public void setProducerProperties(Properties producerProperties) { - this.producerProperties = producerProperties; - } - - public boolean isAutoStartup() { - return autoStartup; - } - - public void setAutoStartup(boolean autoStartup) { - this.autoStartup = autoStartup; - } - - public Tap getTap() { - return tap; - } - - public void setTap(Tap tap) { - this.tap = tap; - } - - private boolean isTap() { - if (tap != null) { - Assert.state(tap.getName() != null, "Tap name not provided"); - Assert.state(!tap.getGroup().equals(group), - "Tap group cannot be the same as module group"); - } - return tap != null; - } - - public Discovery getDiscovery() { - return discovery; - } - - public static class Discovery { - - private boolean enabled = false; - - private String inputServiceId; - - private String outputServiceId; - - public boolean isEnabled() { - return enabled; - } - - public void setEnabled(boolean enabled) { - this.enabled = enabled; - } - - public String getInputServiceId() { - return inputServiceId; - } - - public void setInputServiceId(String inputServiceId) { - this.inputServiceId = inputServiceId; - } - - public String getOutputServiceId() { - return outputServiceId; - } - - public void setOutputServiceId(String outputServiceId) { - this.outputServiceId = outputServiceId; - } - - } - - public static class Tap { - - private String group = "group"; - - private String name; - - public String getName() { - return name; - } - - public void setName(String name) { - this.name = name; - } - - private int index = 0; - - public String getGroup() { - return group; - } - - public void setGroup(String group) { - this.group = group; - } - - public int getIndex() { - return index; - } - - public void setIndex(int index) { - this.index = index; - } - - } - -} diff --git a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/endpoint/ChannelsEndpoint.java b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/endpoint/ChannelsEndpoint.java index c217e4c88..14e9f1307 100644 --- a/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/endpoint/ChannelsEndpoint.java +++ b/spring-cloud-streams/src/main/java/org/springframework/cloud/streams/endpoint/ChannelsEndpoint.java @@ -23,7 +23,7 @@ import java.util.Map; import org.springframework.boot.actuate.endpoint.AbstractEndpoint; import org.springframework.cloud.streams.adapter.ChannelsMetadata; -import org.springframework.cloud.streams.adapter.MessageBusAdapter; +import org.springframework.cloud.streams.adapter.ChannelBindingAdapter; import org.springframework.cloud.streams.adapter.OutputChannelBinding; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; @@ -36,9 +36,9 @@ import org.springframework.web.bind.annotation.RestController; @RestController public class ChannelsEndpoint extends AbstractEndpoint> { - private MessageBusAdapter adapter; + private ChannelBindingAdapter adapter; - public ChannelsEndpoint(MessageBusAdapter adapter) { + public ChannelsEndpoint(ChannelBindingAdapter adapter) { super("channels"); this.adapter = adapter; } diff --git a/spring-cloud-streams/src/main/resources/META-INF/spring.factories b/spring-cloud-streams/src/main/resources/META-INF/spring.factories new file mode 100644 index 000000000..335d864cd --- /dev/null +++ b/spring-cloud-streams/src/main/resources/META-INF/spring.factories @@ -0,0 +1 @@ +org.springframework.boot.autoconfigure.EnableAutoConfiguration:\ diff --git a/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/adapter/DefaultChannelLocatorTests.java b/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/adapter/DefaultChannelLocatorTests.java index 544bb5736..90c1a50e4 100644 --- a/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/adapter/DefaultChannelLocatorTests.java +++ b/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/adapter/DefaultChannelLocatorTests.java @@ -19,7 +19,7 @@ import static org.junit.Assert.assertEquals; import org.junit.Test; import org.springframework.cloud.streams.adapter.DefaultChannelLocator; -import org.springframework.cloud.streams.config.MessageBusProperties; +import org.springframework.cloud.streams.config.ChannelBindingProperties; /** * @author Dave Syer @@ -27,7 +27,7 @@ import org.springframework.cloud.streams.config.MessageBusProperties; */ public class DefaultChannelLocatorTests { - private MessageBusProperties module = new MessageBusProperties(); + private ChannelBindingProperties module = new ChannelBindingProperties(); private DefaultChannelLocator locator = new DefaultChannelLocator(this.module); diff --git a/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfigurationTests.java b/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfigurationTests.java new file mode 100644 index 000000000..e59c1e228 --- /dev/null +++ b/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/config/ChannelBindingAdapterConfigurationTests.java @@ -0,0 +1,150 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.streams.config; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; + +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.support.DefaultListableBeanFactory; +import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.cloud.streams.adapter.ChannelBinding; +import org.springframework.cloud.streams.adapter.ChannelBindingAdapter; +import org.springframework.cloud.streams.adapter.OutputChannelBinding; +import org.springframework.cloud.streams.config.ChannelBindingAdapterConfigurationTests.Empty; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Import; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.annotation.DirtiesContext.ClassMode; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.xd.dirt.integration.bus.local.LocalMessageBus; + +/** + * @author Dave Syer + */ +@RunWith(SpringJUnit4ClassRunner.class) +@SpringApplicationConfiguration(classes = Empty.class) +@DirtiesContext(classMode = ClassMode.AFTER_EACH_TEST_METHOD) +public class ChannelBindingAdapterConfigurationTests { + + @Autowired + private DefaultListableBeanFactory context; + + @Autowired + private ChannelBindingAdapter adapter; + + @Autowired + private ChannelBindingAdapterConfiguration configuration; + + @Autowired + private ChannelBindingProperties module; + + @Before + public void init() { + } + + @Test + public void oneOutput() throws Exception { + this.context.registerSingleton("output", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + assertEquals(1, channels.size()); + assertEquals("group.0", channels.iterator().next().getRemoteName()); + assertEquals("tap:group.0", channels.iterator().next().getTapChannelName()); + } + + private void refresh() { + Collection channels = this.configuration.getOutputChannels(); + for (OutputChannelBinding channel : channels) { + channel.setTapped(true); + } + this.adapter.setOutputChannels(channels); + this.adapter.start(); + } + + @Test + public void oneOutputTopic() throws Exception { + this.context.registerSingleton("output.topic:", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + assertEquals(1, channels.size()); + assertEquals("topic:group.0", channels.iterator().next().getRemoteName()); + assertEquals("tap:group.0", channels.iterator().next().getTapChannelName()); + } + + @Test + public void twoOutputsWithQueue() throws Exception { + this.context.registerSingleton("output", new DirectChannel()); + this.context.registerSingleton("output.queue:foo", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + List names = getChannelNames(channels); + assertEquals(2, channels.size()); + assertTrue(names.contains("group.0")); + assertTrue(names.contains("foo.group.0")); + } + + private List getChannelNames(Collection channels) { + List list = new ArrayList(); + for (ChannelBinding binding : channels) { + list.add(binding.getRemoteName()); + } + return list; + } + + @Test + public void overrideNaturalOutputChannelName() throws Exception { + this.module.setOutputChannelName("bar"); + this.context.registerSingleton("output.queue:foo", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + assertEquals(1, channels.size()); + assertEquals("foo.bar", channels.iterator().next().getRemoteName()); + // TODO: fix this. What should it be? + assertEquals("tap:foo.bar", channels.iterator().next().getTapChannelName()); + } + + @Test + public void overrideNaturalOutputChannelNamedQueueWithTopic() throws Exception { + this.module.setOutputChannelName("queue:bar"); + this.context.registerSingleton("output.topic:foo", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + assertEquals(1, channels.size()); + assertEquals("topic:foo.bar", channels.iterator().next().getRemoteName()); + assertEquals("tap:foo.bar", channels.iterator().next().getTapChannelName()); + } + + @Configuration + @Import(ChannelBindingAdapterConfiguration.class) + protected static class Empty { + @Bean + public LocalMessageBus messageBus() { + return new LocalMessageBus(); + } + } + +} diff --git a/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java b/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java index 84531b3d0..bf94164a6 100644 --- a/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java +++ b/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleOptionsPropertySourceInitializer.java @@ -24,7 +24,6 @@ import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.streams.config.MessageBusProperties; import org.springframework.context.ApplicationContextInitializer; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -53,13 +52,13 @@ import org.springframework.xd.module.options.ModuleOptionsMetadataResolver; * */ @Configuration -@EnableConfigurationProperties(MessageBusProperties.class) +@EnableConfigurationProperties @Order(Ordered.HIGHEST_PRECEDENCE + 10) public class ModuleOptionsPropertySourceInitializer implements - ApplicationContextInitializer { +ApplicationContextInitializer { @Autowired - private MessageBusProperties module = new MessageBusProperties(); + private ModuleProperties module = new ModuleProperties(); @Autowired(required=false) private EnvironmentAwareModuleOptionsMetadataResolver wrapper; @@ -79,8 +78,8 @@ public class ModuleOptionsPropertySourceInitializer implements } private ModuleDefinition getModuleDefinition() { - return ModuleDefinitions.simple(module.getName(), - ModuleType.valueOf(module.getType()), "file:."); + return ModuleDefinitions.simple(this.module.getName(), + ModuleType.valueOf(this.module.getType()), "file:."); } private void insert(ConfigurableEnvironment environment, MapPropertySource source) { @@ -95,9 +94,9 @@ public class ModuleOptionsPropertySourceInitializer implements DelegatingModuleOptionsMetadataResolver delegatingResolver = new DelegatingModuleOptionsMetadataResolver(); delegatingResolver.setDelegates(delegates); ModuleOptionsMetadataResolver resolver = delegatingResolver; - if (wrapper!=null) { - wrapper.setDelegate(delegatingResolver); - resolver = wrapper; + if (this.wrapper!=null) { + this.wrapper.setDelegate(delegatingResolver); + resolver = this.wrapper; } return resolver; } @@ -118,4 +117,12 @@ public class ModuleOptionsPropertySourceInitializer implements } } + @Configuration + protected static class ModulePropertiesConfiguration { + @Bean(name="spring.cloud.channels.CONFIGURATION_PROPERTIES") + public ModuleProperties moduleProperties() { + return new ModuleProperties(); + } + } + } diff --git a/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleProperties.java b/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleProperties.java new file mode 100644 index 000000000..f1365b34e --- /dev/null +++ b/spring-xd-runner/src/main/java/org/springframework/cloud/streams/xd/ModuleProperties.java @@ -0,0 +1,140 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.streams.xd; + +import org.springframework.cloud.streams.config.ChannelBindingProperties; +import org.springframework.util.Assert; +import org.springframework.xd.dirt.integration.bus.BusUtils; + +/** + * @author Dave Syer + * + */ +public class ModuleProperties extends ChannelBindingProperties { + + private String group = "group"; + private String name = "module"; + private int index = 0; + private String type = "processor"; + + private Tap tap; + + public String getType() { + return this.type; + } + + public void setType(String type) { + this.type = type; + } + + public String getGroup() { + return this.group; + } + + public void setGroup(String group) { + this.group = group; + } + + public String getName() { + return this.name; + } + + public void setName(String name) { + this.name = name; + } + + public int getIndex() { + return this.index; + } + + public void setIndex(int index) { + this.index = index; + } + + @Override + public String getInputChannelName() { + if (isTap()) { + return String.format("%s.%s.%s", BusUtils.constructTapPrefix(this.tap.getGroup()), + this.tap.getName(), this.tap.getIndex()); + } + return super.getInputChannelName(); + } + + @Override + public String getTapChannelName() { + return getTapChannelName(this.group); + } + + @Override + public String getTapChannelName(String prefix) { + Assert.isTrue(!this.type .equals("job"), "Job module type not supported."); + // for Stream return channel name with indexed elements + return String + .format("%s.%s.%s", BusUtils.constructTapPrefix(prefix), this.name, this.index); + } + + public Tap getTap() { + return this.tap; + } + + public void setTap(Tap tap) { + this.tap = tap; + } + + private boolean isTap() { + if (this.tap != null) { + Assert.state(this.tap.getName() != null, "Tap name not provided"); + Assert.state(!this.tap.getGroup().equals(this.group), + "Tap group cannot be the same as module group"); + } + return this.tap != null; + } + + public static class Tap { + + private String group = "group"; + + private String name; + + public String getName() { + return this.name; + } + + public void setName(String name) { + this.name = name; + } + + private int index = 0; + + public String getGroup() { + return this.group; + } + + public void setGroup(String group) { + this.group = group; + } + + public int getIndex() { + return this.index; + } + + public void setIndex(int index) { + this.index = index; + } + + } +} diff --git a/spring-xd-runner/src/main/resources/META-INF/spring.factories b/spring-xd-runner/src/main/resources/META-INF/spring.factories index 2606d14dd..609330652 100644 --- a/spring-xd-runner/src/main/resources/META-INF/spring.factories +++ b/spring-xd-runner/src/main/resources/META-INF/spring.factories @@ -1,2 +1,2 @@ org.springframework.cloud.bootstrap.BootstrapConfiguration:\ -org.springframework.cloud.streams.xd.ModuleOptionsPropertySourceInitializer \ No newline at end of file +org.springframework.cloud.streams.xd.ModuleOptionsPropertySourceInitializer diff --git a/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/config/MessageBusAdapterConfigurationTests.java b/spring-xd-runner/src/test/java/org/springframework/cloud/streams/xd/ChannelBindingAdapterConfigurationTests.java similarity index 87% rename from spring-cloud-streams/src/test/java/org/springframework/cloud/streams/config/MessageBusAdapterConfigurationTests.java rename to spring-xd-runner/src/test/java/org/springframework/cloud/streams/xd/ChannelBindingAdapterConfigurationTests.java index ec479c136..f721af599 100644 --- a/spring-cloud-streams/src/test/java/org/springframework/cloud/streams/config/MessageBusAdapterConfigurationTests.java +++ b/spring-xd-runner/src/test/java/org/springframework/cloud/streams/xd/ChannelBindingAdapterConfigurationTests.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.streams.config; +package org.springframework.cloud.streams.xd; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; @@ -30,11 +30,11 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.support.DefaultListableBeanFactory; import org.springframework.boot.test.SpringApplicationConfiguration; import org.springframework.cloud.streams.adapter.ChannelBinding; -import org.springframework.cloud.streams.adapter.MessageBusAdapter; +import org.springframework.cloud.streams.adapter.ChannelBindingAdapter; import org.springframework.cloud.streams.adapter.OutputChannelBinding; -import org.springframework.cloud.streams.config.MessageBusAdapterConfiguration; -import org.springframework.cloud.streams.config.MessageBusProperties; -import org.springframework.cloud.streams.config.MessageBusAdapterConfigurationTests.Empty; +import org.springframework.cloud.streams.config.ChannelBindingAdapterConfiguration; +import org.springframework.cloud.streams.config.ChannelBindingProperties; +import org.springframework.cloud.streams.xd.ChannelBindingAdapterConfigurationTests.Empty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; @@ -50,19 +50,19 @@ import org.springframework.xd.dirt.integration.bus.local.LocalMessageBus; @RunWith(SpringJUnit4ClassRunner.class) @SpringApplicationConfiguration(classes = Empty.class) @DirtiesContext(classMode = ClassMode.AFTER_EACH_TEST_METHOD) -public class MessageBusAdapterConfigurationTests { +public class ChannelBindingAdapterConfigurationTests { @Autowired private DefaultListableBeanFactory context; @Autowired - private MessageBusAdapter adapter; + private ChannelBindingAdapter adapter; @Autowired - private MessageBusAdapterConfiguration configuration; + private ChannelBindingAdapterConfiguration configuration; @Autowired - private MessageBusProperties module; + private ChannelBindingProperties module; @Before public void init() { @@ -79,7 +79,8 @@ public class MessageBusAdapterConfigurationTests { } private void refresh() { - Collection channels = this.configuration.getOutputChannels(); + this.configuration.refresh(); + Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); for (OutputChannelBinding channel : channels) { channel.setTapped(true); } @@ -141,7 +142,7 @@ public class MessageBusAdapterConfigurationTests { } @Configuration - @Import(MessageBusAdapterConfiguration.class) + @Import(ChannelBindingAdapterConfiguration.class) protected static class Empty { @Bean public LocalMessageBus messageBus() {