From 83c499f1cf356a9840b034ceddaf406b8e0b9b70 Mon Sep 17 00:00:00 2001 From: Marcos Costa Pinto Date: Mon, 25 Jul 2016 12:59:08 -0300 Subject: [PATCH] AMQP-625: CCF: Fix onClose Notification JIRA: https://jira.spring.io/browse/AMQP-625 Previouly using CachingConnectionFactory with CacheMode.CHANNEL the connection listener `onClose` method was being notified just on the first time that the connection is closed. Now every time when the connection is closed the connection listener is notified. --- .../amqp/rabbit/connection/CachingConnectionFactory.java | 1 + .../amqp/rabbit/connection/CachingConnectionFactoryTests.java | 4 ++++ 2 files changed, 5 insertions(+) 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); }