From bdfd4a89ebfebcda8383adacc586753904ee432b Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 6 Mar 2018 17:05:19 -0500 Subject: [PATCH] GH-134: Propagate the event publisher to container Fixes https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/134 Requires https://github.com/spring-cloud/spring-cloud-stream/pull/1282 Resolves #135 --- .../stream/binder/rabbit/RabbitMessageChannelBinder.java | 8 +++++++- .../cloud/stream/binder/rabbit/RabbitBinderTests.java | 6 ++++++ 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index 2ebb20174..ed8c24817 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -105,7 +105,7 @@ public class RabbitMessageChannelBinder extends AbstractMessageChannelBinder, ExtendedProducerProperties, RabbitExchangeQueueProvisioner> implements ExtendedPropertiesBinder, - DisposableBean { + DisposableBean { private static final AmqpMessageHeaderErrorMessageStrategy errorMessageStrategy = new AmqpMessageHeaderErrorMessageStrategy(); @@ -374,6 +374,12 @@ public class RabbitMessageChannelBinder listenerContainer.setFailedDeclarationRetryInterval( properties.getExtension().getFailedDeclarationRetryInterval()); } + if (getApplicationEventPublisher() != null) { + listenerContainer.setApplicationEventPublisher(getApplicationEventPublisher()); + } + else if (getApplicationContext() != null) { + listenerContainer.setApplicationEventPublisher(getApplicationContext()); + } listenerContainer.afterPropertiesSet(); AmqpInboundChannelAdapter adapter = new AmqpInboundChannelAdapter(listenerContainer); diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java index 6c2c4ab56..8113456fe 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java @@ -50,6 +50,7 @@ import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitAdmin; import org.springframework.amqp.rabbit.core.RabbitManagementTemplate; import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.amqp.rabbit.listener.AsyncConsumerStartedEvent; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.amqp.support.postprocessor.DelegatingDecompressingPostProcessor; @@ -75,6 +76,7 @@ import org.springframework.cloud.stream.binder.rabbit.provisioning.RabbitExchang import org.springframework.cloud.stream.binder.test.junit.rabbit.RabbitTestSupport; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationListener; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.Lifecycle; import org.springframework.expression.spel.standard.SpelExpression; @@ -162,6 +164,9 @@ public class RabbitBinderTests extends @Test public void testSendAndReceiveBad() throws Exception { RabbitTestBinder binder = getBinder(); + final AtomicReference event = new AtomicReference<>(); + binder.getApplicationContext().addApplicationListener( + (ApplicationListener) e -> event.set(e)); DirectChannel moduleOutputChannel = createBindableChannel("output", new BindingProperties()); DirectChannel moduleInputChannel = createBindableChannel("input", new BindingProperties()); Binding producerBinding = binder.bindProducer("bad.0", moduleOutputChannel, @@ -180,6 +185,7 @@ public class RabbitBinderTests extends }); moduleOutputChannel.send(message); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(event.get()).isNotNull(); producerBinding.unbind(); consumerBinding.unbind(); }