From 9e710f9fa80bc7ad8b97aa493e82f1ad02ef973c Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Tue, 26 Jul 2011 13:45:26 +0100 Subject: [PATCH] AMQP-183: tweak synchronized block to try and avoid race where channel is closed before it is used --- .../connection/CachingConnectionFactory.java | 19 ++++++++++--------- ...veryCachingConnectionIntegrationTests.java | 2 ++ 2 files changed, 12 insertions(+), 9 deletions(-) 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 227731ae..63f51861 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 @@ -284,25 +284,26 @@ public class CachingConnectionFactory extends AbstractConnectionFactory { // Handle getTargetChannel method: return underlying Channel. return this.target; } else if (methodName.equals("isOpen")) { - // Handle isOpen method: we are closed if the target is + // Handle isOpen method: we are closed if the target is closed return this.target != null && this.target.isOpen(); } try { if (this.target == null || !this.target.isOpen()) { this.target = null; - synchronized (targetMonitor) { - if (this.target == null) { - this.target = createBareChannel(transactional); - } - } } - return method.invoke(this.target, args); + synchronized (targetMonitor) { + if (this.target == null) { + this.target = createBareChannel(transactional); + } + return method.invoke(this.target, args); + } } catch (InvocationTargetException ex) { - if (!this.target.isOpen()) { + if (this.target == null || !this.target.isOpen()) { // Basic re-connection logic... + this.target = null; logger.debug("Detected closed channel on exception. Re-initializing: " + target); synchronized (targetMonitor) { - if (!this.target.isOpen()) { + if (this.target == null) { this.target = createBareChannel(transactional); } } diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerRecoveryCachingConnectionIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerRecoveryCachingConnectionIntegrationTests.java index f3c2f387..921b7fcc 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerRecoveryCachingConnectionIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/MessageListenerRecoveryCachingConnectionIntegrationTests.java @@ -240,6 +240,8 @@ public class MessageListenerRecoveryCachingConnectionIntegrationTests { container = createContainer(queue.getName(), new ManualAckListener(latch), connectionFactory); for (int i = 0; i < messageCount; i++) { template.convertAndSend(queue.getName(), i + "foo"); + // Give the listener container a chance to steal the connection from the template + Thread.sleep(200); } int timeout = getTimeout();