diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaRetryDlqBinderOrContainerTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaRetryDlqBinderOrContainerTests.java index 5dc0dbce0..399ea35fe 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaRetryDlqBinderOrContainerTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaRetryDlqBinderOrContainerTests.java @@ -67,11 +67,11 @@ public class KafkaRetryDlqBinderOrContainerTests { @Test void retryAndDlqInRightPlace(@Autowired BindingsLifecycleController controller) throws Exception { - Binding retryInBinder = controller.queryState("retryInBinder-in-0"); + Binding retryInBinder = controller.queryState("retryInBinder-in-0").get(0); assertThat(KafkaTestUtils.getPropertyValue(retryInBinder, "lifecycle.retryTemplate")).isNotNull(); assertThat(KafkaTestUtils.getPropertyValue(retryInBinder, "lifecycle.messageListenerContainer.commonErrorHandler")).isNull(); - Binding retryInContainer = controller.queryState("retryInContainer-in-0"); + Binding retryInContainer = controller.queryState("retryInContainer-in-0").get(0); assertThat(KafkaTestUtils.getPropertyValue(retryInContainer, "lifecycle.retryTemplate")).isNull(); assertThat(KafkaTestUtils.getPropertyValue(retryInContainer, "lifecycle.messageListenerContainer.commonErrorHandler")).isInstanceOf(CommonErrorHandler.class); diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingsLifecycleController.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingsLifecycleController.java index 107fa7c66..6d21eb7e8 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingsLifecycleController.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingsLifecycleController.java @@ -21,6 +21,7 @@ import java.util.Collection; import java.util.List; import java.util.Map; import java.util.stream.Stream; +import java.util.stream.Collectors; import com.fasterxml.jackson.databind.Module; import com.fasterxml.jackson.databind.ObjectMapper; @@ -49,9 +50,9 @@ public class BindingsLifecycleController { @SuppressWarnings("unchecked") public BindingsLifecycleController(List inputBindingLifecycles, - List outputBindingsLifecycles) { + List outputBindingsLifecycles) { Assert.notEmpty(inputBindingLifecycles, - "'inputBindingLifecycles' must not be null or empty"); + "'inputBindingLifecycles' must not be null or empty"); this.inputBindingLifecycles = inputBindingLifecycles; this.outputBindingsLifecycles = outputBindingsLifecycles; @@ -61,7 +62,7 @@ public class BindingsLifecycleController { try { Class javaTimeModuleClass = (Class) - ClassUtils.forName("com.fasterxml.jackson.datatype.jsr310.JavaTimeModule", ClassUtils.getDefaultClassLoader()); + ClassUtils.forName("com.fasterxml.jackson.datatype.jsr310.JavaTimeModule", ClassUtils.getDefaultClassLoader()); Module javaTimeModule = BeanUtils.instantiateClass(javaTimeModuleClass); this.objectMapper.registerModule(javaTimeModule); } @@ -108,23 +109,22 @@ public class BindingsLifecycleController { * @param state the {@link State} you wish to set this binding to */ public void changeState(String bindingName, State state) { - Binding binding = BindingsLifecycleController.this.locateBinding(bindingName); - if (binding != null) { + List> bindingList = BindingsLifecycleController.this.locateBinding(bindingName); + if (!bindingList.isEmpty()) { switch (state) { - case STARTED: - binding.start(); - break; - case STOPPED: - binding.stop(); - break; - case PAUSED: - binding.pause(); - break; - case RESUMED: - binding.resume(); - break; - default: - break; + case STARTED : + bindingList.stream().forEach(Binding::start); + break; + case STOPPED : + bindingList.stream().forEach(Binding::stop); + break; + case PAUSED : + bindingList.stream().forEach(Binding::pause); + break; + case RESUMED : + bindingList.stream().forEach(Binding::resume); + break; + default : break; } } } @@ -146,13 +146,12 @@ public class BindingsLifecycleController { * Queries the individual state of a binding. The returned list * {@link Binding} object could be further interrogated * using {@link Binding#isPaused()} and {@link Binding#isRunning()}. - * @return instance of {@link Binding} object. + * @return collection of {@link Binding} objects. */ - public Binding queryState(String name) { + public List> queryState(String name) { Assert.notNull(name, "'name' must not be null"); return this.locateBinding(name); } - /** * Queries for all input {@link Binding}s. * @return the list of input {@link Binding}s @@ -162,7 +161,7 @@ public class BindingsLifecycleController { List> inputBindings = new ArrayList<>(); for (InputBindingLifecycle inputBindingLifecycle : this.inputBindingLifecycles) { Collection> lifecycleInputBindings = (Collection>) new DirectFieldAccessor( - inputBindingLifecycle).getPropertyValue("inputBindings"); + inputBindingLifecycle).getPropertyValue("inputBindings"); inputBindings.addAll(lifecycleInputBindings); } return inputBindings; @@ -177,17 +176,16 @@ public class BindingsLifecycleController { List> outputBindings = new ArrayList<>(); for (OutputBindingLifecycle inputBindingLifecycle : this.outputBindingsLifecycles) { Collection> lifecycleInputBindings = (Collection>) new DirectFieldAccessor( - inputBindingLifecycle).getPropertyValue("outputBindings"); + inputBindingLifecycle).getPropertyValue("outputBindings"); outputBindings.addAll(lifecycleInputBindings); } return outputBindings; } - private Binding locateBinding(String name) { + private List> locateBinding(String name) { Stream> bindings = Stream.concat(this.gatherInputBindings().stream(), - this.gatherOutputBindings().stream()); - return bindings.filter(binding -> name.equals(binding.getBindingName())).findFirst() - .orElse(null); + this.gatherOutputBindings().stream()); + return bindings.filter(binding -> name.equals(binding.getBindingName())).collect(Collectors.toList()); } /** diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java index 6a31f9dba..5f4238d00 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java @@ -53,7 +53,7 @@ public class BindingsEndpoint { } @ReadOperation - public Binding queryState(@Selector String name) { + public List> queryState(@Selector String name) { return this.lifecycleController.queryState(name); } diff --git a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 34ed2de73..00e38c8d5 100644 --- a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -33,6 +33,8 @@ import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import reactor.core.publisher.EmitterProcessor; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -107,41 +109,80 @@ public class ImplicitFunctionBindingTests { } } - @SuppressWarnings({"rawtypes" }) + @SuppressWarnings({"rawtypes"}) @Test - public void testDisableAutodetect() { + void testDisableAutodetect() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( - TestChannelBinderConfiguration.getCompleteConfiguration(SendToDestinationConfiguration.class)) - .web(WebApplicationType.NONE) - .run("--spring.jmx.enabled=false", "--spring.cloud.stream.function.autodetect=false")) { + TestChannelBinderConfiguration.getCompleteConfiguration(SendToDestinationConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", "--spring.cloud.stream.function.autodetect=false")) { BindingsLifecycleController ctrl = context.getBean(BindingsLifecycleController.class); - Binding input = ctrl.queryState("echo-in-0"); - Binding output = ctrl.queryState("echo-out-0"); - assertThat(input).isNull(); - assertThat(output).isNull(); + List> input = ctrl.queryState("echo-in-0"); + List> output = ctrl.queryState("echo-out-0"); + assertThat(input).isEmpty(); + assertThat(output).isEmpty(); } } - @SuppressWarnings({"rawtypes" }) @Test - public void testBindingControl() { + void testBindingControl() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( - TestChannelBinderConfiguration.getCompleteConfiguration(SendToDestinationConfiguration.class)) - .web(WebApplicationType.NONE) - .run("--spring.jmx.enabled=false")) { + TestChannelBinderConfiguration.getCompleteConfiguration(SendToDestinationConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false")) { BindingsLifecycleController ctrl = context.getBean(BindingsLifecycleController.class); - Binding input = ctrl.queryState("echo-in-0"); - Binding output = ctrl.queryState("echo-out-0"); + Binding input = ctrl.queryState("echo-in-0").get(0); assertThat(input.isRunning()).isTrue(); ctrl.changeState("echo-in-0", State.STOPPED); assertThat(input.isRunning()).isFalse(); } } - @SuppressWarnings({ "unchecked", "rawtypes" }) + @ParameterizedTest + @ValueSource(strings = {"single-destination", "destination1,destination2,destination3"}) + void testGh2658_WithMultipleDestinations(String destination) { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SingleConsumerWithMultipleDestinationConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=consumerMultiple", + "--spring.cloud.stream.bindings.consumerMultiple-in-0.group=group", + "--spring.cloud.stream.bindings.consumerMultiple-in-0.destination=" + destination)) { + + BindingsLifecycleController ctrl = context.getBean(BindingsLifecycleController.class); + List> multipleInput = ctrl.queryState("consumerMultiple-in-0"); + + assertThat(multipleInput).hasSize(destination.split(",").length); + multipleInput.stream().forEach(binding -> assertThat(binding.isRunning()).isTrue()); + + ctrl.changeState("consumerMultiple-in-0", State.STOPPED); + multipleInput.stream().forEach(binding -> assertThat(binding.isRunning()).isFalse()); + + ctrl.changeState("consumerMultiple-in-0", State.STARTED); + multipleInput.stream().forEach(binding -> assertThat(binding.isRunning()).isTrue()); + } + } + + @Test + void testGh2658_queryBindingThatNotExists() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SingleConsumerWithMultipleDestinationConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=consumerMultiple", + "--spring.cloud.stream.bindings.consumerMultiple-in-0.group=group", + "--spring.cloud.stream.bindings.consumerMultiple-in-0.destination=destination")) { + + BindingsLifecycleController ctrl = context.getBean(BindingsLifecycleController.class); + List> inputBindingList = ctrl.queryState("bindingNotExist-in-0"); + assertThat(inputBindingList).isEmpty(); + } + } + + @SuppressWarnings({"unchecked", "rawtypes"}) @Test public void dynamicBindingTestWithFunctionRegistrationAndExplicitDestination() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( @@ -1583,6 +1624,17 @@ public class ImplicitFunctionBindingTests { } } + @EnableAutoConfiguration + public static class SingleConsumerWithMultipleDestinationConfiguration { + + @Bean + public Consumer consumerMultiple() { + return value -> { + System.out.println(value); + }; + } + } + public static class Person { private String name; private int id;