From 756cd41085cd88095aa159b34e797e76396ca68f Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 8 May 2019 13:35:11 -0400 Subject: [PATCH] Fix MessagingGatewaySupport implementors The `AmqpInboundGateway` and `JmsInboundGateway` don't delegate to `super` in their `doStart()/doStop()` implementations causing the problem with an internal `replyMessageCorrelator` missed the proper lifecycle management **Cherry-pick to 5.1.x & 5.0.x** --- .../amqp/inbound/AmqpInboundGateway.java | 20 ++++++++++--------- .../integration/jms/JmsInboundGateway.java | 3 +++ 2 files changed, 14 insertions(+), 9 deletions(-) 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 0808343aa8..a0b9f58949 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 @@ -65,7 +65,7 @@ import com.rabbitmq.client.Channel; */ public class AmqpInboundGateway extends MessagingGatewaySupport { - private static final ThreadLocal attributesHolder = new ThreadLocal(); + private static final ThreadLocal attributesHolder = new ThreadLocal<>(); private final AbstractMessageListenerContainer messageListenerContainer; @@ -208,11 +208,13 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { @Override protected void doStart() { + super.doStart(); this.messageListenerContainer.start(); } @Override protected void doStop() { + super.doStop(); this.messageListenerContainer.stop(); } @@ -271,11 +273,11 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { org.springframework.messaging.Message converted = convert(message, channel); if (converted != null) { AmqpInboundGateway.this.retryTemplate.execute(context -> { - StaticMessageHeaderAccessor.getDeliveryAttempt(converted).incrementAndGet(); - process(message, converted); - return null; - }, - (RecoveryCallback) AmqpInboundGateway.this.recoveryCallback); + StaticMessageHeaderAccessor.getDeliveryAttempt(converted).incrementAndGet(); + process(message, converted); + return null; + }, + (RecoveryCallback) AmqpInboundGateway.this.recoveryCallback); } } } @@ -305,9 +307,9 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { return null; } return getMessageBuilderFactory() - .withPayload(payload) - .copyHeaders(headers) - .build(); + .withPayload(payload) + .copyHeaders(headers) + .build(); } private void process(Message message, org.springframework.messaging.Message messagingMessage) { diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java index d9cb7da185..c16d093430 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java @@ -39,6 +39,7 @@ public class JmsInboundGateway extends MessagingGatewaySupport implements Dispos public JmsInboundGateway(AbstractMessageListenerContainer listenerContainer, ChannelPublishingJmsMessageListener listener) { + this.endpoint = new JmsMessageDrivenEndpoint(listenerContainer, listener); } @@ -78,11 +79,13 @@ public class JmsInboundGateway extends MessagingGatewaySupport implements Dispos @Override protected void doStart() { + super.doStart(); this.endpoint.start(); } @Override protected void doStop() { + super.doStop(); this.endpoint.stop(); }