From ad612a8aeebf28be90659ced8f413ef61d0fd02c Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 23 Dec 2020 12:52:41 -0500 Subject: [PATCH] GH-1289: Confirms and Returns with Routing CF Resolves https://github.com/spring-projects/spring-amqp/issues/1289 `RoutingConnectionFactory` did not support correlated confirms or returns. Target factories (and default) must have the same settings. **cherry-pick to 2.2.x, 2.1.x** GH-1289: Fix test for back port - `CorrelationData` needs an id (`null` by default before 2.3). --- .../AbstractRoutingConnectionFactory.java | 42 +++++++++++++-- ...tePublisherCallbacksIntegrationTests2.java | 54 +++++++++++++++++-- src/reference/asciidoc/amqp.adoc | 7 ++- 3 files changed, 95 insertions(+), 8 deletions(-) diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java index e2a144fd..01f12c90 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/AbstractRoutingConnectionFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2018 the original author or authors. + * Copyright 2002-2020 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -22,6 +22,7 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import org.springframework.amqp.AmqpException; +import org.springframework.beans.factory.InitializingBean; import org.springframework.lang.Nullable; import org.springframework.util.Assert; @@ -35,7 +36,8 @@ import org.springframework.util.Assert; * @author Gary Russell * @since 1.3 */ -public abstract class AbstractRoutingConnectionFactory implements ConnectionFactory, RoutingConnectionFactory { +public abstract class AbstractRoutingConnectionFactory implements ConnectionFactory, RoutingConnectionFactory, + InitializingBean { private final Map targetConnectionFactories = new ConcurrentHashMap(); @@ -46,6 +48,10 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact private boolean lenientFallback = true; + private Boolean confirms; + + private Boolean returns; + /** * Specify the map of target ConnectionFactories, with the lookup key as key. *

