AMQP-797: Defer caching publisher callback channel

JIRA: https://jira.spring.io/browse/AMQP-797

Additional test case (Ack Vs. Nack).
This commit is contained in:
Gary Russell
2018-05-30 19:40:29 -04:00
parent 8a38e6550b
commit 8d41b89216

View File

@@ -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;