diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitAdmin.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitAdmin.java index 5ac532e5..e436dda6 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitAdmin.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitAdmin.java @@ -41,9 +41,12 @@ import com.rabbitmq.client.Channel; * @author Mark Pollack * @author Mark Fisher * @author Dave Syer + * @author Ed Scriven */ public class RabbitAdmin implements AmqpAdmin, ApplicationContextAware, InitializingBean { + protected static final String DEFAULT_EXCHANGE_NAME = ""; + /** Logger available to subclasses */ protected final Log logger = LogFactory.getLog(getClass()); @@ -92,6 +95,10 @@ public class RabbitAdmin implements AmqpAdmin, ApplicationContextAware, Initiali public boolean deleteExchange(final String exchangeName) { return this.rabbitTemplate.execute(new ChannelCallback() { public Boolean doInRabbit(Channel channel) throws Exception { + if (isDeletingDefaultExchange(exchangeName)) { + return true; + } + try { channel.exchangeDelete(exchangeName); } catch (IOException e) { @@ -178,6 +185,10 @@ public class RabbitAdmin implements AmqpAdmin, ApplicationContextAware, Initiali rabbitTemplate.execute(new ChannelCallback() { public Object doInRabbit(Channel channel) throws Exception { if (binding.isDestinationQueue()) { + if (isRemovingImplicitQueueBinding(binding)) { + return null; + } + channel.queueUnbind(binding.getDestination(), binding.getExchange(), binding.getRoutingKey(), binding.getArguments()); } else { @@ -314,8 +325,11 @@ public class RabbitAdmin implements AmqpAdmin, ApplicationContextAware, Initiali if (logger.isDebugEnabled()) { logger.debug("declaring Exchange '" + exchange.getName() + "'"); } - channel.exchangeDeclare(exchange.getName(), exchange.getType(), exchange.isDurable(), - exchange.isAutoDelete(), exchange.getArguments()); + + if (!isDeclaringDefaultExchange(exchange)) { + channel.exchangeDeclare(exchange.getName(), exchange.getType(), exchange.isDurable(), + exchange.isAutoDelete(), exchange.getArguments()); + } } } @@ -340,9 +354,12 @@ public class RabbitAdmin implements AmqpAdmin, ApplicationContextAware, Initiali + ")] to exchange [" + binding.getExchange() + "] with routing key [" + binding.getRoutingKey() + "]"); } + if (binding.isDestinationQueue()) { - channel.queueBind(binding.getDestination(), binding.getExchange(), binding.getRoutingKey(), - binding.getArguments()); + if (!isDeclaringImplicitQueueBinding(binding)) { + channel.queueBind(binding.getDestination(), binding.getExchange(), binding.getRoutingKey(), + binding.getArguments()); + } } else { channel.exchangeBind(binding.getDestination(), binding.getExchange(), binding.getRoutingKey(), binding.getArguments()); @@ -350,4 +367,51 @@ public class RabbitAdmin implements AmqpAdmin, ApplicationContextAware, Initiali } } + private boolean isDeclaringDefaultExchange(Exchange exchange) { + if (isDefaultExchange(exchange.getName())) { + if (logger.isDebugEnabled()) { + logger.debug("Default exchange is pre-declared by server."); + } + return true; + } + return false; + } + + private boolean isDeletingDefaultExchange(String exchangeName) { + if (isDefaultExchange(exchangeName)) { + if (logger.isDebugEnabled()) { + logger.debug("Default exchange cannot be deleted."); + } + return true; + } + return false; + } + + private boolean isDefaultExchange(String exchangeName) { + return DEFAULT_EXCHANGE_NAME.equals(exchangeName); + } + + private boolean isDeclaringImplicitQueueBinding(Binding binding) { + if (isImplicitQueueBinding(binding)) { + if (logger.isDebugEnabled()) { + logger.debug("The default exchange is implicitly bound to every queue, with a routing key equal to the queue name."); + } + return true; + } + return false; + } + + private boolean isRemovingImplicitQueueBinding(Binding binding) { + if (isImplicitQueueBinding(binding)) { + if (logger.isDebugEnabled()) { + logger.debug("Cannot remove implicit default exchange binding to queue."); + } + return true; + } + return false; + } + + private boolean isImplicitQueueBinding(Binding binding) { + return isDefaultExchange(binding.getExchange()) && binding.getDestination().equals(binding.getRoutingKey()); + } } diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitAdminIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitAdminIntegrationTests.java index 2ba1ce08..be158653 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitAdminIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/RabbitAdminIntegrationTests.java @@ -4,13 +4,16 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; import java.io.IOException; -import java.util.concurrent.atomic.AtomicReference; import org.junit.After; import org.junit.Before; import org.junit.Rule; import org.junit.Test; import org.springframework.amqp.AmqpIOException; +import org.springframework.amqp.core.Binding; +import org.springframework.amqp.core.Binding.DestinationType; +import org.springframework.amqp.core.DirectExchange; +import org.springframework.amqp.core.Exchange; import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; import org.springframework.amqp.rabbit.test.BrokerRunning; @@ -101,114 +104,183 @@ public class RabbitAdminIntegrationTests { } @Test - public void testStartupWithAutodelete() throws Exception { + public void testQueueWithAutoDelete() throws Exception { final Queue queue = new Queue("test.queue", false, true, true); context.getBeanFactory().registerSingleton("foo", queue); rabbitAdmin.afterPropertiesSet(); - final AtomicReference connectionHolder = new AtomicReference(); - - RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); - boolean exists = rabbitTemplate.execute(new ChannelCallback() { - public Boolean doInRabbit(Channel channel) throws Exception { - DeclareOk result = channel.queueDeclarePassive(queue.getName()); - connectionHolder.set(channel.getConnection()); - return result != null; - } - }); - assertTrue("Expected Queue to exist", exists); - - assertTrue(queueExists(connectionHolder.get(), queue)); - - exists = rabbitTemplate.execute(new ChannelCallback() { - public Boolean doInRabbit(Channel channel) throws Exception { - DeclareOk result = channel.queueDeclarePassive(queue.getName()); - connectionHolder.set(channel.getConnection()); - return result != null; - } - }); - assertTrue("Expected Queue to exist", exists); + // Queue created on spring startup + rabbitAdmin.initialize(); + assertTrue(queueExists(queue)); + // Stop and broker deletes queue (only verifiable in native API) connectionFactory.destroy(); - // Broker now deletes queue (only verifiable in native API) - assertFalse(queueExists(null, queue)); + assertFalse(queueExists(queue)); - // Broker auto-deleted queue, but it is re-created by the connection listener - exists = rabbitTemplate.execute(new ChannelCallback() { - public Boolean doInRabbit(Channel channel) throws Exception { - DeclareOk result = channel.queueDeclarePassive(queue.getName()); - connectionHolder.set(channel.getConnection()); - return result != null; - } - }); - assertTrue("Expected Queue to exist", exists); + // Start and queue re-created by the connection listener + connectionFactory.createConnection(); + assertTrue(queueExists(queue)); - assertTrue(queueExists(connectionHolder.get(), queue)); + // Queue manually deleted assertTrue(rabbitAdmin.deleteQueue(queue.getName())); - assertFalse(queueExists(null, queue)); + assertFalse(queueExists(queue)); } @Test - public void testStartupWithNonDurable() throws Exception { + public void testQueueWithoutAutoDelete() throws Exception { final Queue queue = new Queue("test.queue", false, false, false); context.getBeanFactory().registerSingleton("foo", queue); rabbitAdmin.afterPropertiesSet(); - final AtomicReference connectionHolder = new AtomicReference(); + // Queue created on Spring startup + rabbitAdmin.initialize(); + assertTrue(queueExists(queue)); - RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); - // Force RabbitAdmin to initialize the queue - boolean exists = rabbitTemplate.execute(new ChannelCallback() { - public Boolean doInRabbit(Channel channel) throws Exception { - DeclareOk result = channel.queueDeclarePassive(queue.getName()); - connectionHolder.set(channel.getConnection()); - return result != null; - } - }); - assertTrue("Expected Queue to exist", exists); - - assertTrue(queueExists(connectionHolder.get(), queue)); - - // simulate broker going down and coming back up... - rabbitAdmin.deleteQueue(queue.getName()); + // Stop and broker retains queue (only verifiable in native API) connectionFactory.destroy(); - assertFalse(queueExists(null, queue)); + assertTrue(queueExists(queue)); - // Broker auto-deleted queue, but it is re-created by the connection listener - exists = rabbitTemplate.execute(new ChannelCallback() { - public Boolean doInRabbit(Channel channel) throws Exception { - DeclareOk result = channel.queueDeclarePassive(queue.getName()); - connectionHolder.set(channel.getConnection()); - return result != null; - } - }); - assertTrue("Expected Queue to exist", exists); + // Start and queue still exists + connectionFactory.createConnection(); + assertTrue(queueExists(queue)); - assertTrue(queueExists(connectionHolder.get(), queue)); + // Queue manually deleted assertTrue(rabbitAdmin.deleteQueue(queue.getName())); - assertFalse(queueExists(null, queue)); + assertFalse(queueExists(queue)); + } + @Test + public void testDeclareExchangeWithDefaultExchange() throws Exception { + Exchange exchange = new DirectExchange(RabbitAdmin.DEFAULT_EXCHANGE_NAME); + + rabbitAdmin.declareExchange(exchange); + + // Pass by virtue of RabbitMQ not firing a 403 reply code + } + + @Test + public void testSpringWithDefaultExchange() throws Exception { + Exchange exchange = new DirectExchange(RabbitAdmin.DEFAULT_EXCHANGE_NAME); + context.getBeanFactory().registerSingleton("foo", exchange); + rabbitAdmin.afterPropertiesSet(); + + rabbitAdmin.initialize(); + + // Pass by virtue of RabbitMQ not firing a 403 reply code + } + + @Test + public void testDeleteExchangeWithDefaultExchange() throws Exception { + boolean result = rabbitAdmin.deleteExchange(RabbitAdmin.DEFAULT_EXCHANGE_NAME); + + assertTrue(result); + } + + @Test + public void testDeclareBindingWithDefaultExchangeImplicitBinding() throws Exception { + Exchange exchange = new DirectExchange(RabbitAdmin.DEFAULT_EXCHANGE_NAME); + String queueName = "test.queue"; + final Queue queue = new Queue(queueName, false, false, false); + rabbitAdmin.declareQueue(queue); + Binding binding = new Binding(queueName, DestinationType.QUEUE, exchange.getName(), queueName, null); + + rabbitAdmin.declareBinding(binding); + + // Pass by virtue of RabbitMQ not firing a 403 reply code for both exchange and binding declaration + assertTrue(queueExists(queue)); + } + + @Test + public void testSpringWithDefaultExchangeImplicitBinding() throws Exception { + Exchange exchange = new DirectExchange(RabbitAdmin.DEFAULT_EXCHANGE_NAME); + context.getBeanFactory().registerSingleton("foo", exchange); + String queueName = "test.queue"; + final Queue queue = new Queue(queueName, false, false, false); + context.getBeanFactory().registerSingleton("bar", queue); + Binding binding = new Binding(queueName, DestinationType.QUEUE, exchange.getName(), queueName, null); + context.getBeanFactory().registerSingleton("baz", binding); + rabbitAdmin.afterPropertiesSet(); + + rabbitAdmin.initialize(); + + // Pass by virtue of RabbitMQ not firing a 403 reply code for both exchange and binding declaration + assertTrue(queueExists(queue)); + } + + @Test + public void testRemoveBindingWithDefaultExchangeImplicitBinding() throws Exception { + String queueName = "test.queue"; + final Queue queue = new Queue(queueName, false, false, false); + rabbitAdmin.declareQueue(queue); + Binding binding = new Binding(queueName, DestinationType.QUEUE, RabbitAdmin.DEFAULT_EXCHANGE_NAME, queueName, null); + + rabbitAdmin.removeBinding(binding); + + // Pass by virtue of RabbitMQ not firing a 403 reply code + } + + @Test + public void testDeclareBindingWithDefaultExchangeNonImplicitBinding() throws Exception { + Exchange exchange = new DirectExchange(RabbitAdmin.DEFAULT_EXCHANGE_NAME); + String queueName = "test.queue"; + final Queue queue = new Queue(queueName, false, false, false); + rabbitAdmin.declareQueue(queue); + Binding binding = new Binding(queueName, DestinationType.QUEUE, exchange.getName(), "test.routingKey", null); + + try { + rabbitAdmin.declareBinding(binding); + } catch (AmqpIOException ex) { + Throwable cause = ex; + Throwable rootCause = null; + while (cause != null) { + rootCause = cause; + cause = cause.getCause(); + } + assertTrue(rootCause.getMessage().contains("reply-code=403")); + assertTrue(rootCause.getMessage().contains("operation not permitted on the default exchange")); + } + } + + @Test + public void testSpringWithDefaultExchangeNonImplicitBinding() throws Exception { + Exchange exchange = new DirectExchange(RabbitAdmin.DEFAULT_EXCHANGE_NAME); + context.getBeanFactory().registerSingleton("foo", exchange); + String queueName = "test.queue"; + final Queue queue = new Queue(queueName, false, false, false); + context.getBeanFactory().registerSingleton("bar", queue); + Binding binding = new Binding(queueName, DestinationType.QUEUE, exchange.getName(), "test.routingKey", null); + context.getBeanFactory().registerSingleton("baz", binding); + rabbitAdmin.afterPropertiesSet(); + + try { + rabbitAdmin.declareBinding(binding); + } catch (AmqpIOException ex) { + Throwable cause = ex; + Throwable rootCause = null; + while (cause != null) { + rootCause = cause; + cause = cause.getCause(); + } + assertTrue(rootCause.getMessage().contains("reply-code=403")); + assertTrue(rootCause.getMessage().contains("operation not permitted on the default exchange")); + } } /** - * Use native Rabbit API to test queue, bypassing all the connection and channel caching and callbacks in Spring - * AMQP. + * Verify that a queue exists using the native Rabbit API to bypass all the connection and + * channel caching and callbacks in Spring AMQP. * - * @param connection the raw connection to use - * @param queue the Queue to test - * @return true if the queue exists + * @param Queue The queue to verify + * @return True if the queue exists */ - private boolean queueExists(Connection connection, Queue queue) throws Exception { - Connection target = connection; - if (target == null) { - ConnectionFactory connectionFactory = new ConnectionFactory(); - connectionFactory.setPort(BrokerTestUtils.getPort()); - target = connectionFactory.newConnection(); - } - Channel channel = target.createChannel(); + private boolean queueExists(final Queue queue) throws Exception { + ConnectionFactory connectionFactory = new ConnectionFactory(); + connectionFactory.setPort(BrokerTestUtils.getPort()); + Connection connection = connectionFactory.newConnection(); + Channel channel = connection.createChannel(); try { DeclareOk result = channel.queueDeclarePassive(queue.getName()); return result != null; @@ -218,10 +290,7 @@ public class RabbitAdminIntegrationTests { } return false; } finally { - if (connection==null) { - target.close(); - } + connection.close(); } } - }