From 82929f9bcf206d8fd4eb2e6f12e57599a67423b8 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 25 Feb 2021 16:04:23 +0100 Subject: [PATCH] GH-2106 Add bridge between destination channels and binding channel in TestChannelBinder Resolves #2106 --- .../binder/test/AbstractDestination.java | 2 +- .../stream/binder/test/OutputDestination.java | 6 +- .../stream/binder/test/TestChannelBinder.java | 24 ++++++- .../test/TestChannelBinderProvisioner.java | 17 ++++- .../cloud/stream/function/ScenarioTests.java | 62 +++++++++++++++++++ 5 files changed, 104 insertions(+), 7 deletions(-) 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 12040958a..ae2cb56d0 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 5f73b083f..4c2f81166 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,6 +23,7 @@ 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; @@ -40,7 +41,6 @@ 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,7 +109,9 @@ public class OutputDestination extends AbstractDestination { if (!this.messageQueues.containsKey(bindingName)) { BlockingQueue> messageQueue = new LinkedTransferQueue<>(); this.messageQueues.put(bindingName, messageQueue); - this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message) message)); + if (((AbstractSubscribableChannel) this.getChannelByName(bindingName)).getSubscriberCount() < 1) { + 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 dc08ba4dc..7a773f4cc 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,12 +20,15 @@ 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; @@ -51,6 +54,7 @@ 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; /** @@ -169,7 +173,25 @@ public class TestChannelBinder extends adapter.setErrorChannel(errorInfrastructure.getErrorChannel()); } - siBinderInputChannel.subscribe(messageListenerContainer); + 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); + } 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 7c91a21ac..ab5ce0aa2 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,6 +19,7 @@ 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; @@ -27,6 +28,9 @@ 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; @@ -43,7 +47,7 @@ import org.springframework.messaging.SubscribableChannel; * */ public class TestChannelBinderProvisioner - implements ProvisioningProvider { + implements ProvisioningProvider, ApplicationContextAware { private final Map provisionedDestinations = new HashMap<>(); @@ -53,6 +57,14 @@ 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 @@ -81,7 +93,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) { @@ -141,5 +153,4 @@ 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 e8df7c64a..9cc5a537f 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,6 +16,7 @@ package org.springframework.cloud.stream.function; +import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; @@ -90,7 +91,68 @@ 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