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
This commit is contained in:
committed by
Oleg Zhurakousky
parent
fe2a099e5b
commit
bdfd4a89eb
@@ -105,7 +105,7 @@ public class RabbitMessageChannelBinder
|
||||
extends AbstractMessageChannelBinder<ExtendedConsumerProperties<RabbitConsumerProperties>,
|
||||
ExtendedProducerProperties<RabbitProducerProperties>, RabbitExchangeQueueProvisioner>
|
||||
implements ExtendedPropertiesBinder<MessageChannel, RabbitConsumerProperties, RabbitProducerProperties>,
|
||||
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);
|
||||
|
||||
@@ -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<AsyncConsumerStartedEvent> event = new AtomicReference<>();
|
||||
binder.getApplicationContext().addApplicationListener(
|
||||
(ApplicationListener<AsyncConsumerStartedEvent>) e -> event.set(e));
|
||||
DirectChannel moduleOutputChannel = createBindableChannel("output", new BindingProperties());
|
||||
DirectChannel moduleInputChannel = createBindableChannel("input", new BindingProperties());
|
||||
Binding<MessageChannel> 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();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user