From fc2397b11d8e78f9187a52fc21d58e1299a45a37 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Thu, 9 Jul 2015 10:01:07 +0100 Subject: [PATCH] Check simple XD samples are working --- .../config/ChannelBindingProperties.java | 6 +- .../cloud/streams/xd/ModuleProperties.java | 27 +++++-- ...annelBindingAdapterConfigurationTests.java | 81 ++++++++++++++++--- .../LogSinkOptionsMetadata.java | 2 +- .../main/resources/config/logger.properties | 1 + .../TimeSourceOptionsMetadata.java | 2 +- .../main/resources/config/ticker.properties | 2 +- .../TimeSourceOptionsMetadata.java | 2 +- .../main/resources/config/ticker.properties | 2 +- 9 files changed, 100 insertions(+), 25 deletions(-) rename spring-xd-samples/sink/src/main/java/{org/springframework/xd/dirt/modules/metadata => config}/LogSinkOptionsMetadata.java (96%) rename spring-xd-samples/source-xml/src/main/java/{org/springframework/xd/dirt/modules/metadata => demo}/TimeSourceOptionsMetadata.java (96%) rename spring-xd-samples/source/src/main/java/{org/springframework/xd/dirt/modules/metadata => config}/TimeSourceOptionsMetadata.java (96%) 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 index c54d0b3ae..a2508ac2f 100644 --- 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 @@ -31,9 +31,11 @@ import com.fasterxml.jackson.annotation.JsonInclude.Include; @JsonInclude(Include.NON_DEFAULT) public class ChannelBindingProperties { - private String outputChannelName = "group.0"; + public static final String DEFAULT_CHANNEL_NAME = "group.0"; - private String inputChannelName = "group.0"; + private String outputChannelName = DEFAULT_CHANNEL_NAME; + + private String inputChannelName = DEFAULT_CHANNEL_NAME; private Properties consumerProperties = new Properties(); 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 index f1365b34e..f847c1384 100644 --- 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 @@ -65,13 +65,28 @@ public class ModuleProperties extends ChannelBindingProperties { this.index = index; } + @Override + public String getOutputChannelName() { + String name = super.getOutputChannelName(); + if (ChannelBindingProperties.DEFAULT_CHANNEL_NAME.equals(name)) { + return BusUtils.constructPipeName(this.group, this.index); + } + return name; + } + @Override public String getInputChannelName() { if (isTap()) { - return String.format("%s.%s.%s", BusUtils.constructTapPrefix(this.tap.getGroup()), - this.tap.getName(), this.tap.getIndex()); + return String.format("%s.%s.%s", + BusUtils.constructTapPrefix(this.tap.getGroup()), this.tap.getName(), + this.tap.getIndex()); } - return super.getInputChannelName(); + String name = super.getInputChannelName(); + if (ChannelBindingProperties.DEFAULT_CHANNEL_NAME.equals(name)) { + return BusUtils.constructPipeName(this.group, this.index > 0 ? this.index - 1 + : this.index); + } + return name; } @Override @@ -81,10 +96,10 @@ public class ModuleProperties extends ChannelBindingProperties { @Override public String getTapChannelName(String prefix) { - Assert.isTrue(!this.type .equals("job"), "Job module type not supported."); + 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); + return String.format("%s.%s.%s", BusUtils.constructTapPrefix(prefix), this.name, + this.index); } public Tap getTap() { diff --git a/spring-xd-runner/src/test/java/org/springframework/cloud/streams/xd/ChannelBindingAdapterConfigurationTests.java b/spring-xd-runner/src/test/java/org/springframework/cloud/streams/xd/ChannelBindingAdapterConfigurationTests.java index f721af599..ad3b64237 100644 --- a/spring-xd-runner/src/test/java/org/springframework/cloud/streams/xd/ChannelBindingAdapterConfigurationTests.java +++ b/spring-xd-runner/src/test/java/org/springframework/cloud/streams/xd/ChannelBindingAdapterConfigurationTests.java @@ -31,9 +31,9 @@ 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.InputChannelBinding; import org.springframework.cloud.streams.adapter.OutputChannelBinding; 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; @@ -62,7 +62,7 @@ public class ChannelBindingAdapterConfigurationTests { private ChannelBindingAdapterConfiguration configuration; @Autowired - private ChannelBindingProperties module; + private ModuleProperties module; @Before public void init() { @@ -72,15 +72,28 @@ public class ChannelBindingAdapterConfigurationTests { public void oneOutput() throws Exception { this.context.registerSingleton("output", new DirectChannel()); refresh(); - Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + Collection channels = this.adapter.getChannelsMetadata() + .getOutputChannels(); + assertEquals(1, channels.size()); + assertEquals("group.0", channels.iterator().next().getRemoteName()); + assertEquals("tap:stream:group.module.0", channels.iterator().next() + .getTapChannelName()); + } + + @Test + public void oneInput() throws Exception { + this.context.registerSingleton("input", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata() + .getInputChannels(); assertEquals(1, channels.size()); assertEquals("group.0", channels.iterator().next().getRemoteName()); - assertEquals("tap:stream:group.module.0", channels.iterator().next().getTapChannelName()); } private void refresh() { this.configuration.refresh(); - Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + Collection channels = this.adapter.getChannelsMetadata() + .getOutputChannels(); for (OutputChannelBinding channel : channels) { channel.setTapped(true); } @@ -92,10 +105,49 @@ public class ChannelBindingAdapterConfigurationTests { public void oneOutputTopic() throws Exception { this.context.registerSingleton("output.topic:", new DirectChannel()); refresh(); - Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + Collection channels = this.adapter.getChannelsMetadata() + .getOutputChannels(); assertEquals(1, channels.size()); assertEquals("topic:group.0", channels.iterator().next().getRemoteName()); - assertEquals("tap:stream:group.module.0", channels.iterator().next().getTapChannelName()); + assertEquals("tap:stream:group.module.0", channels.iterator().next() + .getTapChannelName()); + } + + @Test + public void oneOutputOverrideName() throws Exception { + this.module.setGroup("mine"); + this.module.setName("foo"); + this.module.setIndex(2); + this.context.registerSingleton("output.topic:", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata() + .getOutputChannels(); + assertEquals(1, channels.size()); + assertEquals("topic:mine.2", channels.iterator().next().getRemoteName()); + assertEquals("tap:stream:mine.foo.2", channels.iterator().next() + .getTapChannelName()); + } + + @Test + public void oneInputTopic() throws Exception { + this.context.registerSingleton("input.topic:", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata() + .getInputChannels(); + assertEquals(1, channels.size()); + assertEquals("topic:group.0", channels.iterator().next().getRemoteName()); + } + + @Test + public void oneInputOverrideName() throws Exception { + this.module.setGroup("mine"); + this.module.setIndex(2); + this.context.registerSingleton("input", new DirectChannel()); + refresh(); + Collection channels = this.adapter.getChannelsMetadata() + .getInputChannels(); + assertEquals(1, channels.size()); + assertEquals("mine.1", channels.iterator().next().getRemoteName()); } @Test @@ -103,7 +155,8 @@ public class ChannelBindingAdapterConfigurationTests { this.context.registerSingleton("output", new DirectChannel()); this.context.registerSingleton("output.queue:foo", new DirectChannel()); refresh(); - Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + Collection channels = this.adapter.getChannelsMetadata() + .getOutputChannels(); List names = getChannelNames(channels); assertEquals(2, channels.size()); assertTrue(names.contains("group.0")); @@ -123,11 +176,13 @@ public class ChannelBindingAdapterConfigurationTests { this.module.setOutputChannelName("bar"); this.context.registerSingleton("output.queue:foo", new DirectChannel()); refresh(); - Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + 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:stream:foo.bar.module.0", channels.iterator().next().getTapChannelName()); + assertEquals("tap:stream:foo.bar.module.0", channels.iterator().next() + .getTapChannelName()); } @Test @@ -135,10 +190,12 @@ public class ChannelBindingAdapterConfigurationTests { this.module.setOutputChannelName("queue:bar"); this.context.registerSingleton("output.topic:foo", new DirectChannel()); refresh(); - Collection channels = this.adapter.getChannelsMetadata().getOutputChannels(); + Collection channels = this.adapter.getChannelsMetadata() + .getOutputChannels(); assertEquals(1, channels.size()); assertEquals("topic:foo.bar", channels.iterator().next().getRemoteName()); - assertEquals("tap:stream:foo.bar.module.0", channels.iterator().next().getTapChannelName()); + assertEquals("tap:stream:foo.bar.module.0", channels.iterator().next() + .getTapChannelName()); } @Configuration diff --git a/spring-xd-samples/sink/src/main/java/org/springframework/xd/dirt/modules/metadata/LogSinkOptionsMetadata.java b/spring-xd-samples/sink/src/main/java/config/LogSinkOptionsMetadata.java similarity index 96% rename from spring-xd-samples/sink/src/main/java/org/springframework/xd/dirt/modules/metadata/LogSinkOptionsMetadata.java rename to spring-xd-samples/sink/src/main/java/config/LogSinkOptionsMetadata.java index 13b0e2422..222649e0c 100644 --- a/spring-xd-samples/sink/src/main/java/org/springframework/xd/dirt/modules/metadata/LogSinkOptionsMetadata.java +++ b/spring-xd-samples/sink/src/main/java/config/LogSinkOptionsMetadata.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.xd.dirt.modules.metadata; +package config; import org.hibernate.validator.constraints.NotBlank; diff --git a/spring-xd-samples/sink/src/main/resources/config/logger.properties b/spring-xd-samples/sink/src/main/resources/config/logger.properties index b13888b25..cfa03fa2f 100644 --- a/spring-xd-samples/sink/src/main/resources/config/logger.properties +++ b/spring-xd-samples/sink/src/main/resources/config/logger.properties @@ -1 +1,2 @@ +options_class = config.LogSinkOptionsMetadata base_packages = config \ No newline at end of file diff --git a/spring-xd-samples/source-xml/src/main/java/org/springframework/xd/dirt/modules/metadata/TimeSourceOptionsMetadata.java b/spring-xd-samples/source-xml/src/main/java/demo/TimeSourceOptionsMetadata.java similarity index 96% rename from spring-xd-samples/source-xml/src/main/java/org/springframework/xd/dirt/modules/metadata/TimeSourceOptionsMetadata.java rename to spring-xd-samples/source-xml/src/main/java/demo/TimeSourceOptionsMetadata.java index 6dc884938..93426070e 100644 --- a/spring-xd-samples/source-xml/src/main/java/org/springframework/xd/dirt/modules/metadata/TimeSourceOptionsMetadata.java +++ b/spring-xd-samples/source-xml/src/main/java/demo/TimeSourceOptionsMetadata.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.xd.dirt.modules.metadata; +package demo; import org.springframework.xd.module.options.mixins.MaxMessagesDefaultOneMixin; import org.springframework.xd.module.options.mixins.PeriodicTriggerMixin; diff --git a/spring-xd-samples/source-xml/src/main/resources/config/ticker.properties b/spring-xd-samples/source-xml/src/main/resources/config/ticker.properties index 96c0017ad..f648a982d 100644 --- a/spring-xd-samples/source-xml/src/main/resources/config/ticker.properties +++ b/spring-xd-samples/source-xml/src/main/resources/config/ticker.properties @@ -1 +1 @@ -options_class = org.springframework.xd.dirt.modules.metadata.TimeSourceOptionsMetadata \ No newline at end of file +options_class = demo.TimeSourceOptionsMetadata \ No newline at end of file diff --git a/spring-xd-samples/source/src/main/java/org/springframework/xd/dirt/modules/metadata/TimeSourceOptionsMetadata.java b/spring-xd-samples/source/src/main/java/config/TimeSourceOptionsMetadata.java similarity index 96% rename from spring-xd-samples/source/src/main/java/org/springframework/xd/dirt/modules/metadata/TimeSourceOptionsMetadata.java rename to spring-xd-samples/source/src/main/java/config/TimeSourceOptionsMetadata.java index 6dc884938..fd645b5b4 100644 --- a/spring-xd-samples/source/src/main/java/org/springframework/xd/dirt/modules/metadata/TimeSourceOptionsMetadata.java +++ b/spring-xd-samples/source/src/main/java/config/TimeSourceOptionsMetadata.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.xd.dirt.modules.metadata; +package config; import org.springframework.xd.module.options.mixins.MaxMessagesDefaultOneMixin; import org.springframework.xd.module.options.mixins.PeriodicTriggerMixin; diff --git a/spring-xd-samples/source/src/main/resources/config/ticker.properties b/spring-xd-samples/source/src/main/resources/config/ticker.properties index 464d6dcc4..9fb4a6b77 100644 --- a/spring-xd-samples/source/src/main/resources/config/ticker.properties +++ b/spring-xd-samples/source/src/main/resources/config/ticker.properties @@ -1,2 +1,2 @@ -options_class = org.springframework.xd.dirt.modules.metadata.TimeSourceOptionsMetadata +options_class = config.TimeSourceOptionsMetadata base_packages = config \ No newline at end of file