diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/MissingMessageIdAdvice.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/MissingMessageIdAdvice.java new file mode 100644 index 00000000..0668cf1e --- /dev/null +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/MissingMessageIdAdvice.java @@ -0,0 +1,86 @@ +/* + * Copyright 2002-2012 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 + * + * http://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.retry; + +import java.util.UUID; + +import org.aopalliance.intercept.MethodInterceptor; +import org.aopalliance.intercept.MethodInvocation; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.amqp.AmqpRejectAndDontRequeueException; +import org.springframework.amqp.core.Message; +import org.springframework.amqp.core.MessageProperties; +import org.springframework.amqp.rabbit.listener.ListenerExecutionFailedException; +import org.springframework.retry.policy.RetryContextCache; +import org.springframework.util.Assert; + +/** + * Advice that can be placed in the listener delegate's advice chain to + * enhance the message with an ID if not present. + * If an exception is caught on a redelivered message, rethrows it as an {@link AmqpRejectAndDontRequeueException} + * which signals the container to NOT requeue the message (otherwise we'd have infinite + * immediate retries). + * If so configured, the broker can send the message to a DLE/DLQ. + * Must be placed before the retry interceptor in the advice chain. + * @author Gary Russell + * @since 1.1.2 + * + */ +public class MissingMessageIdAdvice implements MethodInterceptor { + + private final static Log logger = LogFactory.getLog(MissingMessageIdAdvice.class); + + private final RetryContextCache retryContextCache; + + public MissingMessageIdAdvice(RetryContextCache retryContextCache) { + Assert.notNull(retryContextCache, "RetryContextCache must not be null"); + this.retryContextCache = retryContextCache; + } + + public Object invoke(MethodInvocation invocation) throws Throwable { + String id = null; + boolean redelivered = false; + try { + Message message = (Message) invocation.getArguments()[1]; + MessageProperties messageProperties = message.getMessageProperties(); + if (messageProperties.getMessageId() == null) { + id = UUID.randomUUID().toString(); + messageProperties.setMessageId(id); + } + redelivered = messageProperties.isRedelivered(); + return invocation.proceed(); + } + catch (Throwable t) { + if (id != null && redelivered) { + if (logger.isDebugEnabled()) { + logger.debug("Canceling delivery of retried message that has no ID"); + } + throw new ListenerExecutionFailedException("Cannot retry message without an ID", + new AmqpRejectAndDontRequeueException(t)); + } + else { + throw t; + } + } + finally { + if (id != null) { + retryContextCache.remove(id); + } + } + } + +} diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/RejectAndDontRequeueRecoverer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/RejectAndDontRequeueRecoverer.java new file mode 100644 index 00000000..134fa264 --- /dev/null +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/retry/RejectAndDontRequeueRecoverer.java @@ -0,0 +1,39 @@ +/* + * Copyright 2002-2012 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 + * + * http://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.retry; + +import org.springframework.amqp.AmqpRejectAndDontRequeueException; +import org.springframework.amqp.core.Message; +import org.springframework.amqp.rabbit.listener.ListenerExecutionFailedException; + +/** + * MessageRecover that causes the listener container to reject + * the message without requeuing. This enables failed messages + * to be sent to a Dead Letter Exchange/Queue, if the broker is + * so configured. + * + * @author Gary Russell + * @since 1.1.2 + * + */ +public class RejectAndDontRequeueRecoverer implements MessageRecoverer { + + public void recover(Message message, Throwable cause) { + throw new ListenerExecutionFailedException("Retry Policy Exhausted", + new AmqpRejectAndDontRequeueException(cause)); + } + +} diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/retry/MissingIdRetryTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/retry/MissingIdRetryTests.java new file mode 100644 index 00000000..f9af8feb --- /dev/null +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/retry/MissingIdRetryTests.java @@ -0,0 +1,134 @@ +/* + * Copyright 2002-2012 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 + * + * http://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.retry; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import org.aopalliance.aop.Advice; +import org.junit.Test; +import org.springframework.amqp.core.Message; +import org.springframework.amqp.core.MessageProperties; +import org.springframework.amqp.rabbit.config.StatefulRetryOperationsInterceptorFactoryBean; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; +import org.springframework.amqp.rabbit.listener.adapter.MessageListenerAdapter; +import org.springframework.beans.DirectFieldAccessor; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.retry.policy.MapRetryContextCache; +import org.springframework.retry.policy.RetryContextCache; +import org.springframework.retry.support.RetryTemplate; + +/** + * @author Gary Russell + * @since 1.1.2 + * + */ +public class MissingIdRetryTests { + + private volatile CountDownLatch latch; + + @SuppressWarnings("rawtypes") + @Test + public void testWithNoId() throws Exception { + // 2 messsages; each retried once by missing id interceptor + this.latch = new CountDownLatch(4); + ApplicationContext ctx = new ClassPathXmlApplicationContext("retry-context.xml", this.getClass()); + RabbitTemplate template = ctx.getBean(RabbitTemplate.class); + ConnectionFactory connectionFactory = ctx.getBean(ConnectionFactory.class); + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory); + container.setMessageListener(new MessageListenerAdapter(new POJO())); + container.setQueueNames("retry.test.queue"); + + StatefulRetryOperationsInterceptorFactoryBean fb = new StatefulRetryOperationsInterceptorFactoryBean(); + + // use an external template so we can share his cache + RetryTemplate retryTemplate = new RetryTemplate(); + RetryContextCache cache = new MapRetryContextCache(); + retryTemplate.setRetryContextCache(cache); + fb.setRetryOperations(retryTemplate); + + // give him a reference to the retry cache so he can clean it up + MissingMessageIdAdvice missingIdAdvice = new MissingMessageIdAdvice(cache); + + Advice retryInterceptor = fb.getObject(); + // add both advices + container.setAdviceChain(new Advice[] {missingIdAdvice, retryInterceptor}); + container.start(); + + template.convertAndSend("retry.test.exchange", "retry.test.binding", "Hello, world!"); + template.convertAndSend("retry.test.exchange", "retry.test.binding", "Hello, world!"); + assertTrue(latch.await(10, TimeUnit.SECONDS)); + Thread.sleep(2000); + assertEquals(0, ((Map) new DirectFieldAccessor(cache).getPropertyValue("map")).size()); + container.stop(); + } + + @SuppressWarnings("rawtypes") + @Test + public void testWithId() throws Exception { + // 2 messsages; each retried twice by retry interceptor + this.latch = new CountDownLatch(6); + ApplicationContext ctx = new ClassPathXmlApplicationContext("retry-context.xml", this.getClass()); + RabbitTemplate template = ctx.getBean(RabbitTemplate.class); + ConnectionFactory connectionFactory = ctx.getBean(ConnectionFactory.class); + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory); + container.setMessageListener(new MessageListenerAdapter(new POJO())); + container.setQueueNames("retry.test.queue"); + + StatefulRetryOperationsInterceptorFactoryBean fb = new StatefulRetryOperationsInterceptorFactoryBean(); + + // use an external template so we can share his cache + RetryTemplate retryTemplate = new RetryTemplate(); + RetryContextCache cache = new MapRetryContextCache(); + retryTemplate.setRetryContextCache(cache); + fb.setRetryOperations(retryTemplate); + fb.setMessageRecoverer(new RejectAndDontRequeueRecoverer()); + + // give him a reference to the retry cache so he can clean it up + MissingMessageIdAdvice missingIdAdvice = new MissingMessageIdAdvice(cache); + + Advice retryInterceptor = fb.getObject(); + // add both advices + container.setAdviceChain(new Advice[] {missingIdAdvice, retryInterceptor}); + container.start(); + + MessageProperties messageProperties = new MessageProperties(); + messageProperties.setContentType("text/plain"); + messageProperties.setMessageId("foo"); + Message message = new Message("Hello, world!".getBytes(), messageProperties); + template.send("retry.test.exchange", "retry.test.binding", message); + template.send("retry.test.exchange", "retry.test.binding", message); + assertTrue(latch.await(10, TimeUnit.SECONDS)); + Thread.sleep(2000); + assertEquals(0, ((Map) new DirectFieldAccessor(cache).getPropertyValue("map")).size()); + container.stop(); + } + + public class POJO { + public void handleMessage(String foo) { +// System.out.println(foo); + latch.countDown(); + throw new RuntimeException("fail"); + } + } +} diff --git a/spring-rabbit/src/test/resources/org/springframework/amqp/rabbit/retry/retry-context.xml b/spring-rabbit/src/test/resources/org/springframework/amqp/rabbit/retry/retry-context.xml new file mode 100644 index 00000000..5d5d8dd2 --- /dev/null +++ b/spring-rabbit/src/test/resources/org/springframework/amqp/rabbit/retry/retry-context.xml @@ -0,0 +1,22 @@ + + + + + + + + + + + + + + + + + +