diff --git a/build.gradle b/build.gradle index c01edb144c..4119170c41 100644 --- a/build.gradle +++ b/build.gradle @@ -121,14 +121,14 @@ subprojects { subproject -> tomcatVersion = "7.0.55" smack3Version = '3.2.1' smackVersion = '4.0.0' - springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '1.4.1.RELEASE' + springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '1.4.2.RELEASE' springDataMongoVersion = '1.6.0.RELEASE' springDataRedisVersion = '1.4.0.RELEASE' springGemfireVersion = '1.5.0.RELEASE' springSecurityVersion = '3.2.5.RELEASE' springSocialTwitterVersion = '1.1.0.RELEASE' springRetryVersion = '1.1.1.RELEASE' - springVersion = project.hasProperty('springVersion') ? project.springVersion : '4.1.3.RELEASE' + springVersion = project.hasProperty('springVersion') ? project.springVersion : '4.1.4.RELEASE' springWsVersion = '2.2.0.RELEASE' xmlUnitVersion = '1.5' xstreamVersion = '1.4.7' diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java index c4f6ea6058..fb3bd2b6c7 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java @@ -16,11 +16,13 @@ package org.springframework.integration.amqp.support; +import java.lang.reflect.Field; import java.util.ArrayList; import java.util.Date; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; import org.springframework.amqp.core.MessageDeliveryMode; import org.springframework.amqp.core.MessageProperties; @@ -28,6 +30,9 @@ import org.springframework.amqp.support.AmqpHeaders; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.mapping.AbstractHeaderMapper; import org.springframework.integration.mapping.support.JsonHeaders; +import org.springframework.util.ReflectionUtils; +import org.springframework.util.ReflectionUtils.FieldCallback; +import org.springframework.util.ReflectionUtils.FieldFilter; import org.springframework.util.StringUtils; /** @@ -53,6 +58,8 @@ import org.springframework.util.StringUtils; */ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper implements AmqpHeaderMapper { + static final boolean CONSUMER_METADATA_PRESENT; + private static final List STANDARD_HEADER_NAMES = new ArrayList(); static { @@ -79,6 +86,27 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper toHeadersFromRequest(MessageProperties source) { + Map headersFromRequest = super.toHeadersFromRequest(source); + if (CONSUMER_METADATA_PRESENT) { + addConsumerMetadata(source, headersFromRequest); + } + return headersFromRequest; + } + + private void addConsumerMetadata(MessageProperties messageProperties, Map headers) { + String consumerTag = messageProperties.getConsumerTag(); + if (consumerTag != null) { + headers.put(AmqpHeaders.CONSUMER_TAG, consumerTag); + } + String consumerQueue = messageProperties.getConsumerQueue(); + if (consumerQueue != null) { + headers.put(AmqpHeaders.CONSUMER_QUEUE, consumerQueue); + } + } + } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java index 47dc7599c0..205dc17e90 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java @@ -171,6 +171,17 @@ public class DefaultAmqpHeaderMapperTests { assertEquals("test.replyTo2", headerMap.get(AmqpHeaders.SPRING_REPLY_TO_STACK)); } + @Test // INT-3586 requires Spring AMQP 1.4.2 + public void testToHeadersConsumerMetadata() { + DefaultAmqpHeaderMapper headerMapper = new DefaultAmqpHeaderMapper(); + MessageProperties amqpProperties = new MessageProperties(); + amqpProperties.setConsumerTag("consumerTag"); + amqpProperties.setConsumerQueue("consumerQueue"); + Map headerMap = headerMapper.toHeadersFromRequest(amqpProperties); + assertEquals("consumerTag", headerMap.get(AmqpHeaders.CONSUMER_TAG)); + assertEquals("consumerQueue", headerMap.get(AmqpHeaders.CONSUMER_QUEUE)); + } + @Test public void messageIdNotMappedToAmqpProperties() { DefaultAmqpHeaderMapper headerMapper = new DefaultAmqpHeaderMapper(); diff --git a/src/reference/docbook/index.xml b/src/reference/docbook/index.xml index d6adcd4a02..04873323d1 100644 --- a/src/reference/docbook/index.xml +++ b/src/reference/docbook/index.xml @@ -55,8 +55,9 @@ 2012 2013 2014 + 2015 - GoPivotal, Inc. All Rights Reserved. + Pivotal, Inc. All Rights Reserved.