From a5d8e36a4937d8146cfb08a9ff2afe4f211b296f Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 22 Mar 2023 13:16:21 -0400 Subject: [PATCH] GH-2432: Fix Redeclaration of Declarables Resolves https://github.com/spring-projects/spring-amqp/issues/2432 When a container starts; it looks to see if the context contains a queue bean for any of its queues and requests the admin to redeclare the infrastructure. However, if the queue is declared within a `Declarables`, the check is not performed and therefore the queue is not declared. Add a check for `Declarables` beans. **cherry-pick to 2.4.x** --- .../AbstractMessageListenerContainer.java | 13 +- .../listener/QueueDeclarationTests.java | 124 ++++++++++++++++++ 2 files changed, 132 insertions(+), 5 deletions(-) create mode 100644 spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/QueueDeclarationTests.java diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java index c75c8074..0426a6d3 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2022 the original author or authors. + * Copyright 2002-2023 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. @@ -21,9 +21,9 @@ import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.HashSet; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; -import java.util.Map.Entry; import java.util.Properties; import java.util.Set; import java.util.concurrent.CopyOnWriteArrayList; @@ -42,6 +42,7 @@ import org.springframework.amqp.ImmediateAcknowledgeAmqpException; import org.springframework.amqp.core.AcknowledgeMode; import org.springframework.amqp.core.AmqpAdmin; import org.springframework.amqp.core.BatchMessageListener; +import org.springframework.amqp.core.Declarables; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; import org.springframework.amqp.core.MessagePostProcessor; @@ -1977,9 +1978,11 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor ApplicationContext context = this.getApplicationContext(); if (context != null) { Set queueNames = getQueueNamesAsSet(); - Map queueBeans = context.getBeansOfType(Queue.class); - for (Entry entry : queueBeans.entrySet()) { - Queue queue = entry.getValue(); + Collection queueBeans = new LinkedHashSet<>( + context.getBeansOfType(Queue.class, false, false).values()); + Map declarables = context.getBeansOfType(Declarables.class, false, false); + declarables.values().forEach(dec -> queueBeans.addAll(dec.getDeclarablesByType(Queue.class))); + for (Queue queue : queueBeans) { if (isMismatchedQueuesFatal() || (queueNames.contains(queue.getName()) && admin.getQueueProperties(queue.getName()) == null)) { if (logger.isDebugEnabled()) { diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/QueueDeclarationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/QueueDeclarationTests.java new file mode 100644 index 00000000..2c99f8a8 --- /dev/null +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/QueueDeclarationTests.java @@ -0,0 +1,124 @@ +/* + * Copyright 2023 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.listener; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyMap; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; + +import java.io.IOException; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import org.springframework.amqp.core.AmqpAdmin; +import org.springframework.amqp.core.Declarables; +import org.springframework.amqp.core.Queue; +import org.springframework.amqp.rabbit.connection.ChannelProxy; +import org.springframework.amqp.rabbit.connection.Connection; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; +import org.springframework.context.ApplicationContext; + +import com.rabbitmq.client.AMQP; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.Consumer; +import com.rabbitmq.client.impl.recovery.AutorecoveringChannel; + +/** + * @author Gary Russell + * @since 2.4.12 + * + */ +public class QueueDeclarationTests { + + @Test + void redeclareWhenQueue() throws IOException, InterruptedException { + AmqpAdmin admin = mock(AmqpAdmin.class); + ApplicationContext context = mock(ApplicationContext.class); + final CountDownLatch latch = new CountDownLatch(1); + SimpleMessageListenerContainer container = createContainer(admin, latch); + given(context.getBeansOfType(Queue.class, false, false)).willReturn(Map.of("foo", new Queue("test"))); + given(context.getBeansOfType(Declarables.class, false, false)).willReturn(new HashMap<>()); + container.setApplicationContext(context); + container.start(); + + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + verify(admin).initialize(); + container.stop(); + } + + @Test + void redeclareWhenDeclarables() throws IOException, InterruptedException { + AmqpAdmin admin = mock(AmqpAdmin.class); + ApplicationContext context = mock(ApplicationContext.class); + final CountDownLatch latch = new CountDownLatch(1); + SimpleMessageListenerContainer container = createContainer(admin, latch); + given(context.getBeansOfType(Queue.class, false, false)).willReturn(new HashMap<>()); + given(context.getBeansOfType(Declarables.class, false, false)) + .willReturn(Map.of("foo", new Declarables(List.of(new Queue("test"))))); + container.setApplicationContext(context); + container.start(); + + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + verify(admin).initialize(); + container.stop(); + } + + private SimpleMessageListenerContainer createContainer(AmqpAdmin admin, final CountDownLatch latch) + throws IOException { + ConnectionFactory connectionFactory = mock(ConnectionFactory.class); + Connection connection = mock(Connection.class); + ChannelProxy channel = mock(ChannelProxy.class); + Channel rabbitChannel = mock(AutorecoveringChannel.class); + given(channel.getTargetChannel()).willReturn(rabbitChannel); + + given(connectionFactory.createConnection()).willReturn(connection); + given(connection.createChannel(anyBoolean())).willReturn(channel); + final AtomicBoolean isOpen = new AtomicBoolean(true); + willAnswer(i -> isOpen.get()).given(channel).isOpen(); + given(channel.queueDeclarePassive(Mockito.anyString())) + .willAnswer(invocation -> mock(AMQP.Queue.DeclareOk.class)); + given(channel.basicConsume(anyString(), anyBoolean(), anyString(), anyBoolean(), anyBoolean(), + anyMap(), any(Consumer.class))).willReturn("consumerTag"); + + willAnswer(i -> { + latch.countDown(); + return null; + }).given(channel).basicQos(anyInt(), anyBoolean()); + + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory); + container.setQueueNames("test"); + container.setPrefetchCount(2); + container.setAmqpAdmin(admin); + container.afterPropertiesSet(); + return container; + } + +}