From 8f8e1d1294389443074e2a52f5df9a307d24bf3c Mon Sep 17 00:00:00 2001 From: Thomas Badie Date: Tue, 3 Sep 2024 15:00:53 +0100 Subject: [PATCH] GH-2805: Add a checkAfterCompletion to SimpleMessageListenerContainer Fixes: #2805 Issue link: https://github.com/spring-projects/spring-amqp/issues/2805 Run `checkAfterCompletion` to the `SimpleMessageListenerContainer`. That allows to be notified when there is an issue with the commit. --- .../SimpleMessageListenerContainer.java | 4 + .../listener/ExternalTxManagerSMLCTests.java | 122 +++++++++++++++++- .../listener/ExternalTxManagerTests.java | 31 ++++- 3 files changed, 153 insertions(+), 4 deletions(-) diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java index 5e4ae9ce..f4f41b4d 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java @@ -83,6 +83,7 @@ import com.rabbitmq.client.ShutdownSignalException; * @author Tim Bourquin * @author Jeonggi Kim * @author Java4ye + * @author Thomas Badie * * @since 1.0 */ @@ -1013,6 +1014,9 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta catch (WrappedTransactionException e) { // NOSONAR exception flow control throw (Exception) e.getCause(); } + finally { + ConnectionFactoryUtils.checkAfterCompletion(); + } } return doReceiveAndExecute(consumer); diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerSMLCTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerSMLCTests.java index e1fb3e69..44350d83 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerSMLCTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerSMLCTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2019 the original author or authors. + * Copyright 2016-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,10 +16,42 @@ package org.springframework.amqp.rabbit.listener; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyMap; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + import org.springframework.amqp.rabbit.connection.AbstractConnectionFactory; +import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.context.ApplicationEventPublisher; + +import com.rabbitmq.client.AMQP; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.Connection; +import com.rabbitmq.client.ConnectionFactory; +import com.rabbitmq.client.Consumer; +import com.rabbitmq.client.Envelope; /** * @author Gary Russell + * @author Thomas Badie * @since 2.0 * */ @@ -32,4 +64,92 @@ public class ExternalTxManagerSMLCTests extends ExternalTxManagerTests { return container; } + + @Test + public void testMessageListenerTxFail() throws Exception { + ConnectionFactoryUtils.enableAfterCompletionFailureCapture(true); + ConnectionFactory mockConnectionFactory = mock(ConnectionFactory.class); + Connection mockConnection = mock(Connection.class); + final Channel mockChannel = mock(Channel.class); + given(mockChannel.isOpen()).willReturn(true); + given(mockChannel.txSelect()).willReturn(mock(AMQP.Tx.SelectOk.class)); + final AtomicReference commitLatch = new AtomicReference<>(new CountDownLatch(1)); + String exceptionMessage = "Failed to commit."; + willAnswer(invocation -> { + commitLatch.get().countDown(); + throw new IllegalStateException(exceptionMessage); + }).given(mockChannel).txCommit(); + + final CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory(mockConnectionFactory); + cachingConnectionFactory.setExecutor(mock(ExecutorService.class)); + given(mockConnectionFactory.newConnection(any(ExecutorService.class), anyString())).willReturn(mockConnection); + given(mockConnection.isOpen()).willReturn(true); + + willAnswer(invocation -> mockChannel).given(mockConnection).createChannel(); + + final AtomicReference consumer = new AtomicReference(); + final CountDownLatch consumerLatch = new CountDownLatch(1); + + willAnswer(invocation -> { + consumer.set(invocation.getArgument(6)); + consumerLatch.countDown(); + return "consumerTag"; + }).given(mockChannel) + .basicConsume(anyString(), anyBoolean(), anyString(), anyBoolean(), anyBoolean(), anyMap(), + any(Consumer.class)); + + + final CountDownLatch latch = new CountDownLatch(1); + AbstractMessageListenerContainer container = createContainer(cachingConnectionFactory); + container.setMessageListener(message -> { + RabbitTemplate rabbitTemplate = new RabbitTemplate(cachingConnectionFactory); + rabbitTemplate.setChannelTransacted(true); + // should use same channel as container + rabbitTemplate.convertAndSend("foo", "bar", "baz"); + latch.countDown(); + }); + container.setQueueNames("queue"); + container.setChannelTransacted(true); + container.setShutdownTimeout(100); + DummyTxManager transactionManager = new DummyTxManager(); + container.setTransactionManager(transactionManager); + ApplicationEventPublisher applicationEventPublisher = mock(ApplicationEventPublisher.class); + final CountDownLatch applicationEventPublisherLatch = new CountDownLatch(1); + willAnswer(invocation -> { + if (invocation.getArgument(0) instanceof ListenerContainerConsumerFailedEvent) { + applicationEventPublisherLatch.countDown(); + } + return null; + }).given(applicationEventPublisher).publishEvent(any()); + + container.setApplicationEventPublisher(applicationEventPublisher); + container.afterPropertiesSet(); + container.start(); + assertThat(consumerLatch.await(10, TimeUnit.SECONDS)).isTrue(); + + consumer.get().handleDelivery("qux", + new Envelope(1, false, "foo", "bar"), new AMQP.BasicProperties(), + new byte[] { 0 }); + + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + + verify(mockConnection, times(1)).createChannel(); + assertThat(commitLatch.get().await(10, TimeUnit.SECONDS)).isTrue(); + verify(mockChannel).basicAck(anyLong(), anyBoolean()); + verify(mockChannel).txCommit(); + + assertThat(applicationEventPublisherLatch.await(10, TimeUnit.SECONDS)).isTrue(); + verify(applicationEventPublisher).publishEvent(any(ListenerContainerConsumerFailedEvent.class)); + + ArgumentCaptor argumentCaptor + = ArgumentCaptor.forClass(ListenerContainerConsumerFailedEvent.class); + verify(applicationEventPublisher).publishEvent(argumentCaptor.capture()); + assertThat(argumentCaptor.getValue().getThrowable()).hasCauseInstanceOf(IllegalStateException.class); + assertThat(argumentCaptor.getValue().getThrowable()) + .isNotNull().extracting(Throwable::getCause) + .isNotNull().extracting(Throwable::getMessage).isEqualTo(exceptionMessage); + container.stop(); + } + + } diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java index ac4e0c86..dedb34a9 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/ExternalTxManagerTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2019 the original author or authors. + * Copyright 2002-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -69,6 +69,7 @@ import com.rabbitmq.client.Envelope; /** * @author Gary Russell + * @author Thomas Badie * @since 1.1.2 * */ @@ -758,7 +759,7 @@ public abstract class ExternalTxManagerTests { container.stop(); } - private Answer ensureOneChannelAnswer(final Channel onlyChannel, + protected Answer ensureOneChannelAnswer(final Channel onlyChannel, final AtomicReference tooManyChannels) { final AtomicBoolean done = new AtomicBoolean(); return invocation -> { @@ -776,7 +777,7 @@ public abstract class ExternalTxManagerTests { protected abstract AbstractMessageListenerContainer createContainer(AbstractConnectionFactory connectionFactory); @SuppressWarnings("serial") - private static class DummyTxManager extends AbstractPlatformTransactionManager { + protected static class DummyTxManager extends AbstractPlatformTransactionManager { private volatile boolean committed; @@ -804,6 +805,30 @@ public abstract class ExternalTxManagerTests { this.rolledBack = true; this.latch.countDown(); } + + public boolean isCommitted() { + return committed; + } + + public void setCommitted(boolean committed) { + this.committed = committed; + } + + public boolean isRolledBack() { + return rolledBack; + } + + public void setRolledBack(boolean rolledBack) { + this.rolledBack = rolledBack; + } + + public CountDownLatch getLatch() { + return latch; + } + + public void setLatch(CountDownLatch latch) { + this.latch = latch; + } } }