diff --git a/CONTRIBUTING.adoc b/CONTRIBUTING.adoc
index 23425ec6..849b9516 100644
--- a/CONTRIBUTING.adoc
+++ b/CONTRIBUTING.adoc
@@ -51,7 +51,7 @@ _you should see branches on origin as well as upstream, including 'main' and 'ma
== A Day in the Life of a Contributor
-* _Always_ work on topic branches (Typically use the HitHub (or JIRA) issue ID as the branch name).
+* _Always_ work on topic branches (Typically use the GitHub (or JIRA) issue ID as the branch name).
- For example, to create and switch to a new branch for issue #123: `git checkout -b GH-123`
* You might be working on several different topic branches at any given time, but when at a stopping point for one of those branches, commit (a local operation).
* Please follow the "Commit Guidelines" described in https://git-scm.com/book/en/Distributed-Git-Contributing-to-a-Project[this chapter of Pro Git].
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 2ea39028..fdf057a3 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
@@ -54,6 +54,8 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact
private Boolean returns;
+ private boolean consistentConfirmsReturns = true;
+
/**
* Specify the map of target ConnectionFactories, with the lookup key as key.
*
The key can be of arbitrary type; this class implements the
@@ -125,10 +127,13 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact
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");
+
+ if (this.consistentConfirmsReturns) {
+ 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
@@ -230,6 +235,25 @@ public abstract class AbstractRoutingConnectionFactory implements ConnectionFact
return this.targetConnectionFactories.get(key);
}
+ /**
+ * Specify whether to apply a validation enforcing all {@link ConnectionFactory#isPublisherConfirms()} and
+ * {@link ConnectionFactory#isPublisherReturns()} have a consistent value.
+ *
+ * A consistent value means that all ConnectionFactories must have the same value between all
+ * {@link ConnectionFactory#isPublisherConfirms()} and the same value between all
+ * {@link ConnectionFactory#isPublisherReturns()}.
+ *
+ *
+ * Note that in any case the values between {@link ConnectionFactory#isPublisherConfirms()} and
+ * {@link ConnectionFactory#isPublisherReturns()} don't need to be equals between each other.
+ *
+ * @param consistentConfirmsReturns true to validate, false to not validate.
+ * @since 2.4.4
+ */
+ public void setConsistentConfirmsReturns(boolean consistentConfirmsReturns) {
+ this.consistentConfirmsReturns = consistentConfirmsReturns;
+ }
+
/**
* Adds the given {@link ConnectionFactory} and associates it with the given lookup key.
* @param key the lookup key.
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 15504f22..78e6e5e9 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
@@ -94,6 +94,7 @@ import com.rabbitmq.client.impl.recovery.AutorecoveringChannel;
* @author Artem Bilan
* @author Steve Powell
* @author Will Droste
+ * @author Leonardo Ferreira
*/
@ManagedResource
public class CachingConnectionFactory extends AbstractConnectionFactory
@@ -1133,6 +1134,9 @@ public class CachingConnectionFactory extends AbstractConnectionFactory
else if (methodName.equals("isConfirmSelected")) {
return this.confirmSelected;
}
+ else if (methodName.equals("isPublisherConfirms")) {
+ return this.publisherConfirms;
+ }
try {
if (this.target == null || !this.target.isOpen()) {
if (this.target instanceof PublisherCallbackChannel) {
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ChannelProxy.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ChannelProxy.java
index 17a423b7..e0afe0fa 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ChannelProxy.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ChannelProxy.java
@@ -54,4 +54,12 @@ public interface ChannelProxy extends Channel, RawTargetAccess {
return false;
}
+ /**
+ * Return true if publisher confirms are enabled.
+ * @return true if publisherConfirms.
+ */
+ default boolean isPublisherConfirms() {
+ return false;
+ }
+
}
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java
index 15add409..186e0ca3 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/PooledChannelConnectionFactory.java
@@ -222,6 +222,8 @@ public class PooledChannelConnectionFactory extends AbstractConnectionFactory im
return channel.confirmSelect();
case "isConfirmSelected":
return confirmSelected.get();
+ case "isPublisherConfirms":
+ return false;
}
return null;
};
@@ -231,6 +233,7 @@ public class PooledChannelConnectionFactory extends AbstractConnectionFactory im
advisor.addMethodName("isTransactional");
advisor.addMethodName("confirmSelect");
advisor.addMethodName("isConfirmSelected");
+ advisor.addMethodName("isPublisherConfirms");
pf.addAdvisor(advisor);
pf.addInterface(ChannelProxy.class);
proxy.set((Channel) pf.getProxy());
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java
index 0da501fe..69edbc64 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/connection/ThreadChannelConnectionFactory.java
@@ -44,6 +44,7 @@ import com.rabbitmq.client.ShutdownListener;
* {@link #closeThreadChannel()}.
*
* @author Gary Russell
+ * @author Leonardo Ferreira
* @since 2.3
*
*/
@@ -288,6 +289,8 @@ public class ThreadChannelConnectionFactory extends AbstractConnectionFactory im
return channel.confirmSelect();
case "isConfirmSelected":
return confirmSelected.get();
+ case "isPublisherConfirms":
+ return false;
}
return null;
};
@@ -297,6 +300,7 @@ public class ThreadChannelConnectionFactory extends AbstractConnectionFactory im
advisor.addMethodName("isTransactional");
advisor.addMethodName("confirmSelect");
advisor.addMethodName("isConfirmSelected");
+ advisor.addMethodName("isPublisherConfirms");
pf.addAdvisor(advisor);
pf.addInterface(ChannelProxy.class);
return (Channel) pf.getProxy();
diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java
index fa0f8d47..bca2597b 100644
--- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java
+++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java
@@ -147,6 +147,7 @@ import com.rabbitmq.client.ShutdownSignalException;
* @author Mark Norkin
* @author Mohammad Hewedy
* @author Alexey Platonov
+ * @author Leonardo Ferreira
*
* @since 1.0
*/
@@ -257,10 +258,6 @@ public class RabbitTemplate extends RabbitAccessor // NOSONAR type line count
private ErrorHandler replyErrorHandler;
- private volatile Boolean confirmsOrReturnsCapable;
-
- private volatile boolean publisherConfirms;
-
private volatile boolean usingFastReplyTo;
private volatile boolean evaluatedFastReplyTo;
@@ -1263,7 +1260,7 @@ public class RabbitTemplate extends RabbitAccessor // NOSONAR type line count
}
return buildMessageFromDelivery(delivery);
}
- });
+ }, obtainTargetConnectionFactory(this.receiveConnectionFactorySelectorExpression, null));
logReceived(message);
return message;
}
@@ -1960,7 +1957,7 @@ public class RabbitTemplate extends RabbitAccessor // NOSONAR type line count
boolean cancelConsumer = false;
try {
Channel channel = channelHolder.getChannel();
- if (this.confirmsOrReturnsCapable) {
+ if (isPublisherConfirmsOrReturns(connectionFactory)) {
addListener(channel);
}
Message reply = doSendAndReceiveAsListener(exchange, routingKey, message, correlationData, channel,
@@ -2224,12 +2221,10 @@ public class RabbitTemplate extends RabbitAccessor // NOSONAR type line count
private T invokeAction(ChannelCallback action, ConnectionFactory connectionFactory, Channel channel)
throws Exception { // NOSONAR see the callback
- if (this.confirmsOrReturnsCapable == null) {
- determineConfirmsReturnsCapability(connectionFactory);
- }
- if (this.confirmsOrReturnsCapable) {
+ if (isPublisherConfirmsOrReturns(connectionFactory)) {
addListener(channel);
}
+
if (logger.isDebugEnabled()) {
logger.debug(
"Executing callback " + action.getClass().getSimpleName() + " on RabbitMQ Channel: " + channel);
@@ -2351,10 +2346,8 @@ public class RabbitTemplate extends RabbitAccessor // NOSONAR type line count
}
}
- public void determineConfirmsReturnsCapability(ConnectionFactory connectionFactory) {
- this.publisherConfirms = connectionFactory.isPublisherConfirms();
- this.confirmsOrReturnsCapable =
- this.publisherConfirms || connectionFactory.isPublisherReturns();
+ private boolean isPublisherConfirmsOrReturns(ConnectionFactory connectionFactory) {
+ return connectionFactory.isPublisherConfirms() || connectionFactory.isPublisherReturns();
}
/**
@@ -2435,8 +2428,11 @@ public class RabbitTemplate extends RabbitAccessor // NOSONAR type line count
}
private void setupConfirm(Channel channel, Message message, @Nullable CorrelationData correlationDataArg) {
- if ((this.publisherConfirms || this.confirmCallback != null) && channel instanceof PublisherCallbackChannel) {
+ final boolean publisherConfirms = channel instanceof ChannelProxy
+ && ((ChannelProxy) channel).isPublisherConfirms();
+ if ((publisherConfirms || this.confirmCallback != null)
+ && channel instanceof PublisherCallbackChannel) {
long nextPublishSeqNo = channel.getNextPublishSeqNo();
if (nextPublishSeqNo > 0) {
PublisherCallbackChannel publisherCallbackChannel = (PublisherCallbackChannel) channel;
diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateRoutingConnectionFactoryIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateRoutingConnectionFactoryIntegrationTests.java
new file mode 100644
index 00000000..b2ebbbb7
--- /dev/null
+++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitTemplateRoutingConnectionFactoryIntegrationTests.java
@@ -0,0 +1,123 @@
+/*
+ * Copyright 2002-2022 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.
+ * You may obtain a copy of the License at
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.amqp.rabbit.core;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.TimeUnit;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.core.MessageBuilder;
+import org.springframework.amqp.rabbit.connection.AbstractRoutingConnectionFactory;
+import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
+import org.springframework.amqp.rabbit.connection.ConnectionFactory;
+import org.springframework.amqp.rabbit.connection.CorrelationData;
+import org.springframework.amqp.rabbit.connection.PooledChannelConnectionFactory;
+import org.springframework.amqp.rabbit.connection.SimpleRoutingConnectionFactory;
+import org.springframework.amqp.rabbit.junit.BrokerTestUtils;
+import org.springframework.amqp.rabbit.junit.RabbitAvailable;
+import org.springframework.expression.Expression;
+import org.springframework.expression.spel.standard.SpelExpressionParser;
+
+/**
+ * @author Leonardo Ferreira
+ * @since 2.4.4
+ */
+@RabbitAvailable(queues = RabbitTemplateRoutingConnectionFactoryIntegrationTests.ROUTE)
+class RabbitTemplateRoutingConnectionFactoryIntegrationTests {
+
+ public static final String ROUTE = "test.queue.RabbitTemplateRoutingConnectionFactoryIntegrationTests";
+
+ private static RabbitTemplate rabbitTemplate;
+
+ @BeforeAll
+ static void create() {
+ final com.rabbitmq.client.ConnectionFactory cf = new com.rabbitmq.client.ConnectionFactory();
+ cf.setHost("localhost");
+ cf.setPort(BrokerTestUtils.getPort());
+
+ CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory(cf);
+
+ cachingConnectionFactory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED);
+
+ PooledChannelConnectionFactory pooledChannelConnectionFactory = new PooledChannelConnectionFactory(cf);
+
+ Map