The key can be of arbitrary type; this class implements the @@ -58,6 +64,7 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact Assert.noNullElements(targetConnectionFactories.values().toArray(), "'targetConnectionFactories' cannot have null values."); this.targetConnectionFactories.putAll(targetConnectionFactories); + targetConnectionFactories.values().stream().forEach(cf -> checkConfirmsAndReturns(cf)); } /** @@ -69,6 +76,7 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact */ public void setDefaultTargetConnectionFactory(ConnectionFactory defaultTargetConnectionFactory) { this.defaultTargetConnectionFactory = defaultTargetConnectionFactory; + checkConfirmsAndReturns(defaultTargetConnectionFactory); } /** @@ -93,9 +101,37 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact return this.lenientFallback; } + @Override + public boolean isPublisherConfirms() { + return this.confirms; + } + + @Override + public boolean isPublisherReturns() { + return this.returns; + } + + @Override + public void afterPropertiesSet() throws Exception { + Assert.notNull(this.confirms, "At least one target factory (or default) is required"); + } + + private void checkConfirmsAndReturns(ConnectionFactory cf) { + if (this.confirms == null) { + this.confirms = cf.isPublisherConfirms(); + } + if (this.returns == null) { + this.returns = cf.isPublisherReturns(); + } + Assert.isTrue(this.confirms.booleanValue() == cf.isPublisherConfirms(), + "Target connection factories must have the same setting for publisher confirms"); + Assert.isTrue(this.returns.booleanValue() == cf.isPublisherReturns(), + "Target connection factories must have the same setting for publisher returns"); + } + @Override public Connection createConnection() throws AmqpException { - return this.determineTargetConnectionFactory().createConnection(); + return determineTargetConnectionFactory().createConnection(); } /** diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests2.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests2.java index 52ccd877..027034d3 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests2.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplatePublisherCallbacksIntegrationTests2.java @@ -1,5 +1,5 @@ /* - * Copyright 2016 the original author or authors. + * Copyright 2016-2020 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -18,6 +18,7 @@ package org.springframework.amqp.rabbit.core; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; +import static org.junit.jupiter.api.Assertions.assertNotNull; import java.io.IOException; import java.util.concurrent.CountDownLatch; @@ -29,6 +30,8 @@ import org.junit.Rule; import org.junit.Test; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.rabbit.connection.CorrelationData; +import org.springframework.amqp.rabbit.connection.SimpleRoutingConnectionFactory; import org.springframework.amqp.rabbit.junit.BrokerRunning; import org.springframework.amqp.rabbit.junit.BrokerTestUtils; @@ -45,6 +48,8 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests2 { private static final String ROUTE = "test.queue"; + public static final String ROUTE2 = "test.queue.RabbitTemplatePublisherCallbacksIntegrationTests2.route"; + private CachingConnectionFactory connectionFactoryWithConfirmsEnabled; private RabbitTemplate templateWithConfirmsEnabled; @@ -56,8 +61,6 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests2 { public void create() { connectionFactoryWithConfirmsEnabled = new CachingConnectionFactory(); connectionFactoryWithConfirmsEnabled.setHost("localhost"); - // When using publisher confirms, the cache size needs to be large enough - // otherwise channels can be closed before confirms are received. connectionFactoryWithConfirmsEnabled.setChannelCacheSize(100); connectionFactoryWithConfirmsEnabled.setPort(BrokerTestUtils.getPort()); connectionFactoryWithConfirmsEnabled.setPublisherConfirms(true); @@ -94,6 +97,51 @@ public class RabbitTemplatePublisherCallbacksIntegrationTests2 { assertMessageCountEquals(0L); } + @Test + public void routingWithConfirmsNoListener() throws Exception { + routingWithConfirms(false); + } + + @Test + public void routingWithConfirmsListener() throws Exception { + routingWithConfirms(true); + } + + private void routingWithConfirms(boolean listener) throws Exception { + CountDownLatch latch = new CountDownLatch(1); + SimpleRoutingConnectionFactory rcf = new SimpleRoutingConnectionFactory(); + rcf.setDefaultTargetConnectionFactory(this.connectionFactoryWithConfirmsEnabled); + this.templateWithConfirmsEnabled.setConnectionFactory(rcf); + if (listener) { + this.templateWithConfirmsEnabled.setConfirmCallback((correlationData, ack, cause) -> { + latch.countDown(); + }); + } + this.templateWithConfirmsEnabled.setMandatory(true); + CorrelationData corr = new CorrelationData("foo"); + this.templateWithConfirmsEnabled.convertAndSend("", ROUTE2, "foo", corr); + assertTrue(corr.getFuture().get(10, TimeUnit.SECONDS).isAck()); + if (listener) { + assertTrue(latch.await(10, TimeUnit.SECONDS)); + } + corr = new CorrelationData("bar"); + this.templateWithConfirmsEnabled.convertAndSend("", "bad route", "foo", corr); + assertTrue(corr.getFuture().get(10, TimeUnit.SECONDS).isAck()); + assertNotNull(corr.getReturnedMessage()); + } + + @Test + public void routingWithSimpleConfirms() throws Exception { + SimpleRoutingConnectionFactory rcf = new SimpleRoutingConnectionFactory(); + rcf.setDefaultTargetConnectionFactory(this.connectionFactoryWithConfirmsEnabled); + this.templateWithConfirmsEnabled.setConnectionFactory(rcf); + assertTrue(this.templateWithConfirmsEnabled.invoke(template -> { + template.convertAndSend("", ROUTE2, "foo"); + template.waitForConfirmsOrDie(10_000); + return true; + })); + } + private void assertMessageCountEquals(long wanted) throws InterruptedException { long messageCount = determineMessageCount(); int n = 0; diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index 77786f0a..39722be3 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -627,6 +627,9 @@ Doing so enables, for example, listening to queues with the same name but in a d For example, with lookup key qualifier `thing1` and a container listening to queue `thing2`, the lookup key you could register the target connection factory with could be `thing1[thing2]`. +IMPORTANT: The target (and default, if provided) connection factories must have the same settings for publisher confirms and returns. +See <>. + [[queue-affinity]] ===== Queue Affinity and the `LocalizedQueueConnectionFactory` @@ -1115,8 +1118,8 @@ The `Confirm` object is a simple bean with 2 properties: `ack` and `reason` (for The reason is not populated for broker-generated `nack` instances. It is populated for `nack` instances generated by the framework (for example, closing the connection while `ack` instances are outstanding). -In addition, when both confirms and returns are enabled, the `CorrelationData` is populated with the returned message. -It is guaranteed that this occurs before the future is set with the `ack`. +In addition, when both confirms and returns are enabled, the `CorrelationData` is populated with the returned message, as long as the `CorrelationData` has a unique `id`; this is always the case, by default, starting with version 2.3. +It is guaranteed that the returned message is set before the future is set with the `ack`. See also <> for a simpler mechanism for waiting for publisher confirms.