diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java index 96c7bdf4..fbdf5695 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactory.java @@ -549,6 +549,7 @@ public class CachingConnectionFactory extends AbstractConnectionFactory if (!this.checkoutPermits.containsKey(this.connection)) { this.checkoutPermits.put(this.connection, new Semaphore(this.channelCacheSize)); } + this.connection.closeNotified.set(false); getConnectionListener().onCreate(this.connection); } return this.connection; diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactoryTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactoryTests.java index d5b06a38..fdcf8956 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactoryTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/connection/CachingConnectionFactoryTests.java @@ -817,6 +817,7 @@ public class CachingConnectionFactoryTests extends AbstractConnectionFactoryTest final AtomicReference created = new AtomicReference(); final AtomicReference closed = new AtomicReference(); + final AtomicInteger timesClosed = new AtomicInteger(0); AbstractConnectionFactory connectionFactory = createConnectionFactory(mockConnectionFactory); connectionFactory.addConnectionListener(new ConnectionListener() { @@ -828,6 +829,7 @@ public class CachingConnectionFactoryTests extends AbstractConnectionFactoryTest @Override public void onClose(Connection connection) { closed.set(connection); + timesClosed.getAndAdd(1); } }); ((CachingConnectionFactory) connectionFactory).setChannelCacheSize(1); @@ -850,6 +852,7 @@ public class CachingConnectionFactoryTests extends AbstractConnectionFactoryTest when(mockChannel.isOpen()).thenReturn(false); // force a connection refresh channel.basicCancel("foo"); channel.close(); + assertEquals(1, timesClosed.get()); Connection notSame = connectionFactory.createConnection(); assertNotSame(conDelegate, targetDelegate(notSame)); @@ -859,6 +862,7 @@ public class CachingConnectionFactoryTests extends AbstractConnectionFactoryTest connectionFactory.destroy(); verify(mockConnection2, atLeastOnce()).close(anyInt()); assertSame(notSame, closed.get()); + assertEquals(2, timesClosed.get()); verify(mockConnectionFactory, times(2)).newConnection((ExecutorService) null); }