From 807f51f175ac2d943b4f4252fde5c318d51b920e Mon Sep 17 00:00:00 2001 From: Fernando Blanch Date: Tue, 28 Feb 2023 22:04:11 +0100 Subject: [PATCH] Support pause/resume for consumer bindings with multiple destinations from BindingsLifecycleController add test queryng a binding that not exists return empty list remove unnecessary formatting changes Resolves #2660 Resolves #2658 --- .../KafkaRetryDlqBinderOrContainerTests.java | 4 +- .../ImplicitFunctionBindingTests.java | 68 +++++++++++++++++-- .../binding/BindingsLifecycleController.java | 24 +++---- .../stream/endpoint/BindingsEndpoint.java | 2 +- 4 files changed, 76 insertions(+), 22 deletions(-) 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 f5aec50f0..5bd83f38a 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-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 5118decf4..3a618e32e 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -32,6 +32,8 @@ import java.util.function.Supplier; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; @@ -102,7 +104,7 @@ public class ImplicitFunctionBindingTests { } } - @SuppressWarnings({"rawtypes" }) + @SuppressWarnings({"rawtypes"}) @Test void testDisableAutodetect() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( @@ -111,10 +113,10 @@ public class ImplicitFunctionBindingTests { .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(); + var input = ctrl.queryState("echo-in-0"); + var output = ctrl.queryState("echo-out-0"); + assertThat(input).isEmpty(); + assertThat(output).isEmpty(); } } @@ -128,14 +130,55 @@ public class ImplicitFunctionBindingTests { .run("--spring.jmx.enabled=false")) { BindingsLifecycleController ctrl = context.getBean(BindingsLifecycleController.class); - Binding input = ctrl.queryState("echo-in-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); + var 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); + var inputBindingList = ctrl.queryState("bindingNotExist-in-0"); + assertThat(inputBindingList).isEmpty(); + } + } + + @SuppressWarnings({"unchecked", "rawtypes"}) @Test void dynamicBindingTestWithFunctionRegistrationAndExplicitDestination() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( @@ -1563,6 +1606,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; 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 e1970037a..cc85c92f2 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 @@ -20,6 +20,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.List; import java.util.Map; +import java.util.stream.Collectors; import java.util.stream.Stream; import com.fasterxml.jackson.databind.Module; @@ -108,13 +109,13 @@ 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) { + var bindingList = BindingsLifecycleController.this.locateBinding(bindingName); + if (!bindingList.isEmpty()) { switch (state) { - case STARTED -> binding.start(); - case STOPPED -> binding.stop(); - case PAUSED -> binding.pause(); - case RESUMED -> binding.resume(); + case STARTED -> bindingList.stream().forEach(Binding::start); + case STOPPED -> bindingList.stream().forEach(Binding::stop); + case PAUSED -> bindingList.stream().forEach(Binding::pause); + case RESUMED -> bindingList.stream().forEach(Binding::resume); default -> { } } @@ -138,9 +139,9 @@ 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); } @@ -175,11 +176,10 @@ public class BindingsLifecycleController { 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); }