diff --git a/build.gradle b/build.gradle index a5ee1a43c5..02860ad638 100644 --- a/build.gradle +++ b/build.gradle @@ -131,14 +131,14 @@ subprojects { subproject -> romeToolsVersion = '1.9.0' servletApiVersion = '4.0.0' smackVersion = '4.3.1' - springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.1.0.RELEASE' + springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.1.2.BUILD-SNAPSHOT' springDataJpaVersion = '2.1.2.RELEASE' springDataMongoVersion = '2.1.2.RELEASE' springDataRedisVersion = '2.1.2.RELEASE' springGemfireVersion = '2.1.2.RELEASE' springSecurityVersion = '5.1.1.RELEASE' springRetryVersion = '1.2.2.RELEASE' - springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.1.2.RELEASE' + springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.1.3.BUILD-SNAPSHOT' springWsVersion = '3.0.3.RELEASE' tomcatVersion = "9.0.12" xmlUnitVersion = '1.6' diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpPollableMessageChannelSpec.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpPollableMessageChannelSpec.java index ce33c2ed06..5f85c1d40a 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpPollableMessageChannelSpec.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/dsl/AmqpPollableMessageChannelSpec.java @@ -27,6 +27,7 @@ import org.springframework.integration.amqp.channel.AbstractAmqpChannel; import org.springframework.integration.amqp.config.AmqpChannelFactoryBean; import org.springframework.integration.amqp.support.AmqpHeaderMapper; import org.springframework.integration.dsl.MessageChannelSpec; +import org.springframework.lang.Nullable; import org.springframework.util.Assert; /** @@ -62,7 +63,7 @@ public class AmqpPollableMessageChannelSpec requestMessage) { CorrelationData correlationData = null; if (this.correlationDataGenerator != null) { - correlationData = new CorrelationDataWrapper(requestMessage.getHeaders().getId().toString(), + UUID messageId = requestMessage.getHeaders().getId(); + if (messageId == null) { + messageId = NO_ID; + } + correlationData = new CorrelationDataWrapper(messageId.toString(), this.correlationDataGenerator.processMessage(requestMessage), requestMessage); } return correlationData; diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java index 5ddf4f4b94..5d49148f69 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java @@ -37,6 +37,7 @@ import static org.mockito.Mockito.when; import java.io.IOException; import java.lang.reflect.Field; import java.util.List; +import java.util.concurrent.ExecutorService; import java.util.concurrent.atomic.AtomicBoolean; import org.apache.commons.logging.Log; @@ -238,7 +239,8 @@ public class AmqpOutboundChannelAdapterParserTests { Channel mockChannel = mock(Channel.class); when(connectionFactory.createConnection()).thenReturn(mockConnection); - PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel); + PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel, + mock(ExecutorService.class)); when(mockConnection.createChannel(false)).thenReturn(publisherCallbackChannel); MessageChannel requestChannel = context.getBean("toRabbitOnlyWithTemplateChannel", MessageChannel.class); @@ -255,7 +257,8 @@ public class AmqpOutboundChannelAdapterParserTests { Channel mockChannel = mock(Channel.class); when(connectionFactory.createConnection()).thenReturn(mockConnection); - PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel); + PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel, + mock(ExecutorService.class)); when(mockConnection.createChannel(false)).thenReturn(publisherCallbackChannel); MessageChannel requestChannel = context.getBean("withDefaultAmqpTemplateExchangeAndRoutingKey", @@ -272,7 +275,8 @@ public class AmqpOutboundChannelAdapterParserTests { Channel mockChannel = mock(Channel.class); when(connectionFactory.createConnection()).thenReturn(mockConnection); - PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel); + PublisherCallbackChannelImpl publisherCallbackChannel = new PublisherCallbackChannelImpl(mockChannel, + mock(ExecutorService.class)); when(mockConnection.createChannel(false)).thenReturn(publisherCallbackChannel); MessageChannel requestChannel = context.getBean("overrideTemplateAttributesToEmpty", MessageChannel.class); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/IntegrationMessageHeaderAccessor.java b/spring-integration-core/src/main/java/org/springframework/integration/IntegrationMessageHeaderAccessor.java index 11c2593d75..93ee151d30 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/IntegrationMessageHeaderAccessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/IntegrationMessageHeaderAccessor.java @@ -27,7 +27,6 @@ import java.util.concurrent.atomic.AtomicInteger; import org.springframework.integration.acks.AcknowledgmentCallback; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.MessageHeaderAccessor; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; @@ -74,10 +73,11 @@ public class IntegrationMessageHeaderAccessor extends MessageHeaderAccessor { } /** - * Specify a list of headers which should be considered as read only - * and prohibited from being populated in the message. - * @param readOnlyHeaders the list of headers for {@code readOnly} mode. - * Defaults to {@link MessageHeaders#ID} and {@link MessageHeaders#TIMESTAMP}. + * Specify a list of headers which should be considered as read only and prohibited + * from being populated in the message. + * @param readOnlyHeaders the list of headers for {@code readOnly} mode. Defaults to + * {@link org.springframework.messaging.MessageHeaders#ID} and + * {@link org.springframework.messaging.MessageHeaders#TIMESTAMP}. * @since 4.3.2 * @see #isReadOnly(String) */