From 8d41b89216195d2cc541315351be625b270a3dcd Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 30 May 2018 19:40:29 -0400 Subject: [PATCH] AMQP-797: Defer caching publisher callback channel JIRA: https://jira.spring.io/browse/AMQP-797 Additional test case (Ack Vs. Nack). --- ...tePublisherCallbacksIntegrationTests3.java | 28 ++++++++++++++++++- 1 file changed, 27 insertions(+), 1 deletion(-) diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests3.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests3.java index 3aea8909..6efdfcb7 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests3.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests3.java @@ -66,7 +66,7 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests3 { } @Test - public void testDeferredChannelCache() throws Exception { + public void testDeferredChannelCacheNack() throws Exception { final CachingConnectionFactory cf = new CachingConnectionFactory( RabbitAvailableCondition.getBrokerRunning().getConnectionFactory()); cf.setPublisherReturns(true); @@ -97,6 +97,32 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests3 { cf.destroy(); } + @Test + public void testDeferredChannelCacheAck() throws Exception { + final CachingConnectionFactory cf = new CachingConnectionFactory( + RabbitAvailableCondition.getBrokerRunning().getConnectionFactory()); + cf.setPublisherConfirms(true); + final RabbitTemplate template = new RabbitTemplate(cf); + final CountDownLatch confirmLatch = new CountDownLatch(1); + final AtomicInteger cacheCount = new AtomicInteger(); + template.setConfirmCallback((cd, a, c) -> { + cacheCount.set(TestUtils.getPropertyValue(cf, "cachedChannelsNonTransactional", List.class).size()); + confirmLatch.countDown(); + }); + template.setMandatory(true); + Connection conn = cf.createConnection(); + Channel channel1 = conn.createChannel(false); + Channel channel2 = conn.createChannel(false); + channel1.close(); + channel2.close(); + conn.close(); + assertThat(TestUtils.getPropertyValue(cf, "cachedChannelsNonTransactional", List.class).size()).isEqualTo(2); + template.convertAndSend("", QUEUE2, "foo", new MyCD("foo")); + assertThat(confirmLatch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(cacheCount.get()).isEqualTo(1); + cf.destroy(); + } + private static class MyCD extends CorrelationData { private final String payload;