From 87d8e86297f2ef7ac2945e6e579f109dd34e36bb Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 6 Dec 2013 17:22:40 -0500 Subject: [PATCH] INT-3227 Fix Race Condition for WARN Log JIRA: https://jira.springsource.org/browse/INT-3227 In MessagingTemplate.TemporaryReplyChannel, it was possible the `clientHasReceived` boolean could be set by the other thread after normal message receipt, but before the sender tests it for the purpose of emitting a duplicate reply log message. Capture the boolean value of `clientHasReceived` before counting down the latch. --- .../integration/core/MessagingTemplate.java | 54 ++++++++++++++----- 1 file changed, 42 insertions(+), 12 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/core/MessagingTemplate.java b/spring-integration-core/src/main/java/org/springframework/integration/core/MessagingTemplate.java index 41df197255..63d9a48d0d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/core/MessagingTemplate.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/core/MessagingTemplate.java @@ -130,6 +130,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, this.receiveTimeout = receiveTimeout; } + @Override public void setBeanFactory(BeanFactory beanFactory) { if (this.channelResolver == null && beanFactory != null) { this.channelResolver = new BeanFactoryChannelResolver(beanFactory); @@ -148,6 +149,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, this.throwExceptionOnLateReply = throwExceptionOnLateReply; } + @Override public void afterPropertiesSet() { synchronized (this.initializationMonitor) { if (this.initialized) { @@ -157,18 +159,22 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } } + @Override public

void send(final Message

message) { this.send(this.getRequiredDefaultChannel(), message); } + @Override public

void send(final MessageChannel channel, final Message

message) { this.doSend(channel, message); } + @Override public

void send(final String channelName, final Message

message) { this.send(this.resolveChannelName(channelName), message); } + @Override public void convertAndSend(T object) { Message message = this.messageConverter.toMessage(object); if (message != null) { @@ -176,6 +182,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } } + @Override public void convertAndSend(MessageChannel channel, T object) { Message message = this.messageConverter.toMessage(object); if (message != null) { @@ -183,6 +190,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } } + @Override public void convertAndSend(String channelName, T object) { Message message = this.messageConverter.toMessage(object); if (message != null) { @@ -190,6 +198,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } } + @Override public void convertAndSend(T object, MessagePostProcessor postProcessor) { Message message = this.messageConverter.toMessage(object); message = postProcessor.postProcessMessage(message); @@ -198,6 +207,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } } + @Override public void convertAndSend(MessageChannel channel, T object, MessagePostProcessor postProcessor) { Message message = this.messageConverter.toMessage(object); message = postProcessor.postProcessMessage(message); @@ -206,6 +216,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } } + @Override public void convertAndSend(String channelName, T object, MessagePostProcessor postProcessor) { Message message = this.messageConverter.toMessage(object); message = postProcessor.postProcessMessage(message); @@ -214,6 +225,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } } + @Override public

Message

receive() { MessageChannel channel = this.getRequiredDefaultChannel(); Assert.state(channel instanceof PollableChannel, @@ -221,10 +233,12 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, return this.receive((PollableChannel) channel); } + @Override public

Message

receive(final PollableChannel channel) { return this.doReceive(channel); } + @Override public

Message

receive(String channelName) { MessageChannel channel = this.resolveChannelName(channelName); Assert.isInstanceOf(PollableChannel.class, channel, @@ -232,51 +246,61 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, return this.receive((PollableChannel) channel); } + @Override public Object receiveAndConvert() throws MessagingException { Message message = this.receive(); return (message != null) ? this.messageConverter.fromMessage(message) : null; } + @Override public Object receiveAndConvert(PollableChannel channel) throws MessagingException { Message message = this.receive(channel); return (message != null) ? this.messageConverter.fromMessage(message) : null; } + @Override public Object receiveAndConvert(String channelName) throws MessagingException { Message message = this.receive(channelName); return (message != null) ? this.messageConverter.fromMessage(message) : null; } + @Override public Message sendAndReceive(final Message requestMessage) { return this.sendAndReceive(this.getRequiredDefaultChannel(), requestMessage); } + @Override public Message sendAndReceive(final MessageChannel channel, final Message requestMessage) { return this.doSendAndReceive(channel, requestMessage); } + @Override public Message sendAndReceive(final String channelName, final Message requestMessage) { return this.sendAndReceive(this.resolveChannelName(channelName), requestMessage); } + @Override public Object convertSendAndReceive(final Object request) { Message requestMessage = this.messageConverter.toMessage(request); Message replyMessage = this.sendAndReceive(requestMessage); return this.messageConverter.fromMessage(replyMessage); } + @Override public Object convertSendAndReceive(final MessageChannel channel, final Object request) { Message requestMessage = this.messageConverter.toMessage(request); Message replyMessage = this.sendAndReceive(channel, requestMessage); return this.messageConverter.fromMessage(replyMessage); } + @Override public Object convertSendAndReceive(final String channelName, final Object request) { Message requestMessage = this.messageConverter.toMessage(request); Message replyMessage = this.sendAndReceive(channelName, requestMessage); return this.messageConverter.fromMessage(replyMessage); } + @Override public Object convertSendAndReceive(final Object request, MessagePostProcessor requestPostProcessor) { Message requestMessage = this.messageConverter.toMessage(request); requestMessage = requestPostProcessor.postProcessMessage(requestMessage); @@ -284,6 +308,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, return this.messageConverter.fromMessage(replyMessage); } + @Override public Object convertSendAndReceive(final MessageChannel channel, final Object request, MessagePostProcessor requestPostProcessor) { Message requestMessage = this.messageConverter.toMessage(request); requestMessage = requestPostProcessor.postProcessMessage(requestMessage); @@ -291,6 +316,7 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, return this.messageConverter.fromMessage(replyMessage); } + @Override public Object convertSendAndReceive(final String channelName, final Object request, MessagePostProcessor requestPostProcessor) { Message requestMessage = this.messageConverter.toMessage(request); requestMessage = requestPostProcessor.postProcessMessage(requestMessage); @@ -405,10 +431,12 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, } + @Override public Message receive() { return this.receive(-1); } + @Override public Message receive(long timeout) { try { if (this.receiveTimeout < 0) { @@ -430,25 +458,27 @@ public class MessagingTemplate implements MessagingOperations, BeanFactoryAware, return this.message; } + @Override public boolean send(Message message) { return this.send(message, -1); } + @Override public boolean send(Message message, long timeout) { this.message = message; + boolean clientHasReceived = this.clientHasReceived; this.latch.countDown(); - if (this.clientTimedOut || this.clientHasReceived || this.clientWontReceive) { - String exceptionMessage = ""; - if (this.clientTimedOut) { - exceptionMessage = "Reply message being sent, but the receiving thread has already timed out"; - } - else if (this.clientHasReceived) { - exceptionMessage = "Reply message being sent, but the receiving thread has already received a reply"; - } - else if (this.clientWontReceive) { - exceptionMessage = "Reply message being sent, but the receiving thread has already caught an exception and won't receive"; - } - + String exceptionMessage = null; + if (this.clientTimedOut) { + exceptionMessage = "Reply message being sent, but the receiving thread has already timed out"; + } + else if (clientHasReceived) { + exceptionMessage = "Reply message being sent, but the receiving thread has already received a reply"; + } + else if (this.clientWontReceive) { + exceptionMessage = "Reply message being sent, but the receiving thread has already caught an exception and won't receive"; + } + if (exceptionMessage != null) { if (logger.isWarnEnabled()) { logger.warn(exceptionMessage + ":" + message); }