From db692b1d8c0e0c164efca579b3063bd29d8c66c4 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 13 Jun 2017 09:31:37 -0400 Subject: [PATCH] Revert to Spring AMQP 2.0.0.M4 Since we don't have one more Spring AMQP 2.0 Milestone for upcoming Spring Boot Milestone Spring Integration Milestone must be compatible with current IO state. * Replace `AmqpHeaders.RAW_MESSAGE` to the `AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE` constant --- build.gradle | 2 +- .../amqp/inbound/AmqpInboundChannelAdapter.java | 2 +- .../integration/amqp/inbound/AmqpInboundGateway.java | 2 +- .../support/AmqpMessageHeaderErrorMessageStrategy.java | 8 +++++++- .../integration/amqp/inbound/InboundEndpointTests.java | 4 ++-- 5 files changed, 12 insertions(+), 6 deletions(-) diff --git a/build.gradle b/build.gradle index 608c513edb..fee2b24b9d 100644 --- a/build.gradle +++ b/build.gradle @@ -130,7 +130,7 @@ subprojects { subproject -> servletApiVersion = '3.1.0' slf4jVersion = "1.7.25" smackVersion = '4.1.9' - springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.0.0.BUILD-SNAPSHOT' + springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.0.0.M4' springDataJpaVersion = '2.0.0.BUILD-SNAPSHOT' springDataMongoVersion = '2.0.0.BUILD-SNAPSHOT' springDataRedisVersion = '2.0.0.BUILD-SNAPSHOT' 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 cc33b21969..e78eb29cb0 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 @@ -176,7 +176,7 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements AttributeAccessor attributes = attributesHolder.get(); if (attributes != null) { attributes.setAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY, message); - attributes.setAttribute(AmqpHeaders.RAW_MESSAGE, amqpMessage); + attributes.setAttribute(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE, amqpMessage); } } } 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 5b158fa445..0397fa4cc0 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 @@ -232,7 +232,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { AttributeAccessor attributes = attributesHolder.get(); if (attributes != null) { attributes.setAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY, message); - attributes.setAttribute(AmqpHeaders.RAW_MESSAGE, amqpMessage); + attributes.setAttribute(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE, amqpMessage); } } } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java index 35eaca1033..1fcf01913b 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java @@ -38,11 +38,17 @@ import org.springframework.messaging.support.ErrorMessage; */ public class AmqpMessageHeaderErrorMessageStrategy implements ErrorMessageStrategy { + /** + * Header name/retry context variable for the raw received message. + */ + public static final String AMQP_RAW_MESSAGE = AmqpHeaders.PREFIX + "raw_message"; + + @SuppressWarnings("deprecation") @Override public ErrorMessage buildErrorMessage(Throwable throwable, AttributeAccessor context) { Object inputMessage = context.getAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY); Map headers = - Collections.singletonMap(AmqpHeaders.RAW_MESSAGE, context.getAttribute(AmqpHeaders.RAW_MESSAGE)); + Collections.singletonMap(AMQP_RAW_MESSAGE, context.getAttribute(AMQP_RAW_MESSAGE)); return new ErrorMessage(throwable, headers, inputMessage instanceof Message ? (Message) inputMessage : null); } 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 3cb313bfde..25f5b29e67 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 @@ -309,7 +309,7 @@ public class InboundEndpointTests { assertThat(errorMessage.getPayload(), instanceOf(MessagingException.class)); assertThat(((MessagingException) errorMessage.getPayload()).getMessage(), containsString("Dispatcher has no")); org.springframework.amqp.core.Message amqpMessage = errorMessage.getHeaders() - .get(AmqpHeaders.RAW_MESSAGE, org.springframework.amqp.core.Message.class); + .get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE, org.springframework.amqp.core.Message.class); assertThat(amqpMessage, notNullValue()); assertNull(errors.receive(0)); } @@ -334,7 +334,7 @@ public class InboundEndpointTests { assertThat(errorMessage.getPayload(), instanceOf(MessagingException.class)); assertThat(((MessagingException) errorMessage.getPayload()).getMessage(), containsString("Dispatcher has no")); org.springframework.amqp.core.Message amqpMessage = errorMessage.getHeaders() - .get(AmqpHeaders.RAW_MESSAGE, org.springframework.amqp.core.Message.class); + .get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE, org.springframework.amqp.core.Message.class); assertThat(amqpMessage, notNullValue()); assertNull(errors.receive(0)); }