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
Soby Chacko
parent
8059cdabcf
commit
9c9da804e2
@@ -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);
|
||||
|
||||
@@ -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<InputBindingLifecycle> inputBindingLifecycles,
|
||||
List<OutputBindingLifecycle> outputBindingsLifecycles) {
|
||||
List<OutputBindingLifecycle> 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<? extends Module> javaTimeModuleClass = (Class<? extends Module>)
|
||||
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<Binding<?>> 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<Binding<?>> 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<Binding<?>> inputBindings = new ArrayList<>();
|
||||
for (InputBindingLifecycle inputBindingLifecycle : this.inputBindingLifecycles) {
|
||||
Collection<Binding<?>> lifecycleInputBindings = (Collection<Binding<?>>) new DirectFieldAccessor(
|
||||
inputBindingLifecycle).getPropertyValue("inputBindings");
|
||||
inputBindingLifecycle).getPropertyValue("inputBindings");
|
||||
inputBindings.addAll(lifecycleInputBindings);
|
||||
}
|
||||
return inputBindings;
|
||||
@@ -177,17 +176,16 @@ public class BindingsLifecycleController {
|
||||
List<Binding<?>> outputBindings = new ArrayList<>();
|
||||
for (OutputBindingLifecycle inputBindingLifecycle : this.outputBindingsLifecycles) {
|
||||
Collection<Binding<?>> lifecycleInputBindings = (Collection<Binding<?>>) new DirectFieldAccessor(
|
||||
inputBindingLifecycle).getPropertyValue("outputBindings");
|
||||
inputBindingLifecycle).getPropertyValue("outputBindings");
|
||||
outputBindings.addAll(lifecycleInputBindings);
|
||||
}
|
||||
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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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<Binding<?>> input = ctrl.queryState("echo-in-0");
|
||||
List<Binding<?>> 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<Binding<?>> 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<Binding<?>> 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<String> consumerMultiple() {
|
||||
return value -> {
|
||||
System.out.println(value);
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
public static class Person {
|
||||
private String name;
|
||||
private int id;
|
||||
|
||||
Reference in New Issue
Block a user