diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessageChannelTemplate.java b/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessageChannelTemplate.java index 2c4dac235d..cbf139399d 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessageChannelTemplate.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessageChannelTemplate.java @@ -239,12 +239,15 @@ public class MessageChannelTemplate implements InitializingBean { } private Message doSendAndReceive(Message request, MessageChannel channel) { - TemporaryReturnAddress returnAddress = new TemporaryReturnAddress(this.receiveTimeout); - request = MessageBuilder.fromMessage(request).setReplyChannel(returnAddress).build(); + TemporaryReturnAddress replyChannel = new TemporaryReturnAddress(this.receiveTimeout); + request = MessageBuilder.fromMessage(request) + .setReplyChannel(replyChannel) + .setErrorChannel(replyChannel) + .build(); if (!this.doSend(request, channel)) { return null; } - return this.doReceive(returnAddress); + return this.doReceive(replyChannel); } private MessageChannel getRequiredDefaultChannel() { diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessagePublishingErrorHandler.java b/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessagePublishingErrorHandler.java index 2e4ac2e097..024cc9f2bb 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessagePublishingErrorHandler.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/channel/MessagePublishingErrorHandler.java @@ -58,43 +58,44 @@ public class MessagePublishingErrorHandler implements ErrorHandler { } public final void handle(Throwable t) { - Message failedMessage = null; - if (logger.isWarnEnabled()) { - if (t instanceof MessagingException) { - failedMessage = ((MessagingException) t).getFailedMessage(); - logger.warn("failure occurred in messaging task with message: " + failedMessage, t); - } - else { - logger.warn("failure occurred in messaging task", t); - } - } - MessageChannel errorChannel = null; - if (failedMessage != null) { - Object errorChannelHeader = failedMessage.getHeaders().getErrorChannel(); - if (errorChannelHeader != null) { - if (errorChannelHeader instanceof MessageChannel) { - errorChannel = (MessageChannel) errorChannelHeader; - } - else if (errorChannelHeader instanceof String) { - errorChannel = this.channelResolver.resolveChannelName((String) errorChannelHeader); - } - } - } - if (errorChannel == null) { - errorChannel = this.defaultErrorChannel; - } + Message failedMessage = (t instanceof MessagingException) ? + ((MessagingException) t).getFailedMessage() : null; + MessageChannel errorChannel = this.resolveErrorChannel(failedMessage); + boolean sent = false; if (errorChannel != null) { try { if (this.sendTimeout >= 0) { - errorChannel.send(new ErrorMessage(t), this.sendTimeout); + sent = errorChannel.send(new ErrorMessage(t), this.sendTimeout); } else { - errorChannel.send(new ErrorMessage(t)); + sent = errorChannel.send(new ErrorMessage(t)); } } catch (Throwable ignore) { // message will be logged only } } + if (!sent && logger.isErrorEnabled()) { + if (failedMessage != null) { + logger.error("failure occurred in messaging task with message: " + failedMessage, t); + } + else { + logger.error("failure occurred in messaging task", t); + } + } + } + + private MessageChannel resolveErrorChannel(Message failedMessage) { + if (failedMessage == null || failedMessage.getHeaders().getErrorChannel() == null) { + return this.defaultErrorChannel; + } + Object errorChannelHeader = failedMessage.getHeaders().getErrorChannel(); + if (errorChannelHeader instanceof MessageChannel) { + return (MessageChannel) errorChannelHeader; + } + Assert.isInstanceOf(String.class, errorChannelHeader, + "Unsupported error channel header type. Expected MessageChannel or String, but actual type is [" + + errorChannelHeader.getClass() + "]"); + return this.channelResolver.resolveChannelName((String) errorChannelHeader); } }