diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java index ae2cb56d0..12040958a 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/AbstractDestination.java @@ -44,7 +44,7 @@ abstract class AbstractDestination { } SubscribableChannel getChannelByName(String name) { - //name = name.endsWith(".destination") ? name : name + ".destination"; + name = name.endsWith(".destination") ? name : name + ".destination"; for (AbstractSubscribableChannel subscribableChannel : channels) { if (subscribableChannel.getBeanName().equals(name)) { return subscribableChannel; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java index 4c2f81166..5f73b083f 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java @@ -23,7 +23,6 @@ import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedTransferQueue; import java.util.concurrent.TimeUnit; -import org.springframework.integration.channel.AbstractSubscribableChannel; import org.springframework.messaging.Message; import org.springframework.util.StringUtils; @@ -41,6 +40,7 @@ public class OutputDestination extends AbstractDestination { public Message receive(long timeout, String bindingName) { try { + bindingName = bindingName.endsWith(".destination") ? bindingName : bindingName + ".destination"; return this.messageQueues.get(bindingName).poll(timeout, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { @@ -109,9 +109,7 @@ public class OutputDestination extends AbstractDestination { if (!this.messageQueues.containsKey(bindingName)) { BlockingQueue> messageQueue = new LinkedTransferQueue<>(); this.messageQueues.put(bindingName, messageQueue); - if (((AbstractSubscribableChannel) this.getChannelByName(bindingName)).getSubscriberCount() < 1) { - this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message) message)); - } + this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message) message)); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java index 7a773f4cc..dc08ba4dc 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinder.java @@ -20,15 +20,12 @@ import java.util.function.Consumer; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binder.test.TestChannelBinderProvisioner.SpringIntegrationConsumerDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderProvisioner.SpringIntegrationProducerDestination; -import org.springframework.cloud.stream.config.BindingProperties; -import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.core.AttributeAccessor; @@ -54,7 +51,6 @@ import org.springframework.retry.RetryContext; import org.springframework.retry.RetryListener; import org.springframework.retry.support.RetryTemplate; import org.springframework.util.Assert; -import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -173,25 +169,7 @@ public class TestChannelBinder extends adapter.setErrorChannel(errorInfrastructure.getErrorChannel()); } - SubscribableChannel bindingChannel = null; - if (ObjectUtils.isEmpty(this.getApplicationContext().getBeanNamesForAnnotation(EnableBinding.class))) { - BindingServiceProperties bs = this.getApplicationContext().getBean(BindingServiceProperties.class); - for (String bindingName : bs.getBindings().keySet()) { - BindingProperties bindingProperties = bs.getBindingProperties(bindingName); - if (!bindingName.equals(destination.getName()) && destination.getName().equals(bindingProperties.getDestination())) { - BridgeHandler bridge = new BridgeHandler(); - if (this.getApplicationContext().containsBean(bindingName)) { - bindingChannel = this.getApplicationContext().getBean(bindingName, SubscribableChannel.class); - bridge.setOutputChannel(bindingChannel); - siBinderInputChannel.subscribe(bridge); - } - } - } - } - - if (bindingChannel == null) { - siBinderInputChannel.subscribe(messageListenerContainer); - } + siBinderInputChannel.subscribe(messageListenerContainer); return adapter; } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinderProvisioner.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinderProvisioner.java index ab5ce0aa2..7c91a21ac 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinderProvisioner.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/TestChannelBinderProvisioner.java @@ -19,7 +19,6 @@ package org.springframework.cloud.stream.binder.test; import java.util.HashMap; import java.util.Map; -import org.springframework.beans.BeansException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.ConsumerProperties; @@ -28,9 +27,6 @@ import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.cloud.stream.provisioning.ProvisioningException; import org.springframework.cloud.stream.provisioning.ProvisioningProvider; -import org.springframework.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; -import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.PublishSubscribeChannel; @@ -47,7 +43,7 @@ import org.springframework.messaging.SubscribableChannel; * */ public class TestChannelBinderProvisioner - implements ProvisioningProvider, ApplicationContextAware { + implements ProvisioningProvider { private final Map provisionedDestinations = new HashMap<>(); @@ -57,14 +53,6 @@ public class TestChannelBinderProvisioner @Autowired private OutputDestination target; - private GenericApplicationContext applicationContext; - - @Override - public void setApplicationContext(ApplicationContext applicationContext) - throws BeansException { - this.applicationContext = (GenericApplicationContext) applicationContext; - } - /** * Will provision producer destination as an SI {@link PublishSubscribeChannel}.
* This provides convenience of registering additional subscriber (handler in the test @@ -93,7 +81,7 @@ public class TestChannelBinderProvisioner } private SubscribableChannel provisionDestination(String name, boolean pubSub) { - String destinationName = name; // + ".destination"; + String destinationName = name + ".destination"; SubscribableChannel destination = this.provisionedDestinations .get(destinationName); if (destination == null) { @@ -153,4 +141,5 @@ public class TestChannelBinderProvisioner } } + } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java index 9cc5a537f..e8df7c64a 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java @@ -16,7 +16,6 @@ package org.springframework.cloud.stream.function; -import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; @@ -91,68 +90,7 @@ public class ScenarioTests { } } - @Test - public void test2106() { - try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration - .getCompleteConfiguration(ConsumerConfiguration.class, ConsumerConfiguration.class)) - .web(WebApplicationType.NONE).run( - "--spring.cloud.function.definition=consume;echo", - "--spring.cloud.stream.bindings.consume-in-0.destination=input", - "--spring.cloud.stream.bindings.echo-in-0.destination=echoin", - "--spring.cloud.stream.bindings.echo-out-0.destination=echoout", - "--spring.jmx.enabled=false")) { - ConsumerConfiguration configuration = context.getBean(ConsumerConfiguration.class); - - OutputDestination output = context.getBean(OutputDestination.class); - - StreamBridge bridge = context.getBean(StreamBridge.class); - bridge.send("input", "destination"); - bridge.send("input", "destination"); - bridge.send("input", "destination"); - - bridge.send("consume-in-0", "hello"); - bridge.send("consume-in-0", "hello"); - bridge.send("consume-in-0", "hello"); - - bridge.send("echoin", "hello"); - bridge.send("echoin", "hello"); - bridge.send("echoin", "hello"); - - assertThat(configuration.destinationCounter).isEqualTo(3); - assertThat(configuration.bindingCounter).isEqualTo(3); - - assertThat(output.receive(1000, "echoout")).isNotNull(); - assertThat(output.receive(1000, "echoout")).isNotNull(); - assertThat(output.receive(1000, "echoout")).isNotNull(); - assertThat(output.receive(1000, "echoout")).isNull(); - } - } - - @EnableAutoConfiguration - public static class ConsumerConfiguration { - - private int destinationCounter; - - private int bindingCounter; - - @Bean - public Consumer consume() { - return v -> { - if (v.equals("destination")) { - destinationCounter++; - } - else { - bindingCounter++; - } - }; - } - - @Bean - public Function echo() { - return v -> v; - } - } @EnableAutoConfiguration @Configuration