From f7146ad02eedae2ed8c8199b766f66c3f5576f55 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 17 Aug 2022 12:41:28 -0400 Subject: [PATCH] GH-2483: RMQ: Fix Redeclare Multi Routing Keys Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2483 https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/242 added support for binding a queue with multiple routing keys. However, only one binding was added to the `autoRedeclarContext`. Use the `Declarables` wrapper to support multiple instances with the same name. **cherry-pick to 3.2.x** Resolves #2484 --- .../provisioning/RabbitExchangeQueueProvisioner.java | 9 +++++++-- .../cloud/stream/binder/rabbit/RabbitBinderTests.java | 7 ++++++- 2 files changed, 13 insertions(+), 3 deletions(-) diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java index b0a036b95..87750dd82 100644 --- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java +++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java @@ -33,7 +33,9 @@ import org.springframework.amqp.core.Base64UrlNamingStrategy; import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.Binding.DestinationType; import org.springframework.amqp.core.BindingBuilder; +import org.springframework.amqp.core.Declarable; import org.springframework.amqp.core.DeclarableCustomizer; +import org.springframework.amqp.core.Declarables; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.Exchange; import org.springframework.amqp.core.ExchangeBuilder; @@ -632,10 +634,13 @@ public class RabbitExchangeQueueProvisioner addToAutoDeclareContext(rootName + "." + group + ".exchange", exchange); } - private void addToAutoDeclareContext(String name, Object bean) { + private void addToAutoDeclareContext(String name, Declarable bean) { synchronized (this.autoDeclareContext) { if (!this.autoDeclareContext.containsBean(name)) { - this.autoDeclareContext.getBeanFactory().registerSingleton(name, bean); + this.autoDeclareContext.getBeanFactory().registerSingleton(name, new Declarables(bean)); + } + else { + this.autoDeclareContext.getBean(name, Declarables.class).getDeclarables().add(bean); } } } diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java index 1ee65618d..1148e743e 100644 --- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java +++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java @@ -59,6 +59,7 @@ import org.springframework.amqp.core.AnonymousQueue; import org.springframework.amqp.core.Binding.DestinationType; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.Declarable; +import org.springframework.amqp.core.Declarables; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.ExchangeTypes; import org.springframework.amqp.core.MessageDeliveryMode; @@ -163,7 +164,7 @@ public class RabbitBinderTests extends private static final String BIG_EXCEPTION_MESSAGE = new String(new byte[10_000]).replaceAll("\u0000", "x"); @RegisterExtension - private RabbitTestSupport rabbitTestSupport = new RabbitTestSupport(true, RABBITMQ.getAmqpPort(), RABBITMQ.getHttpPort()); + private final RabbitTestSupport rabbitTestSupport = new RabbitTestSupport(true, RABBITMQ.getAmqpPort(), RABBITMQ.getHttpPort()); private int maxStackTraceSize; @@ -602,6 +603,10 @@ public class RabbitBinderTests extends SimpleMessageListenerContainer container = TestUtils.getPropertyValue(endpoint, "messageListenerContainer", SimpleMessageListenerContainer.class); assertThat(container.isRunning()).isTrue(); + ApplicationContext ctx = + TestUtils.getPropertyValue(binder, "binder.provisioningProvider.autoDeclareContext", + ApplicationContext.class); + assertThat(ctx.getBean("infra.binding", Declarables.class).getDeclarables()).hasSize(2); consumerBinding.unbind(); assertThat(container.isRunning()).isFalse(); assertThat(container.getQueueNames()[0]).isEqualTo(group);