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
This commit is contained in:
committed by
Oleg Zhurakousky
parent
c875670b73
commit
807f51f175
@@ -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);
|
||||
|
||||
@@ -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<String> consumerMultiple() {
|
||||
return value -> {
|
||||
System.out.println(value);
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
public static class Person {
|
||||
private String name;
|
||||
private int id;
|
||||
|
||||
@@ -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<Binding<?>> 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<Binding<?>> locateBinding(String name) {
|
||||
Stream<Binding<?>> 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());
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -53,7 +53,7 @@ public class BindingsEndpoint {
|
||||
}
|
||||
|
||||
@ReadOperation
|
||||
public Binding<?> queryState(@Selector String name) {
|
||||
public List<Binding<?>> queryState(@Selector String name) {
|
||||
return this.lifecycleController.queryState(name);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user