From 28894afed7dd02c7fc0a9969f378f6901f5eb048 Mon Sep 17 00:00:00 2001 From: Walliee Date: Fri, 5 Jul 2019 16:02:06 -0400 Subject: [PATCH] Fix pollable consumer auto-startup resolves #1758 --- .../binder/AbstractMessageChannelBinder.java | 2 +- .../stream/binder/PollableConsumerTests.java | 76 +++++++++++++++++++ 2 files changed, 77 insertions(+), 1 deletion(-) 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 1d78efb86..8a2e404dd 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 @@ -487,7 +487,7 @@ public abstract class AbstractMessageChannelBinder> binding = new DefaultBinding>( diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java index b0f9fdd5d..980f2fbab 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/PollableConsumerTests.java @@ -35,13 +35,16 @@ import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.context.ApplicationContext; +import org.springframework.context.Lifecycle; import org.springframework.core.ParameterizedTypeReference; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.acks.AcknowledgmentCallback; import org.springframework.integration.acks.AcknowledgmentCallback.Status; import org.springframework.integration.context.IntegrationContextUtils; +import org.springframework.integration.core.MessageSource; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.converter.SmartMessageConverter; @@ -411,6 +414,50 @@ public class PollableConsumerTests { verify(callback).acknowledge(Status.REQUEUE); } + @SuppressWarnings("unchecked") + @Test + public void testAutoStartupOff() { + TestChannelBinder binder = createBinder(); + binder.setMessageSourceDelegate(new LifecycleMessageSource( + () -> new GenericMessage<>("{\"foo\":\"bar\"}".getBytes()))); + MessageConverterConfigurer configurer = this.context + .getBean(MessageConverterConfigurer.class); + + DefaultPollableMessageSource pollableSource = new DefaultPollableMessageSource( + this.messageConverter); + configurer.configurePolledMessageSource(pollableSource, "foo"); + ExtendedConsumerProperties properties = new ExtendedConsumerProperties<>( + null); + properties.setAutoStartup(false); + + Binding> pollableSourceBinding = binder + .bindPollableConsumer("foo", "bar", pollableSource, properties); + + assertThat(pollableSourceBinding.isRunning()).isFalse(); + } + + @SuppressWarnings("unchecked") + @Test + public void testAutoStartupOn() { + TestChannelBinder binder = createBinder(); + binder.setMessageSourceDelegate(new LifecycleMessageSource( + () -> new GenericMessage<>("{\"foo\":\"bar\"}".getBytes()))); + MessageConverterConfigurer configurer = this.context + .getBean(MessageConverterConfigurer.class); + + DefaultPollableMessageSource pollableSource = new DefaultPollableMessageSource( + this.messageConverter); + configurer.configurePolledMessageSource(pollableSource, "foo"); + ExtendedConsumerProperties properties = new ExtendedConsumerProperties<>( + null); + properties.setAutoStartup(true); + + Binding> pollableSourceBinding = binder + .bindPollableConsumer("foo", "bar", pollableSource, properties); + + assertThat(pollableSourceBinding.isRunning()).isTrue(); + } + private TestChannelBinder createBinder(String... args) { this.context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration()) @@ -433,4 +480,33 @@ public class PollableConsumerTests { } + public static class LifecycleMessageSource implements MessageSource, Lifecycle { + private final MessageSource delegate; + private boolean running = false; + + public LifecycleMessageSource(MessageSource delegate) { + this.delegate = delegate; + } + + @Override + public Message receive() { + return this.delegate.receive(); + } + + @Override + public void start() { + this.running = true; + } + + @Override + public void stop() { + this.running = false; + } + + @Override + public boolean isRunning() { + return this.running; + } + } + }