diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java index 8c1dced5a..e907cbf56 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java @@ -45,8 +45,8 @@ public class ChannelBindingService { public void bindConsumer(MessageChannel inputChannel, String inputChannelName) { String channelBindingTarget = this.channelBindingProperties .getBindingPath(inputChannelName); - if (isChannelPubSub(inputChannelName)) { - this.binder.bindPubSubConsumer(channelBindingTarget, inputChannel, + if (isChannelPubSub(channelBindingTarget)) { + this.binder.bindPubSubConsumer(removePrefix(channelBindingTarget), inputChannel, this.channelBindingProperties.getConsumerProperties()); } else { @@ -58,8 +58,8 @@ public class ChannelBindingService { public void bindProducer(MessageChannel outputChannel, String outputChannelName) { String channelBindingTarget = this.channelBindingProperties .getBindingPath(outputChannelName); - if (isChannelPubSub(outputChannelName)) { - this.binder.bindPubSubProducer(channelBindingTarget, outputChannel, + if (isChannelPubSub(channelBindingTarget)) { + this.binder.bindPubSubProducer(removePrefix(channelBindingTarget), outputChannel, this.channelBindingProperties.getProducerProperties()); } else { @@ -68,10 +68,16 @@ public class ChannelBindingService { } } - private boolean isChannelPubSub(String channelName) { - Assert.isTrue(StringUtils.hasText(channelName), - "Channel name should not be empty/null."); - return channelName.startsWith("topic:"); + private boolean isChannelPubSub(String bindingTarget) { + Assert.isTrue(StringUtils.hasText(bindingTarget), + "Binding target should not be empty/null."); + return bindingTarget.startsWith("topic:"); + } + + private String removePrefix(String bindingTarget) { + Assert.isTrue(StringUtils.hasText(bindingTarget), + "Binding target should not be empty/null."); + return bindingTarget.substring(bindingTarget.indexOf(":") + 1); } public void unbindConsumers(String inputChannelName) { diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithPubSubBindingTargets.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithPubSubBindingTargets.java new file mode 100644 index 000000000..069ca979c --- /dev/null +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/ProcessorBindingTestsWithPubSubBindingTargets.java @@ -0,0 +1,69 @@ +/* + * 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.stream.binder; + +import static org.mockito.Matchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; + +import java.util.Properties; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mockito; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.SpringApplicationConfiguration; +import org.springframework.cloud.stream.annotation.Bindings; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.messaging.Processor; +import org.springframework.cloud.stream.utils.MockBinderConfiguration; +import org.springframework.context.annotation.Import; +import org.springframework.context.annotation.PropertySource; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; + +/** + * @author Marius Bogoevici + */ +@RunWith(SpringJUnit4ClassRunner.class) +@SpringApplicationConfiguration(ProcessorBindingTestsWithPubSubBindingTargets.TestProcessor.class) +public class ProcessorBindingTestsWithPubSubBindingTargets { + + @SuppressWarnings("rawtypes") + @Autowired + private Binder binder; + + @Autowired @Bindings(TestProcessor.class) + private Processor testProcessor; + + @SuppressWarnings("unchecked") + @Test + public void testSourceOutputChannelBound() { + verify(binder).bindPubSubConsumer(eq("testtock.0"), eq(testProcessor.input()), Mockito.any()); + verify(binder).bindPubSubProducer(eq("testtock.1"), eq(testProcessor.output()), Mockito.any()); + verifyNoMoreInteractions(binder); + } + + @EnableBinding(Processor.class) + @EnableAutoConfiguration + @Import(MockBinderConfiguration.class) + @PropertySource("classpath:/org/springframework/cloud/stream/binder/processor-binding-test-pubsub.properties") + public static class TestProcessor { + + } +} diff --git a/spring-cloud-stream/src/test/resources/org/springframework/cloud/stream/binder/processor-binding-test-pubsub.properties b/spring-cloud-stream/src/test/resources/org/springframework/cloud/stream/binder/processor-binding-test-pubsub.properties new file mode 100644 index 000000000..cd9d0bc54 --- /dev/null +++ b/spring-cloud-stream/src/test/resources/org/springframework/cloud/stream/binder/processor-binding-test-pubsub.properties @@ -0,0 +1,2 @@ +spring.cloud.stream.bindings.input=topic:testtock.0 +spring.cloud.stream.bindings.output=topic:testtock.1