From 01afdd6abf9ffee5aad5b3315909e82aa7c0ef8a Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Sat, 4 Jun 2011 10:23:20 +0100 Subject: [PATCH] AMQP-171: fix recovery on connection close - need to clear delivery tags more carefully --- .../amqp/rabbit/listener/BlockingQueueConsumer.java | 2 ++ ...essageListenerRecoveryCachingConnectionIntegrationTests.java | 2 +- 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java index 7ba62acc..a8713d58 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/BlockingQueueConsumer.java @@ -227,6 +227,8 @@ public class BlockingQueueConsumer { logger.debug("Received shutdown signal for consumer tag=" + consumerTag, sig); } shutdown = sig; + // The delivery tags will be invalid if the channel shuts down + deliveryTags.clear(); } @Override 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 7916fdac..f3c2f387 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 @@ -53,7 +53,7 @@ public class MessageListenerRecoveryCachingConnectionIntegrationTests { private SimpleMessageListenerContainer container; @Rule - public Log4jLevelAdjuster logLevels = new Log4jLevelAdjuster(Level.INFO, RabbitTemplate.class, + public Log4jLevelAdjuster logLevels = new Log4jLevelAdjuster(Level.DEBUG, RabbitTemplate.class, SimpleMessageListenerContainer.class, BlockingQueueConsumer.class, CachingConnectionFactory.class); @Rule