diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java index c580b99b6f..c55143d512 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 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. @@ -19,8 +19,10 @@ package org.springframework.integration.amqp.inbound; import java.util.Map; import org.springframework.amqp.core.AcknowledgeMode; +import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.core.ChannelAwareMessageListener; import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; +import org.springframework.amqp.rabbit.listener.exception.ListenerExecutionFailedException; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.amqp.support.converter.SimpleMessageConverter; @@ -28,8 +30,11 @@ import org.springframework.integration.amqp.support.AmqpHeaderMapper; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; import org.springframework.integration.context.OrderlyShutdownCapable; import org.springframework.integration.endpoint.MessageProducerSupport; +import org.springframework.messaging.support.ErrorMessage; import org.springframework.util.Assert; +import com.rabbitmq.client.Channel; + /** * Adapter that receives Messages from an AMQP Queue, converts them into * Spring Integration Messages, and sends the results to a Message Channel. @@ -37,6 +42,7 @@ import org.springframework.util.Assert; * @author Mark Fisher * @author Gary Russell * @author Artem Bilan + * * @since 2.1 */ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements @@ -77,17 +83,7 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements @Override protected void onInit() { - this.messageListenerContainer.setMessageListener((ChannelAwareMessageListener) (message, channel) -> { - Object payload = this.messageConverter.fromMessage(message); - Map headers = - this.headerMapper.toHeadersFromRequest(message.getMessageProperties()); - if (this.messageListenerContainer.getAcknowledgeMode() - == AcknowledgeMode.MANUAL) { - headers.put(AmqpHeaders.DELIVERY_TAG, message.getMessageProperties().getDeliveryTag()); - headers.put(AmqpHeaders.CHANNEL, channel); - } - sendMessage(getMessageBuilderFactory().withPayload(payload).copyHeaders(headers).build()); - }); + this.messageListenerContainer.setMessageListener(new Listener()); this.messageListenerContainer.afterPropertiesSet(); super.onInit(); } @@ -115,13 +111,41 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements } - /** - * {@inheritDoc} - *

No-op - */ @Override public int afterShutdown() { return 0; } + protected class Listener implements ChannelAwareMessageListener { + + @Override + public void onMessage(Message message, Channel channel) throws Exception { + try { + Object payload = AmqpInboundChannelAdapter.this.messageConverter.fromMessage(message); + Map headers = AmqpInboundChannelAdapter.this.headerMapper + .toHeadersFromRequest(message.getMessageProperties()); + if (AmqpInboundChannelAdapter.this.messageListenerContainer.getAcknowledgeMode() + == AcknowledgeMode.MANUAL) { + headers.put(AmqpHeaders.DELIVERY_TAG, message.getMessageProperties().getDeliveryTag()); + headers.put(AmqpHeaders.CHANNEL, channel); + } + sendMessage( + getMessageBuilderFactory() + .withPayload(payload) + .copyHeaders(headers) + .build()); + } + catch (RuntimeException e) { + if (getErrorChannel() != null) { + getMessagingTemplate().send(getErrorChannel(), new ErrorMessage( + new ListenerExecutionFailedException("Message conversion failed", e, message))); + } + else { + throw e; + } + } + } + + } + } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java index 043bf1c615..5f0d44e6e2 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 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. @@ -27,15 +27,19 @@ import org.springframework.amqp.core.MessageProperties; import org.springframework.amqp.rabbit.core.ChannelAwareMessageListener; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; +import org.springframework.amqp.rabbit.listener.exception.ListenerExecutionFailedException; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.amqp.support.converter.SimpleMessageConverter; import org.springframework.integration.amqp.support.AmqpHeaderMapper; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; import org.springframework.integration.gateway.MessagingGatewaySupport; +import org.springframework.messaging.support.ErrorMessage; import org.springframework.util.Assert; import org.springframework.util.StringUtils; +import com.rabbitmq.client.Channel; + /** * Adapter that receives Messages from an AMQP Queue, converts them into * Spring Integration Messages, and sends the results to a Message Channel. @@ -44,6 +48,8 @@ import org.springframework.util.StringUtils; * * @author Mark Fisher * @author Artem Bilan + * @author Gary Russell + * * @since 2.1 */ public class AmqpInboundGateway extends MessagingGatewaySupport { @@ -136,63 +142,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { @Override protected void onInit() throws Exception { - this.messageListenerContainer.setMessageListener((ChannelAwareMessageListener) (message, channel) -> { - Object payload = this.amqpMessageConverter.fromMessage(message); - Map headers = - this.headerMapper.toHeadersFromRequest(message.getMessageProperties()); - if (this.messageListenerContainer.getAcknowledgeMode() == AcknowledgeMode.MANUAL) { - headers.put(AmqpHeaders.DELIVERY_TAG, message.getMessageProperties().getDeliveryTag()); - headers.put(AmqpHeaders.CHANNEL, channel); - } - org.springframework.messaging.Message request = getMessageBuilderFactory() - .withPayload(payload) - .copyHeaders(headers) - .build(); - final org.springframework.messaging.Message reply = sendAndReceiveMessage(request); - if (reply != null) { - Address replyTo; - String replyToProperty = message.getMessageProperties().getReplyTo(); - if (replyToProperty != null) { - replyTo = new Address(replyToProperty); - } - else { - replyTo = AmqpInboundGateway.this.defaultReplyTo; - } - - MessagePostProcessor messagePostProcessor = message1 -> { - MessageProperties messageProperties = message1.getMessageProperties(); - String contentEncoding = messageProperties.getContentEncoding(); - long contentLength = messageProperties.getContentLength(); - String contentType = messageProperties.getContentType(); - this.headerMapper.fromHeadersToReply(reply.getHeaders(), messageProperties); - // clear the replyTo from the original message since we are using it now - messageProperties.setReplyTo(null); - // reset the content-* properties as determined by the MessageConverter - if (StringUtils.hasText(contentEncoding)) { - messageProperties.setContentEncoding(contentEncoding); - } - messageProperties.setContentLength(contentLength); - if (contentType != null) { - messageProperties.setContentType(contentType); - } - return message1; - }; - - if (replyTo != null) { - this.amqpTemplate.convertAndSend(replyTo.getExchangeName(), - replyTo.getRoutingKey(), reply.getPayload(), messagePostProcessor); - } - else { - if (!this.amqpTemplateExplicitlySet) { - throw new IllegalStateException("There is no 'replyTo' message property " + - "and the `defaultReplyTo` hasn't been configured."); - } - else { - this.amqpTemplate.convertAndSend(reply.getPayload(), messagePostProcessor); - } - } - } - }); + this.messageListenerContainer.setMessageListener(new Listener()); this.messageListenerContainer.afterPropertiesSet(); if (!this.amqpTemplateExplicitlySet) { ((RabbitTemplate) this.amqpTemplate).afterPropertiesSet(); @@ -210,4 +160,86 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { this.messageListenerContainer.stop(); } + protected class Listener implements ChannelAwareMessageListener { + + @Override + public void onMessage(Message message, Channel channel) throws Exception { + boolean error = false; + Map headers = null; + Object payload = null; + try { + payload = AmqpInboundGateway.this.amqpMessageConverter.fromMessage(message); + headers = AmqpInboundGateway.this.headerMapper.toHeadersFromRequest(message.getMessageProperties()); + if (AmqpInboundGateway.this.messageListenerContainer.getAcknowledgeMode() == AcknowledgeMode.MANUAL) { + headers.put(AmqpHeaders.DELIVERY_TAG, message.getMessageProperties().getDeliveryTag()); + headers.put(AmqpHeaders.CHANNEL, channel); + } + } + catch (RuntimeException e) { + if (getErrorChannel() != null) { + AmqpInboundGateway.this.messagingTemplate.send(getErrorChannel(), new ErrorMessage( + new ListenerExecutionFailedException("Message conversion failed", e, message))); + } + else { + throw e; + } + error = true; + } + + if (!error) { + final org.springframework.messaging.Message reply = + sendAndReceiveMessage( + getMessageBuilderFactory() + .withPayload(payload) + .copyHeaders(headers) + .build()); + if (reply != null) { + Address replyTo; + String replyToProperty = message.getMessageProperties().getReplyTo(); + if (replyToProperty != null) { + replyTo = new Address(replyToProperty); + } + else { + replyTo = AmqpInboundGateway.this.defaultReplyTo; + } + + MessagePostProcessor messagePostProcessor = message1 -> { + MessageProperties messageProperties = message1.getMessageProperties(); + String contentEncoding = messageProperties.getContentEncoding(); + long contentLength = messageProperties.getContentLength(); + String contentType = messageProperties.getContentType(); + AmqpInboundGateway.this.headerMapper.fromHeadersToReply(reply.getHeaders(), messageProperties); + // clear the replyTo from the original message since we are using it now + messageProperties.setReplyTo(null); + // reset the content-* properties as determined by the MessageConverter + if (StringUtils.hasText(contentEncoding)) { + messageProperties.setContentEncoding(contentEncoding); + } + messageProperties.setContentLength(contentLength); + if (contentType != null) { + messageProperties.setContentType(contentType); + } + return message1; + }; + + if (replyTo != null) { + AmqpInboundGateway.this.amqpTemplate.convertAndSend(replyTo.getExchangeName(), + replyTo.getRoutingKey(), reply.getPayload(), messagePostProcessor); + } + else { + if (!AmqpInboundGateway.this.amqpTemplateExplicitlySet) { + throw new IllegalStateException("There is no 'replyTo' message property " + + "and the `defaultReplyTo` hasn't been configured."); + } + else { + AmqpInboundGateway.this.amqpTemplate.convertAndSend(reply.getPayload(), + messagePostProcessor); + } + } + } + } + } + + } + } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java index 276f8de59e..84b6366b07 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java @@ -19,6 +19,8 @@ package org.springframework.integration.amqp.inbound; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; @@ -46,6 +48,8 @@ import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; +import org.springframework.amqp.support.converter.MessageConversionException; +import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.amqp.support.converter.SimpleMessageConverter; import org.springframework.beans.factory.BeanFactory; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; @@ -210,6 +214,71 @@ public class InboundEndpointTests { assertTrue(sendLatch.await(10, TimeUnit.SECONDS)); } + @Test + public void testAdapterConversionError() throws Exception { + Connection connection = mock(Connection.class); + doAnswer(invocation -> mock(Channel.class)).when(connection).createChannel(anyBoolean()); + ConnectionFactory connectionFactory = mock(ConnectionFactory.class); + when(connectionFactory.createConnection()).thenReturn(connection); + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(); + container.setConnectionFactory(connectionFactory); + AmqpInboundChannelAdapter adapter = new AmqpInboundChannelAdapter(container); + QueueChannel outputChannel = new QueueChannel(); + adapter.setOutputChannel(outputChannel); + QueueChannel errorChannel = new QueueChannel(); + adapter.setErrorChannel(errorChannel); + adapter.setMessageConverter(new MessageConverter() { + + @Override + public org.springframework.amqp.core.Message toMessage(Object object, MessageProperties messageProperties) + throws MessageConversionException { + throw new MessageConversionException("intended"); + } + + @Override + public Object fromMessage(org.springframework.amqp.core.Message message) throws MessageConversionException { + return null; + } + + }); + adapter.afterPropertiesSet(); + ((ChannelAwareMessageListener) container.getMessageListener()).onMessage(null, null); + assertNull(outputChannel.receive(0)); + assertNotNull(errorChannel.receive(0)); + } + + @Test + public void testGatewayConversionError() throws Exception { + Connection connection = mock(Connection.class); + doAnswer(invocation -> mock(Channel.class)).when(connection).createChannel(anyBoolean()); + ConnectionFactory connectionFactory = mock(ConnectionFactory.class); + when(connectionFactory.createConnection()).thenReturn(connection); + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(); + container.setConnectionFactory(connectionFactory); + AmqpInboundGateway adapter = new AmqpInboundGateway(container); + QueueChannel outputChannel = new QueueChannel(); + adapter.setRequestChannel(outputChannel); + QueueChannel errorChannel = new QueueChannel(); + adapter.setErrorChannel(errorChannel); + adapter.setMessageConverter(new MessageConverter() { + + @Override + public org.springframework.amqp.core.Message toMessage(Object object, MessageProperties messageProperties) + throws MessageConversionException { + throw new MessageConversionException("intended"); + } + + @Override + public Object fromMessage(org.springframework.amqp.core.Message message) throws MessageConversionException { + return null; + } + + }); + adapter.afterPropertiesSet(); + ((ChannelAwareMessageListener) container.getMessageListener()).onMessage(null, null); + assertNull(outputChannel.receive(0)); + assertNotNull(errorChannel.receive(0)); + } public static class Foo { diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java index 18da883d6a..194eacb5a6 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java @@ -32,6 +32,7 @@ import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.InitializingBean; +import org.springframework.integration.core.MessagingTemplate; import org.springframework.integration.gateway.MessagingGatewaySupport; import org.springframework.integration.support.DefaultMessageBuilderFactory; import org.springframework.integration.support.MessageBuilderFactory; @@ -330,8 +331,9 @@ public class ChannelPublishingJmsMessageListener if (errorChannel == null) { throw e; } - errorChannel.send(this.gatewayDelegate.buildErrorMessage( - new MessagingException("Inbound conversion failed for: " + jmsMessage, e))); + this.gatewayDelegate.getMessagingTemplate().send(errorChannel, + this.gatewayDelegate.buildErrorMessage( + new MessagingException("Inbound conversion failed for: " + jmsMessage, e))); errors = true; } if (!errors) { @@ -514,10 +516,14 @@ public class ChannelPublishingJmsMessageListener return super.sendAndReceiveMessage(request); } - public ErrorMessage buildErrorMessage(Throwable throwable) { + protected ErrorMessage buildErrorMessage(Throwable throwable) { return super.buildErrorMessage(null, throwable); } + protected MessagingTemplate getMessagingTemplate() { + return this.messagingTemplate; + } + @Override public String getComponentType() { if (ChannelPublishingJmsMessageListener.this.expectReply) {