From 3ebea011bfbaca8e0fab10106ff4bf16e5bd26f4 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 13 Jun 2018 19:59:54 -0400 Subject: [PATCH] GH-1389 Added support for output binding control Added missing support for visualization and control of output binding Resolves #1389 --- .../binder/AbstractMessageChannelBinder.java | 2 +- .../cloud/stream/binder/DefaultBinding.java | 10 +++++++- .../cloud/stream/binding/Bindable.java | 13 ++++++++++ .../stream/binding/BindableProxyFactory.java | 13 +++++++++- .../binding/OutputBindingLifecycle.java | 9 ++++++- .../binding/SingleBindingTargetBindable.java | 10 +++++++- .../BindingsEndpointAutoConfiguration.java | 5 ++-- .../stream/endpoint/BindingsEndpoint.java | 25 ++++++++++++++++--- 8 files changed, 77 insertions(+), 10 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index 930a3da67..db1c4eebd 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -176,7 +176,7 @@ public abstract class AbstractMessageChannelBinder binding = new DefaultBinding(destination, null, outputChannel, + Binding binding = new DefaultBinding(destination, outputChannel, producerMessageHandler instanceof Lifecycle ? (Lifecycle) producerMessageHandler : null) { @Override diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java index a3bcecc7b..7d71f4238 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java @@ -55,6 +55,8 @@ public class DefaultBinding implements Binding { private boolean paused; + private boolean restartable; + /** * Creates an instance that associates a given name, group and binding target with an * optional {@link Lifecycle} component, which will be stopped during unbinding. @@ -70,6 +72,12 @@ public class DefaultBinding implements Binding { this.group = group; this.target = target; this.lifecycle = lifecycle; + this.restartable = StringUtils.hasText(group); + } + + public DefaultBinding(String name, T target, Lifecycle lifecycle) { + this(name, null, target, lifecycle); + this.restartable = true; } public String getName() { @@ -104,7 +112,7 @@ public class DefaultBinding implements Binding { @Override public final synchronized void start() { if (!this.isRunning()) { - if (this.lifecycle != null && StringUtils.hasText(this.group)) { + if (this.lifecycle != null && this.restartable) { this.lifecycle.start(); } else { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/Bindable.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/Bindable.java index 8f4e7dc3c..5dd7a931a 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/Bindable.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/Bindable.java @@ -54,9 +54,22 @@ public interface Bindable { /** * Binds all the outputs associated with this instance. + * @deprecated as of 2.0 in favor of {@link #createAndBindOutputs(BindingService)} */ + @Deprecated default void bindOutputs(BindingService adapter) {} + /** + * Binds all the outputs associated with this instance. + * @param adapter instance of {@link BindingService} + * @return collection of {@link Binding}s + * + * @since 2.0 + */ + default Collection> createAndBindOutputs(BindingService adapter) { + return Collections.>emptyList(); + } + /** * Unbinds all the inputs associated with this instance. */ diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindableProxyFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindableProxyFactory.java index 6c3198c76..ef1dc244f 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindableProxyFactory.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindableProxyFactory.java @@ -237,8 +237,18 @@ public class BindableProxyFactory implements MethodInterceptor, FactoryBean> createAndBindOutputs(BindingService bindingService) { + List> bindings = new ArrayList<>(); if (log.isDebugEnabled()) { log.debug(String.format("Binding outputs for %s:%s", this.namespace, this.type)); } @@ -249,9 +259,10 @@ public class BindableProxyFactory implements MethodInterceptor, FactoryBean> outputBindings; + public OutputBindingLifecycle(BindingService bindingService, Map bindables) { super(bindingService, bindables); } @@ -43,7 +50,7 @@ public class OutputBindingLifecycle extends AbstractBindingLifecycle { @Override void doStartWithBindable(Bindable bindable) { - bindable.bindOutputs(bindingService); + this.outputBindings = bindable.createAndBindOutputs(bindingService); } @Override diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/SingleBindingTargetBindable.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/SingleBindingTargetBindable.java index a57ea0397..e4c7a7736 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/SingleBindingTargetBindable.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/SingleBindingTargetBindable.java @@ -17,10 +17,13 @@ package org.springframework.cloud.stream.binding; import java.util.Arrays; +import java.util.Collection; import java.util.Collections; import java.util.HashSet; import java.util.Set; +import org.springframework.cloud.stream.binder.Binding; + /** * A {@link Bindable} component that wraps a generic output binding target. Useful for * binding targets outside the {@link org.springframework.cloud.stream.annotation.Input} @@ -42,7 +45,12 @@ public class SingleBindingTargetBindable implements Bindable { @Override public void bindOutputs(BindingService bindingService) { - bindingService.bindProducer(bindingTarget, name); + this.createAndBindOutputs(bindingService); + } + + @Override + public Collection> createAndBindOutputs(BindingService bindingService) { + return Collections.singletonList(bindingService.bindProducer(bindingTarget, name)); } @Override diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingsEndpointAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingsEndpointAutoConfiguration.java index 025c7e395..ccb5f1b2b 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingsEndpointAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingsEndpointAutoConfiguration.java @@ -24,6 +24,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.cloud.stream.binding.BindingService; import org.springframework.cloud.stream.binding.InputBindingLifecycle; +import org.springframework.cloud.stream.binding.OutputBindingLifecycle; import org.springframework.cloud.stream.endpoint.BindingsEndpoint; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -40,7 +41,7 @@ import org.springframework.context.annotation.Configuration; public class BindingsEndpointAutoConfiguration { @Bean - public BindingsEndpoint bindingsEndpoint(List inputBindings) { - return new BindingsEndpoint(inputBindings); + public BindingsEndpoint bindingsEndpoint(List inputBindings, List outputBindings) { + return new BindingsEndpoint(inputBindings, outputBindings); } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java index ba4b4ba39..032d6b50c 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/endpoint/BindingsEndpoint.java @@ -19,6 +19,7 @@ package org.springframework.cloud.stream.endpoint; import java.util.ArrayList; import java.util.Collection; import java.util.List; +import java.util.stream.Stream; import com.fasterxml.jackson.databind.ObjectMapper; @@ -29,6 +30,7 @@ import org.springframework.boot.actuate.endpoint.annotation.Selector; import org.springframework.boot.actuate.endpoint.annotation.WriteOperation; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binding.InputBindingLifecycle; +import org.springframework.cloud.stream.binding.OutputBindingLifecycle; import org.springframework.util.Assert; /** @@ -45,11 +47,14 @@ public class BindingsEndpoint { private final List inputBindingLifecycles; + private final List outputBindingsLifecycles; + private final ObjectMapper objectMapper; - public BindingsEndpoint(List inputBindingLifecycles) { + public BindingsEndpoint(List inputBindingLifecycles, List outputBindingsLifecycles) { Assert.notEmpty(inputBindingLifecycles, "'inputBindingLifecycles' must not be null or empty"); this.inputBindingLifecycles = inputBindingLifecycles; + this.outputBindingsLifecycles = outputBindingsLifecycles; this.objectMapper = new ObjectMapper(); } @@ -78,7 +83,9 @@ public class BindingsEndpoint { @ReadOperation public List queryStates() { - return objectMapper.convertValue(gatherInputBindings(), List.class); + List> bindings = new ArrayList<>(gatherInputBindings()); + bindings.addAll(gatherOutputBindings()); + return objectMapper.convertValue(bindings, List.class); } @ReadOperation @@ -98,8 +105,20 @@ public class BindingsEndpoint { return inputBindings; } + @SuppressWarnings("unchecked") + private List> gatherOutputBindings() { + List> outputBindings = new ArrayList<>(); + for (OutputBindingLifecycle inputBindingLifecycle : this.outputBindingsLifecycles) { + Collection> lifecycleInputBindings = (Collection>) new DirectFieldAccessor( + inputBindingLifecycle).getPropertyValue("outputBindings"); + outputBindings.addAll(lifecycleInputBindings); + } + return outputBindings; + } + private Binding locateBinding(String name) { - return BindingsEndpoint.this.gatherInputBindings().stream() + Stream> bindings = Stream.concat(this.gatherInputBindings().stream(), this.gatherOutputBindings().stream()); + return bindings .filter(binding -> name.equals(binding.getName())) .findFirst() .orElse(null);