From 80d679a9b0537cf71371502afeb6f52b851b31d4 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 23 Jul 2019 14:01:05 -0400 Subject: [PATCH] GH-2748: More bean definitions into exceptions Fixes https://github.com/spring-projects/spring-integration/issues/2748 * Refactor more `MessageHandlingException`s to include `this` into an exception message * Revert using `MessagingException` in some places which really are not about messaging. This helps to wrap them into `MessageHandlingException` later in the `MessageHandler` for the `BeanDefinition` reference * Remove `volatile` from configuration properties in the affected classes * Remove already deprecated `JmsOutboundGateway.setPriority()` * Add `resource` and `source` for `BeanDefinition` in the `AbstractChannelAdapterParser` & `AbstractInboundGatewayParser` * Document the feature --- .../outbound/AsyncAmqpOutboundGateway.java | 5 +- .../AbstractCorrelatingMessageHandler.java | 2 +- .../aggregator/BarrierMessageHandler.java | 3 +- .../xml/AbstractChannelAdapterParser.java | 5 +- .../xml/AbstractInboundGatewayParser.java | 10 +- .../context/IntegrationObjectSupport.java | 2 +- .../AbstractMessageProducingHandler.java | 3 +- .../integration/handler/DelayHandler.java | 3 +- .../MethodInvokingMessageProcessor.java | 4 +- .../transformer/ContentEnricher.java | 56 ++-- .../file/FileWritingMessageHandler.java | 5 +- .../AbstractRemoteFileOutboundGateway.java | 9 +- .../file/splitter/FileSplitter.java | 3 +- .../FileOutboundGatewayParserTests.java | 4 +- .../file/splitter/FileSplitterTests.java | 2 +- .../ftp/gateway/FtpOutboundGateway.java | 23 +- ...ovyScriptPayloadMessageProcessorTests.java | 16 +- ...actHttpRequestExecutingMessageHandler.java | 30 +- .../HttpRequestExecutingMessageHandler.java | 8 +- .../ip/tcp/TcpOutboundGateway.java | 11 +- .../ip/tcp/TcpSendingMessageHandler.java | 14 +- .../ip/tcp/connection/TcpMessageMapper.java | 8 +- .../ip/udp/DatagramPacketMessageMapper.java | 14 +- .../ip/udp/UnicastSendingMessageHandler.java | 2 +- .../jdbc/StoredProcOutboundGateway.java | 8 +- .../integration/jms/JmsOutboundGateway.java | 304 ++++++++---------- .../jms/config/JmsOutboundGatewayParser.java | 2 +- .../jms/dsl/JmsOutboundGatewaySpec.java | 2 +- .../NotificationPublishingMessageHandler.java | 8 +- .../jmx/OperationInvokingMessageHandler.java | 9 +- .../mail/MailSendingMessageHandler.java | 20 +- .../outbound/AbstractMqttMessageHandler.java | 5 +- .../mqtt/outbound/MqttPahoMessageHandler.java | 11 +- .../support/DefaultPahoMessageConverter.java | 3 +- .../RedisStoreWritingMessageHandler.java | 2 +- .../integration/rmi/RmiOutboundGateway.java | 14 +- ...stractScriptExecutingMessageProcessor.java | 12 +- .../inbound/StompInboundChannelAdapter.java | 2 +- .../WebSocketInboundChannelAdapter.java | 8 +- .../WebSocketOutboundMessageHandler.java | 6 +- .../XmlValidatingMessageSelector.java | 18 +- .../xml/splitter/XPathMessageSplitter.java | 10 +- .../ChatMessageSendingMessageHandler.java | 14 +- .../PresenceSendingMessageHandler.java | 12 +- src/reference/asciidoc/configuration.adoc | 7 + src/reference/asciidoc/whats-new.adoc | 3 + 46 files changed, 352 insertions(+), 370 deletions(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java index 4ab7e479b1..d649afd1db 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java @@ -64,7 +64,7 @@ public class AsyncAmqpOutboundGateway extends AbstractAmqpOutboundEndpoint { @Override protected RabbitTemplate getRabbitTemplate() { - return this.template.getRabbitTemplate(); + return this.template.getRabbitTemplate(); } @Override @@ -112,7 +112,8 @@ public class AsyncAmqpOutboundGateway extends AbstractAmqpOutboundEndpoint { catch (Exception e) { Exception exceptionToLogAndSend = e; if (!(e instanceof MessagingException)) { - exceptionToLogAndSend = new MessageHandlingException(this.requestMessage, e); + exceptionToLogAndSend = new MessageHandlingException(this.requestMessage, + "failed to handle a message in the [" + AsyncAmqpOutboundGateway.this + ']', e); if (replyMessageBuilder != null) { exceptionToLogAndSend = new MessagingException(replyMessageBuilder.build(), exceptionToLogAndSend); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index f05e5fe07b..c9102d8974 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -454,7 +454,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP } catch (InterruptedException e) { Thread.currentThread().interrupt(); - throw new MessageHandlingException(message, "Interrupted getting lock", e); + throw new MessageHandlingException(message, "Interrupted getting lock in the [" + this + ']', e); } try { noOutput = processMessageForGroup(message, correlationKey, groupIdUuid, lock); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/BarrierMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/BarrierMessageHandler.java index 74a1d169b2..da566ef981 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/BarrierMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/BarrierMessageHandler.java @@ -169,7 +169,8 @@ public class BarrierMessageHandler extends AbstractReplyProducingMessageHandler } catch (InterruptedException e) { Thread.currentThread().interrupt(); - throw new MessageHandlingException(requestMessage, "Interrupted while waiting for release", e); + throw new MessageHandlingException(requestMessage, + "Interrupted while waiting for release in the [" + this + ']', e); } finally { this.inProcess.remove(key); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractChannelAdapterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractChannelAdapterParser.java index 3c9a860e48..8b02d2ea76 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractChannelAdapterParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractChannelAdapterParser.java @@ -48,6 +48,7 @@ public abstract class AbstractChannelAdapterParser extends AbstractBeanDefinitio @Override protected final String resolveId(Element element, AbstractBeanDefinition definition, ParserContext parserContext) throws BeanDefinitionStoreException { + String id = element.getAttribute(ID_ATTRIBUTE); if (!element.hasAttribute("channel")) { // the created channel will get the 'id', so the adapter's bean name includes a suffix @@ -82,6 +83,8 @@ public abstract class AbstractChannelAdapterParser extends AbstractBeanDefinitio propertyValues.add("role", new TypedStringValue(role)); } } + beanDefinition.setResource(parserContext.getReaderContext().getResource()); + beanDefinition.setSource(IntegrationNamespaceUtils.createElementDescription(element)); return beanDefinition; } @@ -89,14 +92,12 @@ public abstract class AbstractChannelAdapterParser extends AbstractBeanDefinitio if (parserContext.isNested()) { return null; } - return IntegrationNamespaceUtils.createDirectChannel(element, parserContext); } /** * Subclasses must implement this method to parse the adapter element. * The name of the MessageChannel bean is provided. - * * @param element The element. * @param parserContext The parser context. * @param channelName The channel name. diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractInboundGatewayParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractInboundGatewayParser.java index 346256b701..7e7f634900 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractInboundGatewayParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractInboundGatewayParser.java @@ -56,6 +56,14 @@ public abstract class AbstractInboundGatewayParser extends AbstractSimpleBeanDef && !attributeName.equals("reply-channel") && super.isEligibleAttribute(attributeName); } + @Override + protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { + super.doParse(element, parserContext, builder); + AbstractBeanDefinition beanDefinition = builder.getRawBeanDefinition(); + beanDefinition.setResource(parserContext.getReaderContext().getResource()); + beanDefinition.setSource(IntegrationNamespaceUtils.createElementDescription(element)); + } + @Override protected final void postProcess(BeanDefinitionBuilder builder, Element element) { String requestChannelRef = element.getAttribute("request-channel"); @@ -69,7 +77,7 @@ public abstract class AbstractInboundGatewayParser extends AbstractSimpleBeanDef if (StringUtils.hasText(errorChannel)) { builder.addPropertyValue("errorChannelName", errorChannel); } - this.doPostProcess(builder, element); + doPostProcess(builder, element); } /** diff --git a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java index 76f503ce97..300f220b0f 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationObjectSupport.java @@ -153,7 +153,7 @@ public abstract class IntegrationObjectSupport implements BeanNameAware, NamedCo StringBuilder sb = new StringBuilder("bean '") .append(this.beanName).append("'"); if (!this.beanName.equals(getComponentName())) { - sb.append("for component '").append(getComponentName()).append("'"); + sb.append(" for component '").append(getComponentName()).append("'"); } if (description != null) { sb.append("; defined in: '").append(description).append("'"); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java index bb91257308..6dae2b8b85 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageProducingHandler.java @@ -475,7 +475,8 @@ public abstract class AbstractMessageProducingHandler extends AbstractMessageHan } catch (Exception e) { Exception exceptionToLog = - IntegrationUtils.wrapInHandlingExceptionIfNecessary(requestMessage, () -> null, e); + IntegrationUtils.wrapInHandlingExceptionIfNecessary(requestMessage, + () -> "failed to send error message in the [" + this + ']', e); logger.error("Failed to send async reply", exceptionToLog); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java index 10293a737c..edf1f67878 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java @@ -46,7 +46,6 @@ import org.springframework.jmx.export.annotation.ManagedResource; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; -import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessagingException; import org.springframework.messaging.core.DestinationResolver; import org.springframework.messaging.support.ErrorMessage; @@ -381,7 +380,7 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement } } else { - throw new MessageHandlingException(message, "Error occurred during 'delay' value determination", + throw new IllegalStateException("Error occurred during 'delay' value determination", delayValueException); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/MethodInvokingMessageProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/MethodInvokingMessageProcessor.java index b01bc252ed..2b7bfbcabb 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/MethodInvokingMessageProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/MethodInvokingMessageProcessor.java @@ -109,8 +109,8 @@ public class MethodInvokingMessageProcessor extends AbstractMessageProcessor< } catch (Exception ex) { throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "error occurred during processing message in 'MethodInvokingMessageProcessor' [" + this + - "]", ex); + () -> "error occurred during processing message in 'MethodInvokingMessageProcessor' [" + this + ']', + ex); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java b/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java index 7061298c9e..c35c93afe5 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java @@ -62,50 +62,52 @@ public class ContentEnricher extends AbstractReplyProducingMessageHandler implem */ private final SpelExpressionParser parser = new SpelExpressionParser(new SpelParserConfiguration(true, true)); - private volatile Map nullResultPropertyExpressions = new HashMap<>(); + private Map nullResultPropertyExpressions = new HashMap<>(); - private volatile Map> nullResultHeaderExpressions = - new HashMap<>(); + private Map> nullResultHeaderExpressions = new HashMap<>(); - private volatile Map propertyExpressions = new HashMap<>(); + private Map propertyExpressions = new HashMap<>(); - private volatile Map> headerExpressions = - new HashMap<>(); + private Map> headerExpressions = new HashMap<>(); private EvaluationContext sourceEvaluationContext; private EvaluationContext targetEvaluationContext; - private volatile boolean shouldClonePayload = false; + private boolean shouldClonePayload = false; private Expression requestPayloadExpression; - private volatile MessageChannel requestChannel; + private MessageChannel requestChannel; - private volatile String requestChannelName; + private String requestChannelName; - private volatile MessageChannel replyChannel; + private MessageChannel replyChannel; - private volatile String replyChannelName; + private String replyChannelName; - private volatile MessageChannel errorChannel; + private MessageChannel errorChannel; - private volatile String errorChannelName; + private String errorChannelName; - private volatile Gateway gateway = null; + private Gateway gateway; - private volatile Long requestTimeout; + private Long requestTimeout; - private volatile Long replyTimeout; + private Long replyTimeout; public void setNullResultPropertyExpressions(Map nullResultPropertyExpressions) { - Map localMap = new HashMap<>(nullResultPropertyExpressions.size()); - for (Map.Entry entry : nullResultPropertyExpressions.entrySet()) { + this.nullResultPropertyExpressions = convertExpressions(nullResultPropertyExpressions); + } + + private Map convertExpressions(Map expressions) { + Map localMap = new HashMap<>(expressions.size()); + for (Map.Entry entry : expressions.entrySet()) { String key = entry.getKey(); Expression value = entry.getValue(); localMap.put(this.parser.parseExpression(key), value); } - this.nullResultPropertyExpressions = localMap; + return localMap; } public void setNullResultHeaderExpressions(Map> nullResultHeaderExpressions) { @@ -122,13 +124,7 @@ public class ContentEnricher extends AbstractReplyProducingMessageHandler implem Assert.notEmpty(propertyExpressions, "propertyExpressions must not be empty"); Assert.noNullElements(propertyExpressions.keySet().toArray(), "propertyExpressions keys must not be empty"); Assert.noNullElements(propertyExpressions.values().toArray(), "propertyExpressions values must not be empty"); - Map localMap = new HashMap<>(propertyExpressions.size()); - for (Map.Entry entry : propertyExpressions.entrySet()) { - String key = entry.getKey(); - Expression value = entry.getValue(); - localMap.put(this.parser.parseExpression(key), value); - } - this.propertyExpressions = localMap; + this.propertyExpressions = convertExpressions(propertyExpressions); } /** @@ -337,7 +333,8 @@ public class ContentEnricher extends AbstractReplyProducingMessageHandler implem } } - for (Map.Entry> entry : this.nullResultHeaderExpressions.entrySet()) { + for (Map.Entry> entry : + this.nullResultHeaderExpressions.entrySet()) { if (checkReadOnlyHeaders && (MessageHeaders.ID.equals(entry.getKey()) || MessageHeaders.TIMESTAMP.equals(entry.getKey()))) { throw new BeanInitializationException( @@ -362,7 +359,8 @@ public class ContentEnricher extends AbstractReplyProducingMessageHandler implem targetPayload = ReflectionUtils.invokeMethod(cloneMethod, requestPayload); } catch (Exception e) { - throw new MessageHandlingException(requestMessage, "Failed to clone payload object", e); + throw new MessageHandlingException(requestMessage, + "Failed to clone payload object in the [" + this + ']', e); } } else { @@ -400,7 +398,7 @@ public class ContentEnricher extends AbstractReplyProducingMessageHandler implem return targetPayload; } else { - Map targetHeaders = new HashMap( + Map targetHeaders = new HashMap<>( this.nullResultHeaderExpressions.size()); for (Map.Entry> entry : this.nullResultHeaderExpressions .entrySet()) { diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java b/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java index 4bf526469d..a996dfec8a 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java @@ -503,7 +503,8 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand boolean exists = resultFile.exists(); if (exists && FileExistsMode.FAIL.equals(this.fileExistsMode)) { throw new MessageHandlingException(requestMessage, - "The destination file already exists at '" + resultFile.getAbsolutePath() + "'."); + "Failed to process message in the [" + this + + "]. The destination file already exists at '" + resultFile.getAbsolutePath() + "'."); } Object timestamp = requestMessage.getHeaders().get(FileHeaders.SET_MODIFIED); @@ -528,7 +529,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } catch (Exception e) { throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(requestMessage, - () -> "failed to write Message payload to file", e); + () -> "failed to write Message payload to file in the [" + this + ']', e); } } diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/gateway/AbstractRemoteFileOutboundGateway.java b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/gateway/AbstractRemoteFileOutboundGateway.java index fbdf7a395c..3822eb6692 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/gateway/AbstractRemoteFileOutboundGateway.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/gateway/AbstractRemoteFileOutboundGateway.java @@ -543,9 +543,9 @@ public abstract class AbstractRemoteFileOutboundGateway extends AbstractReply payload = session.readRaw(remoteFilePath); } catch (IOException e) { - throw new MessageHandlingException(requestMessage, "Failed to get the remote file [" - + remoteFilePath - + "] as a stream", e); + throw new MessageHandlingException(requestMessage, + "Error handling message in the [" + this + + "]. Failed to get the remote file [" + remoteFilePath + "] as a stream", e); } } else { @@ -984,7 +984,8 @@ public abstract class AbstractRemoteFileOutboundGateway extends AbstractReply } } else if (!FileExistsMode.IGNORE.equals(existsMode)) { - throw new MessageHandlingException(message, "Local file " + localFile + " already exists"); + throw new MessageHandlingException(message, + "Error handling message in the [" + this + "]. Local file " + localFile + " already exists"); } else { if (logger.isDebugEnabled()) { diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java b/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java index a9cff1673a..3f419e9163 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java @@ -180,7 +180,8 @@ public class FileSplitter extends AbstractMessageSplitter { filePath = (String) payload; } catch (FileNotFoundException e) { - throw new MessageHandlingException(message, "failed to read file [" + payload + "]", e); + throw new MessageHandlingException(message, + "Error handing message in the [" + this + "]. Failed to read file [" + payload + "]", e); } } else if (payload instanceof File) { diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileOutboundGatewayParserTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileOutboundGatewayParserTests.java index 192e3e8a01..48c3fe4f9b 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileOutboundGatewayParserTests.java +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/config/FileOutboundGatewayParserTests.java @@ -181,7 +181,7 @@ public class FileOutboundGatewayParserTests { assertThatExceptionOfType(MessageHandlingException.class) .isThrownBy(() -> messagingTemplate.sendAndReceive(new GenericMessage<>("String content:"))) - .withMessageStartingWith("The destination file already exists at '"); + .withMessageContaining("The destination file already exists at '"); } /** @@ -208,7 +208,7 @@ public class FileOutboundGatewayParserTests { assertThatExceptionOfType(MessageHandlingException.class) .isThrownBy(() -> messagingTemplate.sendAndReceive(new GenericMessage<>("String content:"))) - .withMessageStartingWith("The destination file already exists at '"); + .withMessageContaining("The destination file already exists at '"); } /** diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java index d7c4db595c..66e47c5781 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java @@ -162,7 +162,7 @@ public class FileSplitterTests { } catch (Exception e) { assertThat(e.getCause()).isInstanceOf(FileNotFoundException.class); - assertThat(e.getMessage()).contains("failed to read file [bar]"); + assertThat(e.getMessage()).contains("Failed to read file [bar]"); } this.input2.send(new GenericMessage<>(new Date())); receive = this.output.receive(10000); diff --git a/spring-integration-ftp/src/main/java/org/springframework/integration/ftp/gateway/FtpOutboundGateway.java b/spring-integration-ftp/src/main/java/org/springframework/integration/ftp/gateway/FtpOutboundGateway.java index b4c50d89e3..eb59e4d811 100644 --- a/spring-integration-ftp/src/main/java/org/springframework/integration/ftp/gateway/FtpOutboundGateway.java +++ b/spring-integration-ftp/src/main/java/org/springframework/integration/ftp/gateway/FtpOutboundGateway.java @@ -18,6 +18,7 @@ package org.springframework.integration.ftp.gateway; import java.io.File; import java.io.IOException; +import java.io.UncheckedIOException; import java.util.ArrayList; import java.util.Collection; import java.util.List; @@ -39,7 +40,6 @@ import org.springframework.integration.file.remote.session.SessionFactory; import org.springframework.integration.ftp.session.FtpFileInfo; import org.springframework.integration.ftp.session.FtpRemoteFileTemplate; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandlingException; /** * Outbound Gateway for performing remote file operations via FTP/FTPS. @@ -63,6 +63,7 @@ public class FtpOutboundGateway extends AbstractRemoteFileOutboundGateway sessionFactory, MessageSessionCallback messageSessionCallback) { + this(new FtpRemoteFileTemplate(sessionFactory), messageSessionCallback); ((FtpRemoteFileTemplate) getRemoteFileTemplate()).setExistsMode(FtpRemoteFileTemplate.ExistsMode.NLST); } @@ -75,6 +76,7 @@ public class FtpOutboundGateway extends AbstractRemoteFileOutboundGateway remoteFileTemplate, MessageSessionCallback messageSessionCallback) { + super(remoteFileTemplate, messageSessionCallback); } @@ -179,7 +181,7 @@ public class FtpOutboundGateway extends AbstractRemoteFileOutboundGateway> asFileInfoList(Collection files) { - List> canonicalFiles = new ArrayList>(); + List> canonicalFiles = new ArrayList<>(); for (FTPFile file : files) { canonicalFiles.add(new FtpFileInfo(file)); } @@ -200,19 +202,18 @@ public class FtpOutboundGateway extends AbstractRemoteFileOutboundGateway ls(Message message, Session session, String dir) throws IOException { - return doInWorkingDirectory(message, session, - () -> super.ls(message, session, dir)); + return doInWorkingDirectory(message, session, () -> super.ls(message, session, dir)); } @Override protected List nlst(Message message, Session session, String dir) throws IOException { - return doInWorkingDirectory(message, session, - () -> super.nlst(message, session, dir)); + return doInWorkingDirectory(message, session, () -> super.nlst(message, session, dir)); } @Override protected File get(Message message, Session session, String remoteDir, String remoteFilePath, String remoteFilename, FTPFile fileInfoParam) throws IOException { + return doInWorkingDirectory(message, session, () -> super.get(message, session, remoteDir, remoteFilePath, remoteFilename, fileInfoParam)); } @@ -220,19 +221,20 @@ public class FtpOutboundGateway extends AbstractRemoteFileOutboundGateway mGet(Message message, Session session, String remoteDirectory, String remoteFilename) throws IOException { + return doInWorkingDirectory(message, session, () -> super.mGet(message, session, remoteDirectory, remoteFilename)); } @Override protected boolean rm(Message message, Session session, String remoteFilePath) throws IOException { - return doInWorkingDirectory(message, session, - () -> super.rm(message, session, remoteFilePath)); + return doInWorkingDirectory(message, session, () -> super.rm(message, session, remoteFilePath)); } @Override protected boolean mv(Message message, Session session, String remoteFilePath, String remoteFileNewPath) throws IOException { + return doInWorkingDirectory(message, session, () -> super.mv(message, session, remoteFilePath, remoteFileNewPath)); } @@ -244,7 +246,7 @@ public class FtpOutboundGateway extends AbstractRemoteFileOutboundGateway super.put(message, session, subDirectory)); } catch (IOException e) { - throw new MessageHandlingException(message, "Cannot handle PUT command", e); + throw new UncheckedIOException(e); } } @@ -255,12 +257,13 @@ public class FtpOutboundGateway extends AbstractRemoteFileOutboundGateway super.mPut(message, session, localDir)); } catch (IOException e) { - throw new MessageHandlingException(message, "Cannot handle MPUT command", e); + throw new UncheckedIOException(e); } } private V doInWorkingDirectory(Message message, Session session, Callable task) throws IOException { + Expression workDirExpression = this.workingDirExpression; FTPClient ftpClient = (FTPClient) session.getClientInstance(); String currentWorkingDirectory = null; diff --git a/spring-integration-groovy/src/test/java/org/springframework/integration/groovy/GroovyScriptPayloadMessageProcessorTests.java b/spring-integration-groovy/src/test/java/org/springframework/integration/groovy/GroovyScriptPayloadMessageProcessorTests.java index 7f5a2c8c37..7ad3d3c33e 100644 --- a/spring-integration-groovy/src/test/java/org/springframework/integration/groovy/GroovyScriptPayloadMessageProcessorTests.java +++ b/spring-integration-groovy/src/test/java/org/springframework/integration/groovy/GroovyScriptPayloadMessageProcessorTests.java @@ -17,7 +17,7 @@ package org.springframework.integration.groovy; import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.fail; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; import java.util.Collections; import java.util.concurrent.atomic.AtomicInteger; @@ -83,7 +83,7 @@ public class GroovyScriptPayloadMessageProcessorTests { assertThat(result.toString()).isEqualTo("spam is bucket foo is bar"); } - @Test //INT-2567 + @Test public void testBindingOverwrite() { Binding binding = new Binding() { @@ -94,16 +94,12 @@ public class GroovyScriptPayloadMessageProcessorTests { }; Message message = MessageBuilder.withPayload("foo").build(); processor = new GroovyCommandMessageProcessor(binding); - try { - processor.processMessage(message); - fail("Expected RuntimeException"); - } - catch (Exception e) { - assertThat(e.getCause().getMessage()).isEqualTo("intentional"); - } + assertThatExceptionOfType(RuntimeException.class) + .isThrownBy(() -> processor.processMessage(message)) + .withMessage("intentional"); } - @Test //INT-2567 + @Test public void testBindingOverwriteWithContext() { final String defaultValue = "default"; Binding binding = new Binding() { diff --git a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java index 9ad7a5bb4f..e619d978d9 100644 --- a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java +++ b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java @@ -19,6 +19,7 @@ package org.springframework.integration.http.outbound; import java.net.URI; import java.net.URISyntaxException; import java.nio.charset.Charset; +import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; @@ -92,25 +93,25 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac private boolean trustedSpel; - private volatile boolean encodeUri = true; + private boolean encodeUri = true; - private volatile Expression httpMethodExpression = new ValueExpression<>(HttpMethod.POST); + private Expression httpMethodExpression = new ValueExpression<>(HttpMethod.POST); - private volatile boolean expectReply = true; + private boolean expectReply = true; - private volatile Expression expectedResponseTypeExpression; + private Expression expectedResponseTypeExpression; - private volatile boolean extractPayload = true; + private boolean extractPayload = true; - private volatile boolean extractPayloadExplicitlySet = false; + private boolean extractPayloadExplicitlySet = false; - private volatile Charset charset = Charset.forName("UTF-8"); + private Charset charset = StandardCharsets.UTF_8; - private volatile boolean transferCookies = false; + private boolean transferCookies = false; - private volatile HeaderMapper headerMapper = DefaultHttpHeaderMapper.outboundMapper(); + private HeaderMapper headerMapper = DefaultHttpHeaderMapper.outboundMapper(); - private volatile Expression uriVariablesExpression; + private Expression uriVariablesExpression; public AbstractHttpRequestExecutingMessageHandler(Expression uriExpression) { Assert.notNull(uriExpression, "URI Expression is required"); @@ -166,7 +167,7 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac * @param charset The charset. */ public void setCharset(String charset) { - Assert.isTrue(Charset.isSupported(charset), "unsupported charset '" + charset + "'"); + Assert.isTrue(Charset.isSupported(charset), () -> "unsupported charset '" + charset + "'"); this.charset = Charset.forName(charset); } @@ -293,7 +294,7 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac private URI generateUri(Message requestMessage) { Object uri = this.uriExpression.getValue(this.evaluationContext, requestMessage); Assert.state(uri instanceof String || uri instanceof URI, - "'uriExpression' evaluation must result in a 'String' or 'URI' instance, not: " + () -> "'uriExpression' evaluation must result in a 'String' or 'URI' instance, not: " + (uri == null ? "null" : uri.getClass())); Map uriVariables = determineUriVariables(requestMessage); UriComponentsBuilder uriComponentsBuilder = @@ -305,7 +306,7 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac return this.encodeUri ? uriComponents.encode().toUri() : new URI(uriComponents.toUriString()); } catch (URISyntaxException e) { - throw new MessageHandlingException(requestMessage, "Invalid URI [" + uri + "]", e); + throw new MessageHandlingException(requestMessage, "Invalid URI [" + uri + "] in the [" + this + ']', e); } } @@ -351,8 +352,7 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac Object cookies = headers.remove(keyName); headers.put(HttpHeaders.COOKIE, cookies); if (logger.isDebugEnabled()) { - logger.debug("Converted Set-Cookie header to Cookie for: " - + cookies); + logger.debug("Converted Set-Cookie header to Cookie for: " + cookies); } } } diff --git a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java index 0be5d6f4d5..1704fba7ee 100755 --- a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java +++ b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java @@ -69,7 +69,7 @@ public class HttpRequestExecutingMessageHandler extends AbstractHttpRequestExecu * @param uri The URI. */ public HttpRequestExecutingMessageHandler(URI uri) { - this(new ValueExpression(uri)); + this(new ValueExpression<>(uri)); } /** @@ -159,12 +159,14 @@ public class HttpRequestExecutingMessageHandler extends AbstractHttpRequestExecu (ParameterizedTypeReference) expectedResponseType); } else { - httpResponse = this.restTemplate.exchange(uri, httpMethod, httpRequest, (Class) expectedResponseType); + httpResponse = this.restTemplate.exchange(uri, httpMethod, httpRequest, + (Class) expectedResponseType); } return getReply(httpResponse); } catch (RestClientException e) { - throw new MessageHandlingException(requestMessage, "HTTP request execution failed for URI [" + uri + "]", e); + throw new MessageHandlingException(requestMessage, + "HTTP request execution failed for URI [" + uri + "] in the [" + this + ']', e); } } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java index e43269dbf7..6b85bca77b 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java @@ -79,6 +79,7 @@ public class TcpOutboundGateway extends AbstractReplyProducingMessageHandler private Expression remoteTimeoutExpression = new ValueExpression<>(DEFAULT_REMOTE_TIMEOUT); private long requestTimeout = 10000; + private EvaluationContext evaluationContext = new StandardEvaluationContext(); private boolean evaluationContextSet; @@ -168,7 +169,7 @@ public class TcpOutboundGateway extends AbstractReplyProducingMessageHandler } catch (InterruptedException e) { Thread.currentThread().interrupt(); - throw new MessageHandlingException(requestMessage, "Interrupted", e); + throw new MessageHandlingException(requestMessage, "Interrupted in the [" + this + ']', e); } finally { cleanUp(haveSemaphore, connection, connectionId); @@ -232,9 +233,7 @@ public class TcpOutboundGateway extends AbstractReplyProducingMessageHandler } if (haveSemaphore) { this.semaphore.release(); - if (logger.isDebugEnabled()) { - logger.debug("released semaphore"); - } + logger.debug("released semaphore"); } } @@ -377,7 +376,6 @@ public class TcpOutboundGateway extends AbstractReplyProducingMessageHandler /** * Sender blocks here until the reply is received, or we time out * @return The return message or null if we time out - * @throws Exception */ public Message getReply() { try { @@ -397,7 +395,8 @@ public class TcpOutboundGateway extends AbstractReplyProducingMessageHandler */ logger.debug("second chance"); try { - this.secondChanceLatch.await(TcpOutboundGateway.this.secondChanceDelay, TimeUnit.SECONDS); // NOSONAR + this.secondChanceLatch + .await(TcpOutboundGateway.this.secondChanceDelay, TimeUnit.SECONDS); // NOSONAR } catch (@SuppressWarnings("unused") InterruptedException e) { Thread.currentThread().interrupt(); diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandler.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandler.java index 533d2046ed..de2c808623 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandler.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandler.java @@ -36,6 +36,7 @@ import org.springframework.integration.ip.tcp.connection.TcpSender; import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; +import org.springframework.scheduling.TaskScheduler; import org.springframework.util.Assert; /** @@ -86,7 +87,7 @@ public class TcpSendingMessageHandler extends AbstractMessageHandler implements } catch (Exception e) { logger.error("Error creating connection", e); - throw new MessageHandlingException(message, "Failed to obtain a connection", e); + throw new MessageHandlingException(message, "Failed to obtain a connection in the [" + this + ']', e); } return connection; } @@ -121,7 +122,7 @@ public class TcpSendingMessageHandler extends AbstractMessageHandler implements logger.error("Error sending message", ex); connection.close(); throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "Error sending message", ex); + () -> "Error sending message in the [" + this + ']', ex); } finally { if (this.isSingleUse) { // close after replying @@ -132,7 +133,7 @@ public class TcpSendingMessageHandler extends AbstractMessageHandler implements else { logger.error("Unable to find outbound socket for " + message); MessageHandlingException messageHandlingException = - new MessageHandlingException(message, "Unable to find outbound socket"); + new MessageHandlingException(message, "Unable to find outbound socket in the [" + this + ']'); publishNoConnectionEvent(messageHandlingException, connectionId); throw messageHandlingException; } @@ -190,7 +191,7 @@ public class TcpSendingMessageHandler extends AbstractMessageHandler implements } throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "Failed to handle message using " + connectionId, ex); + () -> "Failed to handle message in the [" + this + "] using " + connectionId, ex); } return connection; } @@ -265,8 +266,9 @@ public class TcpSendingMessageHandler extends AbstractMessageHandler implements ClientModeConnectionManager manager = new ClientModeConnectionManager(this.clientConnectionFactory); this.clientModeConnectionManager = manager; - Assert.state(getTaskScheduler() != null, "Client mode requires a task scheduler"); - this.scheduledFuture = getTaskScheduler().scheduleAtFixedRate(manager, this.retryInterval); + TaskScheduler taskScheduler = getTaskScheduler(); + Assert.state(taskScheduler != null, "Client mode requires a task scheduler"); + this.scheduledFuture = taskScheduler.scheduleAtFixedRate(manager, this.retryInterval); } } } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpMessageMapper.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpMessageMapper.java index afa303b208..1951a9f042 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpMessageMapper.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpMessageMapper.java @@ -16,6 +16,7 @@ package org.springframework.integration.ip.tcp.connection; +import java.io.UncheckedIOException; import java.io.UnsupportedEncodingException; import java.util.Map; @@ -37,7 +38,6 @@ import org.springframework.integration.support.MutableMessageHeaders; import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessageHeaders; import org.springframework.util.Assert; import org.springframework.util.InvalidMimeTypeException; @@ -268,13 +268,13 @@ public class TcpMessageMapper implements bytes = ((String) payload).getBytes(this.charset); } catch (UnsupportedEncodingException e) { - throw new MessageHandlingException(message, e); + throw new UncheckedIOException(e); } } else { - throw new MessageHandlingException(message, + throw new IllegalArgumentException( "When using a byte array serializer, the socket mapper expects " + - "either a byte array or String payload, but received: " + payload.getClass()); + "either a byte array or String payload, but received: " + payload.getClass()); } return bytes; } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java index d1cb1c572d..fa990e5e11 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/udp/DatagramPacketMessageMapper.java @@ -16,6 +16,7 @@ package org.springframework.integration.ip.udp; +import java.io.UncheckedIOException; import java.io.UnsupportedEncodingException; import java.net.DatagramPacket; import java.nio.ByteBuffer; @@ -37,7 +38,6 @@ import org.springframework.integration.support.MessageBuilderFactory; import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.MessagingException; import org.springframework.util.Assert; @@ -85,10 +85,10 @@ public class DatagramPacketMessageMapper implements InboundMessageMapper "Failed to send UDP packet", ex); + () -> "Failed to send UDP packet in the [" + this + ']', ex); } finally { if (countdownLatch != null) { diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/StoredProcOutboundGateway.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/StoredProcOutboundGateway.java index 9bc24509fc..7a36ff1513 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/StoredProcOutboundGateway.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/StoredProcOutboundGateway.java @@ -20,7 +20,6 @@ import java.util.Map; import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandlingException; import org.springframework.util.Assert; /** @@ -107,10 +106,8 @@ public class StoredProcOutboundGateway extends AbstractReplyProducingMessageHand payload = resultMap.values().iterator().next(); } else if (this.expectSingleResult && resultMap.size() > 1) { - - throw new MessageHandlingException(requestMessage, - "Stored Procedure/Function call returned more than " - + "1 result object and expectSingleResult was 'true'. "); + throw new IllegalStateException("Stored Procedure/Function call returned more than " + + "1 result object and expectSingleResult was 'true'."); } else { payload = resultMap; @@ -121,5 +118,4 @@ public class StoredProcOutboundGateway extends AbstractReplyProducingMessageHand } - } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java index 528c755b36..aad7f7d650 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java @@ -49,6 +49,7 @@ import org.springframework.core.convert.ConversionService; import org.springframework.expression.Expression; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.MessageTimeoutException; +import org.springframework.integration.StaticMessageHeaderAccessor; import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor; import org.springframework.integration.jms.util.JmsAdapterUtils; @@ -82,93 +83,94 @@ import org.springframework.util.concurrent.SettableListenableFuture; */ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler implements Lifecycle, MessageListener { + /** + * A default receive timeout in milliseconds. + */ + public static final long DEFAULT_RECEIVE_TIMEOUT = 5000L; + private final Object initializationMonitor = new Object(); private final AtomicLong correlationId = new AtomicLong(); private final String gatewayCorrelation = UUID.randomUUID().toString(); - private final Map> replies = - new ConcurrentHashMap>(); + private final Map> replies = new ConcurrentHashMap<>(); - private final ConcurrentHashMap earlyOrLateReplies = - new ConcurrentHashMap(); + private final ConcurrentHashMap earlyOrLateReplies = new ConcurrentHashMap<>(); private final Map>> futures = - new ConcurrentHashMap>>(); + new ConcurrentHashMap<>(); private final Object lifeCycleMonitor = new Object(); - private volatile Destination requestDestination; + private Destination requestDestination; - private volatile String requestDestinationName; + private String requestDestinationName; - private volatile ExpressionEvaluatingMessageProcessor requestDestinationExpressionProcessor; + private ExpressionEvaluatingMessageProcessor requestDestinationExpressionProcessor; - private volatile Destination replyDestination; + private Destination replyDestination; - private volatile String replyDestinationName; + private String replyDestinationName; - private volatile ExpressionEvaluatingMessageProcessor replyDestinationExpressionProcessor; + private ExpressionEvaluatingMessageProcessor replyDestinationExpressionProcessor; - private volatile DestinationResolver destinationResolver = new DynamicDestinationResolver(); + private DestinationResolver destinationResolver = new DynamicDestinationResolver(); - private volatile boolean requestPubSubDomain; + private boolean requestPubSubDomain; - private volatile boolean replyPubSubDomain; + private boolean replyPubSubDomain; - private volatile long receiveTimeout = 5000; + private long receiveTimeout = DEFAULT_RECEIVE_TIMEOUT; - private volatile int deliveryMode = javax.jms.Message.DEFAULT_DELIVERY_MODE; + private int deliveryMode = javax.jms.Message.DEFAULT_DELIVERY_MODE; - private volatile long timeToLive = javax.jms.Message.DEFAULT_TIME_TO_LIVE; + private long timeToLive = javax.jms.Message.DEFAULT_TIME_TO_LIVE; - private volatile int defaultPriority = javax.jms.Message.DEFAULT_PRIORITY; + private int defaultPriority = javax.jms.Message.DEFAULT_PRIORITY; - private volatile boolean explicitQosEnabled; + private boolean explicitQosEnabled; private ConnectionFactory connectionFactory; - private volatile MessageConverter messageConverter = new SimpleMessageConverter(); + private MessageConverter messageConverter = new SimpleMessageConverter(); - private volatile JmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper(); + private JmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper(); - private volatile String correlationKey; + private String correlationKey; - private volatile boolean extractRequestPayload = true; + private boolean extractRequestPayload = true; - private volatile boolean extractReplyPayload = true; + private boolean extractReplyPayload = true; - private volatile boolean initialized; + private GatewayReplyListenerContainer replyContainer; - private volatile GatewayReplyListenerContainer replyContainer; + private ReplyContainerProperties replyContainerProperties; - private volatile ReplyContainerProperties replyContainerProperties; + private boolean useReplyContainer; - private volatile boolean useReplyContainer; + private boolean requiresReply; + + private long idleReplyContainerTimeout; private volatile boolean active; + private volatile boolean initialized; + private volatile ScheduledFuture reaper; - private volatile boolean requiresReply; - - private long lastSend; - - private volatile long idleReplyContainerTimeout; - private volatile boolean wasStopped; - private ScheduledFuture idleTask; + private volatile ScheduledFuture idleTask; + + private volatile long lastSend; /** * Set whether message delivery should be persistent or non-persistent, * specified as a boolean value ("true" or "false"). This will set the delivery * mode accordingly to either "PERSISTENT" (1) or "NON_PERSISTENT" (2). *

The default is "true", i.e. delivery mode "PERSISTENT". - * * @param deliveryPersistent true for a persistent delivery. - * * @see javax.jms.DeliveryMode#PERSISTENT * @see javax.jms.DeliveryMode#NON_PERSISTENT */ @@ -179,7 +181,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Set the JMS ConnectionFactory that this gateway should use. * This is a required property. - * * @param connectionFactory The connection factory. */ public void setConnectionFactory(ConnectionFactory connectionFactory) { @@ -189,7 +190,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Set the JMS Destination to which request Messages should be sent. * Either this or one of 'requestDestinationName' or 'requestDestinationExpression' is required. - * * @param requestDestination The request destination. */ public void setRequestDestination(Destination requestDestination) { @@ -202,7 +202,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Set the name of the JMS Destination to which request Messages should be sent. * Either this or one of 'requestDestination' or 'requestDestinationExpression' is required. - * * @param requestDestinationName The request destination name. */ public void setRequestDestinationName(String requestDestinationName) { @@ -213,20 +212,18 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp * Set the SpEL Expression to be used for determining the request Destination instance * or request destination name. Either this or one of 'requestDestination' or * 'requestDestinationName' is required. - * * @param requestDestinationExpression The request destination expression. */ public void setRequestDestinationExpression(Expression requestDestinationExpression) { Assert.notNull(requestDestinationExpression, "'requestDestinationExpression' must not be null"); this.requestDestinationExpressionProcessor = - new ExpressionEvaluatingMessageProcessor(requestDestinationExpression); + new ExpressionEvaluatingMessageProcessor<>(requestDestinationExpression); setPrimaryExpression(requestDestinationExpression); } /** * Set the JMS Destination from which reply Messages should be received. * If none is provided, this gateway will create a {@link TemporaryQueue} per invocation. - * * @param replyDestination The reply destination. */ public void setReplyDestination(Destination replyDestination) { @@ -239,7 +236,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Set the name of the JMS Destination from which reply Messages should be received. * If none is provided, this gateway will create a {@link TemporaryQueue} per invocation. - * * @param replyDestinationName The reply destination name. */ public void setReplyDestinationName(String replyDestinationName) { @@ -250,19 +246,18 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp * Set the SpEL Expression to be used for determining the reply Destination instance * or reply destination name. Either this or one of 'replyDestination' or * 'replyDestinationName' is required. - * * @param replyDestinationExpression The reply destination expression. */ public void setReplyDestinationExpression(Expression replyDestinationExpression) { Assert.notNull(replyDestinationExpression, "'replyDestinationExpression' must not be null"); - this.replyDestinationExpressionProcessor = new ExpressionEvaluatingMessageProcessor(replyDestinationExpression); + this.replyDestinationExpressionProcessor = + new ExpressionEvaluatingMessageProcessor<>(replyDestinationExpression); } /** * Provide the {@link DestinationResolver} to use when resolving either a * 'requestDestinationName' or 'replyDestinationName' value. The default * is an instance of {@link DynamicDestinationResolver}. - * * @param destinationResolver The destination resolver. */ public void setDestinationResolver(DestinationResolver destinationResolver) { @@ -273,7 +268,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp * Specify whether the request destination is a Topic. This value is * necessary when providing a destination name for a Topic rather than * a destination reference. - * * @param requestPubSubDomain true if the request destination is a Topic. */ public void setRequestPubSubDomain(boolean requestPubSubDomain) { @@ -284,7 +278,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp * Specify whether the reply destination is a Topic. This value is * necessary when providing a destination name for a Topic rather than * a destination reference. - * * @param replyPubSubDomain true if the reply destination is a Topic. */ public void setReplyPubSubDomain(boolean replyPubSubDomain) { @@ -294,7 +287,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Set the max timeout value for the MessageConsumer's receive call when * waiting for a reply. The default value is 5 seconds. - * * @param receiveTimeout The receive timeout. */ public void setReceiveTimeout(long receiveTimeout) { @@ -304,23 +296,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Specify the default JMS priority to use when sending request Messages with * no {@link IntegrationMessageHeaderAccessor#PRIORITY} header. - * * The value should be within the range of 0-9. - * - * @param priority The priority. - * @deprecated in favor of {@link #setDefaultPriority(int)}. - */ - @Deprecated - public void setPriority(int priority) { - this.defaultPriority = priority; - } - - /** - * Specify the default JMS priority to use when sending request Messages with - * no {@link IntegrationMessageHeaderAccessor#PRIORITY} header. - * - * The value should be within the range of 0-9. - * * @param priority The priority. * @since 5.1.2 */ @@ -331,7 +307,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Specify the timeToLive for each sent Message. * The default value indicates no expiration. - * * @param timeToLive The time to live. */ public void setTimeToLive(long timeToLive) { @@ -341,7 +316,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Specify whether explicit QoS settings are enabled * (deliveryMode, priority, and timeToLive). - * * @param explicitQosEnabled true to enable explicit QoS. */ public void setExplicitQosEnabled(boolean explicitQosEnabled) { @@ -375,9 +349,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp * Provide a {@link MessageConverter} strategy to use for converting the * Spring Integration request Message into a JMS Message and for converting * the JMS reply Messages back into Spring Integration Messages. - *

- * The default is {@link SimpleMessageConverter}. - * + *

The default is {@link SimpleMessageConverter}. * @param messageConverter The message converter. */ public void setMessageConverter(MessageConverter messageConverter) { @@ -388,7 +360,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Provide a {@link JmsHeaderMapper} implementation for mapping the * Spring Integration Message Headers to/from JMS Message properties. - * * @param headerMapper The header mapper. */ public void setHeaderMapper(JmsHeaderMapper headerMapper) { @@ -403,9 +374,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp * the base for JMS Message creation. Since the JMS Message is created by the * MessageConverter, this really manages what is sent to the {@link MessageConverter}: * the entire Spring Integration Message or only its payload. - *
- * Default is 'true' - * + * Default is 'true'. * @param extractRequestPayload true to extract the request payload. */ public void setExtractRequestPayload(boolean extractRequestPayload) { @@ -418,7 +387,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp * created from the JMS Reply Message's body (via MessageConverter). * Otherwise, the entire JMS Message will become the payload of the * Spring Integration Message. - * * @param extractReplyPayload true to extract the reply payload. */ public void setExtractReplyPayload(boolean extractReplyPayload) { @@ -428,11 +396,10 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Specify the Spring Integration reply channel. If this property is not * set the gateway will check for a 'replyChannel' header on the request. - * * @param replyChannel The reply channel. */ public void setReplyChannel(MessageChannel replyChannel) { - this.setOutputChannel(replyChannel); + setOutputChannel(replyChannel); } /** @@ -489,7 +456,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp return this.requestDestination; } if (this.requestDestinationName != null) { - return this.resolveRequestDestination(this.requestDestinationName, session); + return resolveRequestDestination(this.requestDestinationName, session); } if (this.requestDestinationExpressionProcessor != null) { Object result = this.requestDestinationExpressionProcessor.processMessage(message); @@ -497,7 +464,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp return (Destination) result; } if (result instanceof String) { - return this.resolveRequestDestination((String) result, session); + return resolveRequestDestination((String) result, session); } throw new MessageDeliveryException(message, "Evaluation of requestDestinationExpression failed " + @@ -510,8 +477,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp private Destination resolveRequestDestination(String reqDestinationName, Session session) throws JMSException { Assert.notNull(this.destinationResolver, "DestinationResolver is required when relying upon the 'requestDestinationName' property."); - return this.destinationResolver.resolveDestinationName( - session, reqDestinationName, this.requestPubSubDomain); + return this.destinationResolver.resolveDestinationName(session, reqDestinationName, this.requestPubSubDomain); } private Destination determineReplyDestination(Message message, Session session) throws JMSException { @@ -519,7 +485,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp return this.replyDestination; } if (this.replyDestinationName != null) { - return this.resolveReplyDestination(this.replyDestinationName, session); + return resolveReplyDestination(this.replyDestinationName, session); } if (this.replyDestinationExpressionProcessor != null) { Object result = this.replyDestinationExpressionProcessor.processMessage(message); @@ -527,7 +493,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp return (Destination) result; } if (result instanceof String) { - return this.resolveReplyDestination((String) result, session); + return resolveReplyDestination((String) result, session); } throw new MessageDeliveryException(message, "Evaluation of replyDestinationExpression failed to produce a Destination or destination name. " + @@ -539,8 +505,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp private Destination resolveReplyDestination(String repDestinationName, Session session) throws JMSException { Assert.notNull(this.destinationResolver, "DestinationResolver is required when relying upon the 'replyDestinationName' property."); - return this.destinationResolver.resolveDestinationName( - session, repDestinationName, this.replyPubSubDomain); + return this.destinationResolver.resolveDestinationName(session, repDestinationName, this.replyPubSubDomain); } @Override @@ -575,17 +540,15 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp boolean hasAReplyDest = this.replyDestination != null || this.replyDestinationName != null || this.replyDestinationExpressionProcessor != null; if (this.useReplyContainer && (this.correlationKey == null && hasAReplyDest)) { - if (logger.isWarnEnabled()) { - logger.warn("The gateway cannot use a reply listener container with a specified " + - "destination(Name/Expression) " + - "without a 'correlation-key'; " + - "a container will NOT be used; " + - "to avoid this problem, set the 'correlation-key' attribute; " + - "some consumers, including the Spring Integration , " + - "support the use of the value 'JMSCorrelationID' " + - "for this purpose. Alternatively, do not specify a reply destination " + - "and a temporary queue will be used for replies."); - } + logger.warn("The gateway cannot use a reply listener container with a specified " + + "destination(Name/Expression) " + + "without a 'correlation-key'; " + + "a container will NOT be used; " + + "to avoid this problem, set the 'correlation-key' attribute; " + + "some consumers, including the Spring Integration , " + + "support the use of the value 'JMSCorrelationID' " + + "for this purpose. Alternatively, do not specify a reply destination " + + "and a temporary queue will be used for replies."); this.useReplyContainer = false; } if (this.useReplyContainer) { @@ -730,7 +693,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp if (this.replyContainer != null) { this.replyContainer.shutdown(); this.wasStopped = true; - this.deleteDestinationIfTemporary(this.replyContainer.getDestination()); + deleteDestinationIfTemporary(this.replyContainer.getDestination()); if (this.reaper != null) { this.reaper.cancel(false); } @@ -764,7 +727,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp this.lastSend = System.currentTimeMillis(); if (!this.replyContainer.isRunning()) { if (logger.isDebugEnabled()) { - logger.debug(this.getComponentName() + ": Starting reply container."); + logger.debug(getComponentName() + ": Starting reply container."); } this.replyContainer.start(); this.idleTask = getTaskScheduler().scheduleAtFixedRate(new IdleContainerStopper(), @@ -772,7 +735,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } } } - reply = this.sendAndReceiveWithContainer(requestMessage); + reply = sendAndReceiveWithContainer(requestMessage); } if (reply == null) { if (this.requiresReply) { @@ -792,7 +755,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } } catch (JMSException e) { - throw new MessageHandlingException(requestMessage, e); + throw new MessageHandlingException(requestMessage, "failed to handle a message in the [" + this + ']', e); } } @@ -801,7 +764,8 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp if (this.extractReplyPayload) { result = this.messageConverter.fromMessage(jmsReply); if (logger.isDebugEnabled()) { - logger.debug("converted JMS Message [" + jmsReply + "] to integration Message payload [" + result + "]"); + logger.debug("converted JMS Message [" + jmsReply + "] to integration Message payload [" + result + + "]"); } } Map jmsReplyHeaders = this.headerMapper.toHeaders(jmsReply); @@ -819,11 +783,11 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } private Object sendAndReceiveWithContainer(Message requestMessage) throws JMSException { - Connection connection = this.createConnection(); // NOSONAR - closed in ConnectionFactoryUtils. + Connection connection = createConnection(); // NOSONAR - closed in ConnectionFactoryUtils. Session session = null; Destination replyTo = this.replyContainer.getReplyDestination(); try { - session = this.createSession(connection); + session = createSession(connection); // convert to JMS Message Object objectToSend = requestMessage; @@ -841,13 +805,13 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp logger.debug("ReplyTo: " + replyTo); } - Integer priority = new IntegrationMessageHeaderAccessor(requestMessage).getPriority(); + Integer priority = StaticMessageHeaderAccessor.getPriority(requestMessage); if (priority == null) { priority = this.defaultPriority; } Destination destination = determineRequestDestination(requestMessage, session); - Object reply = null; + Object reply; if (this.correlationKey == null) { /* * Remove any existing correlation id that was mapped from the inbound message @@ -875,7 +839,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } private javax.jms.Message sendAndReceiveWithoutContainer(Message requestMessage) throws JMSException { - Connection connection = this.createConnection(); // NOSONAR - closed in ConnectionFactoryUtils. + Connection connection = createConnection(); // NOSONAR - closed in ConnectionFactoryUtils. Session session = null; Destination replyTo = null; try { @@ -891,19 +855,19 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp // map headers this.headerMapper.fromHeaders(requestMessage.getHeaders(), jmsRequest); - replyTo = this.determineReplyDestination(requestMessage, session); + replyTo = determineReplyDestination(requestMessage, session); jmsRequest.setJMSReplyTo(replyTo); connection.start(); if (logger.isDebugEnabled()) { logger.debug("ReplyTo: " + replyTo); } - Integer priority = new IntegrationMessageHeaderAccessor(requestMessage).getPriority(); + Integer priority = StaticMessageHeaderAccessor.getPriority(requestMessage); if (priority == null) { priority = this.defaultPriority; } - javax.jms.Message replyMessage = null; - Destination destination = this.determineRequestDestination(requestMessage, session); + javax.jms.Message replyMessage; + Destination destination = determineRequestDestination(requestMessage, session); if (this.correlationKey != null) { replyMessage = doSendAndReceiveWithGeneratedCorrelationId(destination, jmsRequest, replyTo, session, priority); @@ -935,7 +899,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp try { messageProducer = session.createProducer(reqDestination); Assert.state(this.correlationKey != null, "correlationKey must not be null"); - String messageSelector = null; + String messageSelector; if (!this.correlationKey.equals("JMSCorrelationID*") || jmsRequest.getJMSCorrelationID() == null) { String correlation = UUID.randomUUID().toString().replaceAll("'", "''"); if (this.correlationKey.equals("JMSCorrelationID")) { @@ -987,17 +951,27 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp private javax.jms.Message doSendAndReceiveWithMessageIdCorrelation(Destination reqDestination, javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority) throws JMSException { - if (replyTo instanceof Topic && logger.isWarnEnabled()) { - logger.warn("Relying on the MessageID for correlation is not recommended when using a Topic as the replyTo Destination " + - "because that ID can only be provided to a MessageSelector after the request Message has been sent thereby " + - "creating a race condition where a fast response might be sent before the MessageConsumer has been created. " + - "Consider providing a value to the 'correlationKey' property of this gateway instead. Then the MessageConsumer " + + if (replyTo instanceof Topic) { + logger.warn("Relying on the MessageID for correlation is not recommended when using a Topic as the " + + "replyTo" + + " " + + "Destination " + + "because that ID can only be provided to a MessageSelector after the request Message has been " + + "sent" + + " " + + "thereby " + + "creating a race condition where a fast response might be sent before the MessageConsumer has " + + "been" + + " " + + "created. " + + "Consider providing a value to the 'correlationKey' property of this gateway instead. Then the " + + "MessageConsumer " + "will be created before the request Message is sent."); } MessageProducer messageProducer = null; try { messageProducer = session.createProducer(reqDestination); - this.sendRequestMessage(jmsRequest, messageProducer, priority); + sendRequestMessage(jmsRequest, messageProducer, priority); String messageId = jmsRequest.getJMSMessageID().replaceAll("'", "''"); String messageSelector = "JMSCorrelationID = '" + messageId + "'"; return retryableReceiveReply(session, replyTo, messageSelector); @@ -1014,16 +988,17 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp */ private javax.jms.Message retryableReceiveReply(Session session, Destination replyTo, String messageSelector) throws JMSException { + Connection consumerConnection = null; //NOSONAR Session consumerSession = session; MessageConsumer messageConsumer = null; - JMSException exception = null; + JMSException exception; boolean isTemporaryReplyTo = replyTo instanceof TemporaryQueue || replyTo instanceof TemporaryTopic; long replyTimeout = isTemporaryReplyTo ? Long.MIN_VALUE : this.receiveTimeout < 0 - ? Long.MAX_VALUE - : System.currentTimeMillis() + this.receiveTimeout; + ? Long.MAX_VALUE + : System.currentTimeMillis() + this.receiveTimeout; try { do { try { @@ -1088,7 +1063,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp MessageProducer messageProducer = null; try { messageProducer = session.createProducer(reqDestination); - correlation = this.gatewayCorrelation + "_" + Long.toString(this.correlationId.incrementAndGet()); + correlation = this.gatewayCorrelation + "_" + this.correlationId.incrementAndGet(); if (this.correlationKey.equals("JMSCorrelationID")) { jmsRequest.setJMSCorrelationID(correlation); } @@ -1107,14 +1082,14 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp SettableListenableFuture> future = null; boolean async = isAsync(); if (!async) { - replyQueue = new LinkedBlockingQueue(1); + replyQueue = new LinkedBlockingQueue<>(1); this.replies.put(correlation, replyQueue); } else { future = createFuture(correlation); } - this.sendRequestMessage(jmsRequest, messageProducer, priority); + sendRequestMessage(jmsRequest, messageProducer, priority); if (async) { return future; @@ -1139,14 +1114,14 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp try { messageProducer = session.createProducer(reqDestination); - LinkedBlockingQueue replyQueue = new LinkedBlockingQueue(1); + LinkedBlockingQueue replyQueue = new LinkedBlockingQueue<>(1); this.sendRequestMessage(jmsRequest, messageProducer, priority); correlation = jmsRequest.getJMSMessageID(); if (logger.isDebugEnabled()) { - logger.debug(this.getComponentName() + " Sent message with correlationId " + correlation); + logger.debug(getComponentName() + " Sent message with correlationId " + correlation); } this.replies.put(correlation, replyQueue); @@ -1173,8 +1148,9 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } } - private javax.jms.Message obtainReplyFromContainer(String correlnId, + private javax.jms.Message obtainReplyFromContainer(String correlationId, LinkedBlockingQueue replyQueue) { + javax.jms.Message reply = null; if (this.receiveTimeout < 0) { @@ -1191,29 +1167,28 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } if (logger.isDebugEnabled()) { if (reply == null) { - logger.debug(this.getComponentName() + " Timed out waiting for reply with CorrelationId " - + correlnId); + logger.debug(getComponentName() + " Timed out waiting for reply with CorrelationId " + + correlationId); } else { - logger.debug(this.getComponentName() + " Obtained reply with CorrelationId " + correlnId); + logger.debug(getComponentName() + " Obtained reply with CorrelationId " + correlationId); } } return reply; } - private SettableListenableFuture> createFuture(final String correlnId) { - SettableListenableFuture> future = - new SettableListenableFuture>(); - this.futures.put(correlnId, future); + private SettableListenableFuture> createFuture(final String correlationId) { + SettableListenableFuture> future = new SettableListenableFuture<>(); + this.futures.put(correlationId, future); if (this.receiveTimeout > 0) { - getTaskScheduler().schedule((Runnable) () -> expire(correlnId), + getTaskScheduler().schedule(() -> expire(correlationId), new Date(System.currentTimeMillis() + this.receiveTimeout)); } return future; } - private void expire(String correlnId) { - final SettableListenableFuture> future = this.futures.remove(correlnId); + private void expire(String correlationId) { + SettableListenableFuture> future = this.futures.remove(correlationId); if (future != null) { try { if (getRequiresReply()) { @@ -1221,7 +1196,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } else { if (logger.isDebugEnabled()) { - logger.debug("Reply expired and reply not required for " + correlnId); + logger.debug("Reply expired and reply not required for " + correlationId); } } } @@ -1233,6 +1208,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp private void sendRequestMessage(javax.jms.Message jmsRequest, MessageProducer messageProducer, int priority) throws JMSException { + if (this.explicitQosEnabled) { messageProducer.send(jmsRequest, this.deliveryMode, priority, this.timeToLive); } @@ -1266,7 +1242,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Create a new JMS Connection for this JMS gateway. - * * @return The connection. * @throws JMSException Any JMSException. */ @@ -1276,7 +1251,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp /** * Create a new JMS Session using the provided Connection. - * * @param connection The connection. * @return The session. * @throws JMSException Any JMSException. @@ -1290,7 +1264,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp String correlation = null; try { if (logger.isTraceEnabled()) { - logger.trace(this.getComponentName() + " Received " + message); + logger.trace(getComponentName() + " Received " + message); } if (this.correlationKey == null || this.correlationKey.equals("JMSCorrelationID") || @@ -1315,50 +1289,50 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp } } - private void onMessageAsync(javax.jms.Message message, String correlnId) throws JMSException { - SettableListenableFuture> future = this.futures.remove(correlnId); + private void onMessageAsync(javax.jms.Message message, String correlationId) throws JMSException { + SettableListenableFuture> future = this.futures.remove(correlationId); if (future != null) { message.setJMSCorrelationID(null); future.set(buildReply(message)); } else { - logger.warn("Late reply for " + correlnId); + logger.warn("Late reply for " + correlationId); } } - private void onMessageSync(javax.jms.Message message, String correlnId) { + private void onMessageSync(javax.jms.Message message, String correlationId) { try { - LinkedBlockingQueue queue = this.replies.get(correlnId); + LinkedBlockingQueue queue = this.replies.get(correlationId); if (queue == null) { if (this.correlationKey != null) { Log debugLogger = LogFactory.getLog("si.jmsgateway.debug"); if (debugLogger.isDebugEnabled()) { Object siMessage = this.messageConverter.fromMessage(message); debugLogger.debug("No pending reply for " + siMessage + " with correlationId: " - + correlnId + " pending replies: " + this.replies.keySet()); + + correlationId + " pending replies: " + this.replies.keySet()); } throw new IllegalStateException("No sender waiting for reply"); } synchronized (this.earlyOrLateReplies) { - queue = this.replies.get(correlnId); + queue = this.replies.get(correlationId); if (queue == null) { if (logger.isDebugEnabled()) { - logger.debug("Reply for correlationId " + correlnId + " received early or late"); + logger.debug("Reply for correlationId " + correlationId + " received early or late"); } - this.earlyOrLateReplies.put(correlnId, new TimedReply(message)); + this.earlyOrLateReplies.put(correlationId, new TimedReply(message)); } } } if (queue != null) { if (logger.isDebugEnabled()) { - logger.debug("Received reply with correlationId " + correlnId); + logger.debug("Received reply with correlationId " + correlationId); } queue.add(message); } } catch (Exception e) { if (logger.isWarnEnabled()) { - logger.warn("Failed to consume reply with correlationId " + correlnId, e); + logger.warn("Failed to consume reply with correlationId " + correlationId, e); } } } @@ -1453,6 +1427,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp this.replyDestination = null; super.recoverAfterListenerSetupFailure(); } + } private static final class TimedReply { @@ -1472,6 +1447,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp private javax.jms.Message getReply() { return this.reply; } + } private class LateReplyReaper implements Runnable { @@ -1535,29 +1511,29 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler imp public static class ReplyContainerProperties { - private volatile Boolean sessionTransacted; + private Boolean sessionTransacted; - private volatile Integer sessionAcknowledgeMode; + private Integer sessionAcknowledgeMode; - private volatile String sessionAcknowledgeModeName; + private String sessionAcknowledgeModeName; - private volatile Long receiveTimeout; + private Long receiveTimeout; - private volatile Long recoveryInterval; + private Long recoveryInterval; - private volatile Integer cacheLevel; + private Integer cacheLevel; - private volatile Integer concurrentConsumers; + private Integer concurrentConsumers; - private volatile Integer maxConcurrentConsumers; + private Integer maxConcurrentConsumers; - private volatile Integer maxMessagesPerTask; + private Integer maxMessagesPerTask; - private volatile Integer idleConsumerLimit; + private Integer idleConsumerLimit; - private volatile Integer idleTaskExecutionLimit; + private Integer idleTaskExecutionLimit; - private volatile Executor taskExecutor; + private Executor taskExecutor; public String getSessionAcknowledgeModeName() { return this.sessionAcknowledgeModeName; diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundGatewayParser.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundGatewayParser.java index 9b11da6843..cbd6b6b1a2 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundGatewayParser.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsOutboundGatewayParser.java @@ -65,7 +65,7 @@ public class JmsOutboundGatewayParser extends AbstractConsumerEndpointParser { IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "request-pub-sub-domain"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-pub-sub-domain"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "time-to-live"); - IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "priority"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "priority", "defaultPriority"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "explicit-qos-enabled"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "requires-reply"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "idle-reply-listener-timeout", diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundGatewaySpec.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundGatewaySpec.java index 42773cdd21..452c869a41 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundGatewaySpec.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/dsl/JmsOutboundGatewaySpec.java @@ -240,7 +240,7 @@ public class JmsOutboundGatewaySpec extends MessageHandlerSpec message) { - try { - this.delegate.publish(this.notificationMapper.fromMessage(message)); - } - catch (Exception e) { - throw new MessageHandlingException(message, "Failed to handle", e); - } + this.delegate.publish(this.notificationMapper.fromMessage(message)); } diff --git a/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/OperationInvokingMessageHandler.java b/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/OperationInvokingMessageHandler.java index 1230a995de..f7a95eaf8f 100644 --- a/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/OperationInvokingMessageHandler.java +++ b/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/OperationInvokingMessageHandler.java @@ -17,6 +17,7 @@ package org.springframework.integration.jmx; import java.io.IOException; +import java.io.UncheckedIOException; import java.util.Arrays; import java.util.Collections; import java.util.HashMap; @@ -141,13 +142,13 @@ public class OperationInvokingMessageHandler extends AbstractReplyProducingMessa } return result; } - catch (JMException e) { + catch (JMException ex) { throw new MessageHandlingException(requestMessage, "failed to invoke JMX operation '" + operation + "' on MBean [" + objectName + "]" + " with " + - paramsFromMessage.size() + " parameters: " + paramsFromMessage, e); + paramsFromMessage.size() + " parameters [" + paramsFromMessage + "] in the [" + this + ']', ex); } - catch (IOException e) { - throw new MessageHandlingException(requestMessage, "IOException on MBeanServerConnection", e); + catch (IOException ex) { + throw new UncheckedIOException(ex); } } diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/MailSendingMessageHandler.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/MailSendingMessageHandler.java index ec059703be..457d636853 100644 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/MailSendingMessageHandler.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/MailSendingMessageHandler.java @@ -29,7 +29,6 @@ import org.springframework.mail.javamail.JavaMailSender; import org.springframework.mail.javamail.MimeMailMessage; import org.springframework.mail.javamail.MimeMessageHelper; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessageHeaders; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; @@ -91,7 +90,7 @@ public class MailSendingMessageHandler extends AbstractMessageHandler { @SuppressWarnings("unchecked") private MailMessage convertMessageToMailMessage(Message message) { - MailMessage mailMessage = null; + MailMessage mailMessage; Object payload = message.getPayload(); if (payload instanceof MimeMessage) { mailMessage = new MimeMailMessage((MimeMessage) payload); @@ -100,12 +99,12 @@ public class MailSendingMessageHandler extends AbstractMessageHandler { mailMessage = (MailMessage) payload; } else if (payload instanceof byte[]) { - mailMessage = this.createMailMessageFromByteArrayMessage((Message) message); + mailMessage = createMailMessageFromByteArrayMessage((Message) message); } else if (payload instanceof String) { String contentType = (String) message.getHeaders().get(MailHeaders.CONTENT_TYPE); if (StringUtils.hasText(contentType)) { - mailMessage = this.createMailMessageWithContentType((Message) message, contentType); + mailMessage = createMailMessageWithContentType((Message) message, contentType); } else { mailMessage = new SimpleMailMessage(); @@ -113,11 +112,11 @@ public class MailSendingMessageHandler extends AbstractMessageHandler { } } else { - throw new MessageHandlingException(message, "Unable to create MailMessage from payload type [" + throw new IllegalArgumentException("Unable to create MailMessage from payload type [" + message.getPayload().getClass().getName() + "], " + "expected MimeMessage, MailMessage, byte array or String."); } - this.applyHeadersToMailMessage(mailMessage, message.getHeaders()); + applyHeadersToMailMessage(mailMessage, message.getHeaders()); return mailMessage; } @@ -131,8 +130,7 @@ public class MailSendingMessageHandler extends AbstractMessageHandler { return new MimeMailMessage(mimeMessage); } catch (Exception e) { - throw new org.springframework.messaging.MessagingException("Failed to create MimeMessage with contentType: " - + contentType, e); + throw new IllegalStateException("Failed to create MimeMessage with contentType: " + contentType, e); } } @@ -166,7 +164,7 @@ public class MailSendingMessageHandler extends AbstractMessageHandler { if (subject != null) { mailMessage.setSubject(subject); } - String[] to = this.retrieveHeaderValueAsStringArray(headers, MailHeaders.TO); + String[] to = retrieveHeaderValueAsStringArray(headers, MailHeaders.TO); if (to != null) { mailMessage.setTo(to); } @@ -174,11 +172,11 @@ public class MailSendingMessageHandler extends AbstractMessageHandler { Assert.state(!ObjectUtils.isEmpty(((SimpleMailMessage) mailMessage).getTo()), "No recipient has been provided on the MailMessage or the 'MailHeaders.TO' header."); } - String[] cc = this.retrieveHeaderValueAsStringArray(headers, MailHeaders.CC); + String[] cc = retrieveHeaderValueAsStringArray(headers, MailHeaders.CC); if (cc != null) { mailMessage.setCc(cc); } - String[] bcc = this.retrieveHeaderValueAsStringArray(headers, MailHeaders.BCC); + String[] bcc = retrieveHeaderValueAsStringArray(headers, MailHeaders.BCC); if (bcc != null) { mailMessage.setBcc(bcc); } diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java index c5834cef22..d50a79f481 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java @@ -28,7 +28,6 @@ import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter; import org.springframework.integration.mqtt.support.MqttHeaders; import org.springframework.integration.mqtt.support.MqttMessageConverter; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.converter.MessageConverter; import org.springframework.util.Assert; @@ -44,7 +43,7 @@ import org.springframework.util.Assert; public abstract class AbstractMqttMessageHandler extends AbstractMessageHandler implements Lifecycle { private static final MessageProcessor DEFAULT_TOPIC_PROCESSOR = - m -> m.getHeaders().get(MqttHeaders.TOPIC, String.class); + (message) -> message.getHeaders().get(MqttHeaders.TOPIC, String.class); private final AtomicBoolean running = new AtomicBoolean(); @@ -262,7 +261,7 @@ public abstract class AbstractMqttMessageHandler extends AbstractMessageHandler Object mqttMessage = this.converter.fromMessage(message, Object.class); String topic = this.topicProcessor.processMessage(message); if (topic == null && this.defaultTopic == null) { - throw new MessageHandlingException(message, + throw new IllegalStateException( "No topic could be determined from the message and no default topic defined"); } publish(topic == null ? this.defaultTopic : topic, mqttMessage, message); diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java index 7af216bd17..0d2690e99a 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java @@ -47,6 +47,9 @@ import org.springframework.util.Assert; public class MqttPahoMessageHandler extends AbstractMqttMessageHandler implements MqttCallback, ApplicationEventPublisherAware { + /** + * The default completion timeout in milliseconds. + */ public static final long DEFAULT_COMPLETION_TIMEOUT = 30000L; private long completionTimeout = DEFAULT_COMPLETION_TIMEOUT; @@ -174,9 +177,7 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler incrementClientInstance(); this.client.setCallback(this); this.client.connect(connectionOptions).waitForCompletion(this.completionTimeout); - if (logger.isDebugEnabled()) { - logger.debug("Client connected"); - } + logger.debug("Client connected"); } catch (MqttException e) { if (this.client != null) { @@ -191,7 +192,7 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler @Override protected void publish(String topic, Object mqttMessage, Message message) { - Assert.isInstanceOf(MqttMessage.class, mqttMessage); + Assert.isInstanceOf(MqttMessage.class, mqttMessage, "The 'mqttMessage' must be an instance of 'MqttMessage'"); try { IMqttDeliveryToken token = checkConnection() .publish(topic, (MqttMessage) mqttMessage); @@ -205,7 +206,7 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler } } catch (MqttException e) { - throw new MessageHandlingException(message, "Failed to publish to MQTT", e); + throw new MessageHandlingException(message, "Failed to publish to MQTT in the [" + this + ']', e); } } diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java index b8625ebbdf..91cb51d255 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java @@ -272,8 +272,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa return this.bytesMessageMapper.fromMessage(message); } catch (Exception e) { - throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "Failed to map outbound message", e); + throw new IllegalStateException("Failed to map outbound message", e); } } else { diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisStoreWritingMessageHandler.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisStoreWritingMessageHandler.java index 60c77de45d..40a26391e4 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisStoreWritingMessageHandler.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisStoreWritingMessageHandler.java @@ -316,7 +316,7 @@ public class RedisStoreWritingMessageHandler extends AbstractMessageHandler { } catch (Exception ex) { throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "Failed to store Message data in Redis collection", ex); + () -> "Failed to store Message data into Redis collection in the [" + this + ']', ex); } } diff --git a/spring-integration-rmi/src/main/java/org/springframework/integration/rmi/RmiOutboundGateway.java b/spring-integration-rmi/src/main/java/org/springframework/integration/rmi/RmiOutboundGateway.java index dac256d983..5776160498 100644 --- a/spring-integration-rmi/src/main/java/org/springframework/integration/rmi/RmiOutboundGateway.java +++ b/spring-integration-rmi/src/main/java/org/springframework/integration/rmi/RmiOutboundGateway.java @@ -23,8 +23,6 @@ import org.springframework.integration.handler.AbstractReplyProducingMessageHand import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandlingException; -import org.springframework.messaging.MessagingException; -import org.springframework.remoting.RemoteAccessException; import org.springframework.remoting.rmi.RmiProxyFactoryBean; /** @@ -75,19 +73,15 @@ public class RmiOutboundGateway extends AbstractReplyProducingMessageHandler { @Override public final Object handleRequestMessage(Message requestMessage) { if (!(requestMessage.getPayload() instanceof Serializable)) { - throw new MessageHandlingException(requestMessage, - this.getComponentName() + " expects a Serializable payload type " + + throw new IllegalStateException( + "Expected a Serializable payload type " + "but encountered [" + requestMessage.getPayload().getClass().getName() + "]"); } try { return this.proxy.exchange(requestMessage); } - catch (MessagingException e) { - throw new MessageHandlingException(requestMessage, e); - } - catch (RemoteAccessException e) { - throw new MessageHandlingException(requestMessage, - "Remote failure in RmiOutboundGateway: " + getComponentName(), e); + catch (Exception ex) { + throw new MessageHandlingException(requestMessage, "Remote failure in the [" + this + ']', ex); } } diff --git a/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/AbstractScriptExecutingMessageProcessor.java b/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/AbstractScriptExecutingMessageProcessor.java index 222029db38..0f2503efc2 100644 --- a/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/AbstractScriptExecutingMessageProcessor.java +++ b/spring-integration-scripting/src/main/java/org/springframework/integration/scripting/AbstractScriptExecutingMessageProcessor.java @@ -23,7 +23,6 @@ import org.springframework.beans.factory.BeanClassLoaderAware; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.integration.handler.MessageProcessor; -import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.scripting.ScriptSource; @@ -85,14 +84,9 @@ public abstract class AbstractScriptExecutingMessageProcessor @Override @Nullable public final T processMessage(Message message) { - try { - ScriptSource source = getScriptSource(message); - Map variables = this.scriptVariableGenerator.generateScriptVariables(message); - return executeScript(source, variables); - } - catch (Exception e) { - throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, () -> "Failed to execute script.", e); - } + ScriptSource source = getScriptSource(message); + Map variables = this.scriptVariableGenerator.generateScriptVariables(message); + return executeScript(source, variables); } /** diff --git a/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapter.java b/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapter.java index b40ed4fb7e..5317597910 100644 --- a/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapter.java +++ b/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/inbound/StompInboundChannelAdapter.java @@ -282,7 +282,7 @@ public class StompInboundChannelAdapter extends MessageProducerSupport implement public void handleException(StompSession session, StompCommand command, StompHeaders headers, byte[] payload, Throwable exception) { - String exceptionMessage = "STOMP Frame handling error."; + String exceptionMessage = "STOMP Frame handling error in the [" + StompInboundChannelAdapter.this + ']'; MessageChannel errorChannel = getErrorChannel(); diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java index 8aeeed7674..e02d99a871 100644 --- a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java @@ -116,6 +116,7 @@ public class WebSocketInboundChannelAdapter extends MessageProducerSupport public WebSocketInboundChannelAdapter(IntegrationWebSocketContainer webSocketContainer, SubProtocolHandlerRegistry protocolHandlerRegistry) { + Assert.notNull(webSocketContainer, "'webSocketContainer' must not be null"); Assert.notNull(protocolHandlerRegistry, "'protocolHandlerRegistry' must not be null"); this.webSocketContainer = webSocketContainer; @@ -128,7 +129,8 @@ public class WebSocketInboundChannelAdapter extends MessageProducerSupport } catch (Exception e) { throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "Failed to handle and process message.", e); + () -> "Failed to handle and process message in the [" + + WebSocketInboundChannelAdapter.this + ']', e); } }); } @@ -241,6 +243,7 @@ public class WebSocketInboundChannelAdapter extends MessageProducerSupport @Override public void afterSessionEnded(WebSocketSession session, CloseStatus closeStatus) throws Exception { // NOSONAR Thrown from the delegate + if (isActive()) { this.subProtocolHandlerRegistry.findProtocolHandler(session) .afterSessionEnded(session, closeStatus, this.subProtocolHandlerChannel); @@ -250,6 +253,7 @@ public class WebSocketInboundChannelAdapter extends MessageProducerSupport @Override public void onMessage(WebSocketSession session, WebSocketMessage webSocketMessage) throws Exception { // NOSONAR Thrown from the delegate + if (isActive()) { this.subProtocolHandlerRegistry.findProtocolHandler(session) .handleMessageFromClient(session, webSocketMessage, this.subProtocolHandlerChannel); @@ -350,7 +354,7 @@ public class WebSocketInboundChannelAdapter extends MessageProducerSupport } catch (Exception e) { throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "Error sending connect ack message", e); + () -> "Error sending connect ack message in the [" + this + ']', e); } } diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/outbound/WebSocketOutboundMessageHandler.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/outbound/WebSocketOutboundMessageHandler.java index 9abdc45c86..cd6d9091f3 100644 --- a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/outbound/WebSocketOutboundMessageHandler.java +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/outbound/WebSocketOutboundMessageHandler.java @@ -71,9 +71,9 @@ public class WebSocketOutboundMessageHandler extends AbstractMessageHandler { private final boolean client; - private volatile List messageConverters; + private List messageConverters; - private volatile boolean mergeWithDefaultConverters = false; + private boolean mergeWithDefaultConverters = false; public WebSocketOutboundMessageHandler(IntegrationWebSocketContainer webSocketContainer) { this(webSocketContainer, new SubProtocolHandlerRegistry(new PassThruSubProtocolHandler())); @@ -168,7 +168,7 @@ public class WebSocketOutboundMessageHandler extends AbstractMessageHandler { } } catch (Exception e) { - throw new MessageHandlingException(message, "Failed to handle", e); + throw new MessageHandlingException(message, "Failed to handle message in the [" + this + ']', e); } } diff --git a/spring-integration-xml/src/main/java/org/springframework/integration/xml/selector/XmlValidatingMessageSelector.java b/spring-integration-xml/src/main/java/org/springframework/integration/xml/selector/XmlValidatingMessageSelector.java index 7febfb075d..f6e216e4ad 100644 --- a/spring-integration-xml/src/main/java/org/springframework/integration/xml/selector/XmlValidatingMessageSelector.java +++ b/spring-integration-xml/src/main/java/org/springframework/integration/xml/selector/XmlValidatingMessageSelector.java @@ -17,6 +17,7 @@ package org.springframework.integration.xml.selector; import java.io.IOException; +import java.io.UncheckedIOException; import java.util.Arrays; import org.apache.commons.logging.Log; @@ -30,7 +31,6 @@ import org.springframework.integration.xml.AggregatedXmlMessageValidationExcepti import org.springframework.integration.xml.DefaultXmlPayloadConverter; import org.springframework.integration.xml.XmlPayloadConverter; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHandlingException; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; @@ -65,7 +65,7 @@ public class XmlValidatingMessageSelector implements MessageSelector { } - private final Log logger = LogFactory.getLog(getClass()); + private static final Log LOGGER = LogFactory.getLog(XmlValidatingMessageSelector.class); private final XmlValidator xmlValidator; @@ -110,7 +110,6 @@ public class XmlValidatingMessageSelector implements MessageSelector { /** * Specify the Converter to use when converting payloads prior to validation. - * * @param converter The payload converter. */ public void setConverter(XmlPayloadConverter converter) { @@ -120,21 +119,22 @@ public class XmlValidatingMessageSelector implements MessageSelector { @Override public boolean accept(Message message) { - SAXParseException[] validationExceptions = null; + SAXParseException[] validationExceptions; try { validationExceptions = this.xmlValidator.validate(this.converter.convertToSource(message.getPayload())); } - catch (Exception e) { - throw new MessageHandlingException(message, e); + catch (IOException e) { + throw new UncheckedIOException(e); } boolean validationSuccess = ObjectUtils.isEmpty(validationExceptions); if (!validationSuccess) { + String exceptionMessage = "Message was rejected due to XML Validation errors"; if (this.throwExceptionOnRejection) { - throw new MessageRejectedException(message, "Message was rejected due to XML Validation errors", + throw new MessageRejectedException(message, exceptionMessage, new AggregatedXmlMessageValidationException(Arrays.asList(validationExceptions))); } - else if (this.logger.isInfoEnabled()) { - this.logger.info("Message was rejected due to XML Validation errors", + else if (LOGGER.isInfoEnabled()) { + LOGGER.info(exceptionMessage, new AggregatedXmlMessageValidationException(Arrays.asList(validationExceptions))); } } diff --git a/spring-integration-xml/src/main/java/org/springframework/integration/xml/splitter/XPathMessageSplitter.java b/spring-integration-xml/src/main/java/org/springframework/integration/xml/splitter/XPathMessageSplitter.java index 4c0f770b10..eca16e93ab 100644 --- a/spring-integration-xml/src/main/java/org/springframework/integration/xml/splitter/XPathMessageSplitter.java +++ b/spring-integration-xml/src/main/java/org/springframework/integration/xml/splitter/XPathMessageSplitter.java @@ -210,13 +210,13 @@ public class XPathMessageSplitter extends AbstractMessageSplitter { protected Object splitMessage(Message message) { try { Object payload = message.getPayload(); - Object result = null; + Object result; if (payload instanceof Node) { result = splitNode((Node) payload); } else { Document document = this.xmlPayloadConverter.convertToDocument(payload); - Assert.notNull(document, "unsupported payload type [" + payload.getClass().getName() + "]"); + Assert.notNull(document, () -> "unsupported payload type [" + payload.getClass().getName() + "]"); result = splitDocument(document); } return result; @@ -226,7 +226,7 @@ public class XPathMessageSplitter extends AbstractMessageSplitter { } catch (Exception ex) { throw IntegrationUtils.wrapInHandlingExceptionIfNecessary(message, - () -> "Failed to split Message payload", ex); + () -> "Failed to split Message payload in the [" + this + ']', ex); } } @@ -369,9 +369,7 @@ public class XPathMessageSplitter extends AbstractMessageSplitter { private final Iterator delegate; - TransformFunctionIterator(Iterator delegate, - Function function) { - + TransformFunctionIterator(Iterator delegate, Function function) { super(null, delegate, function); this.delegate = delegate; } diff --git a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/ChatMessageSendingMessageHandler.java b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/ChatMessageSendingMessageHandler.java index a78d4c544c..a06294967e 100644 --- a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/ChatMessageSendingMessageHandler.java +++ b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/ChatMessageSendingMessageHandler.java @@ -44,13 +44,14 @@ import org.springframework.util.StringUtils; * @author Mario Gray * @author Oleg Zhurakousky * @author Artem Bilan + * * @since 2.0 */ public class ChatMessageSendingMessageHandler extends AbstractXmppConnectionAwareMessageHandler { private static final Pattern XML_PATTERN = Pattern.compile("<(\\S[^>\\s]*)[^>]*>[^<]*"); - private volatile XmppHeaderMapper headerMapper = new DefaultXmppHeaderMapper(); + private XmppHeaderMapper headerMapper = new DefaultXmppHeaderMapper(); private ExtensionElementProvider extensionProvider; @@ -85,7 +86,8 @@ public class ChatMessageSendingMessageHandler extends AbstractXmppConnectionAwar @Override protected void handleMessageInternal(Message message) { - Assert.isTrue(isInitialized(), getComponentName() + "#" + this.getComponentType() + " must be initialized"); + Assert.isTrue(isInitialized(), + () -> getComponentName() + "#" + this.getComponentType() + " must be initialized"); try { Object payload = message.getPayload(); org.jivesoftware.smack.packet.Message xmppMessage = null; @@ -94,7 +96,7 @@ public class ChatMessageSendingMessageHandler extends AbstractXmppConnectionAwar } else { String to = message.getHeaders().get(XmppHeaders.TO, String.class); - Assert.state(StringUtils.hasText(to), "The '" + XmppHeaders.TO + "' header must not be null"); + Assert.state(StringUtils.hasText(to), () -> "The '" + XmppHeaders.TO + "' header must not be null"); xmppMessage = buildXmppMessage(message, payload, to); } @@ -109,10 +111,10 @@ public class ChatMessageSendingMessageHandler extends AbstractXmppConnectionAwar } catch (InterruptedException e) { Thread.currentThread().interrupt(); - throw new MessageHandlingException(message, "Interrupted", e); + throw new MessageHandlingException(message, "Thread interrupted in the [" + this + ']', e); } catch (Exception e) { - throw new MessageHandlingException(message, "Failed to handle", e); + throw new MessageHandlingException(message, "Failed to handle message in the [" + this + ']', e); } } @@ -144,7 +146,7 @@ public class ChatMessageSendingMessageHandler extends AbstractXmppConnectionAwar } } else { - throw new MessageHandlingException(message, + throw new IllegalStateException( "Only payloads of type java.lang.String, org.jivesoftware.smack.packet.Message " + "or org.jivesoftware.smack.packet.ExtensionElement " + "are supported. Received [" + payload.getClass().getName() + diff --git a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/PresenceSendingMessageHandler.java b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/PresenceSendingMessageHandler.java index af2bd246c2..386958d07a 100644 --- a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/PresenceSendingMessageHandler.java +++ b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/outbound/PresenceSendingMessageHandler.java @@ -31,9 +31,10 @@ import org.springframework.util.Assert; * @author Josh Long * @author Oleg Zhurakousky * @author Artem Bilan + * * @since 2.0 */ -public class PresenceSendingMessageHandler extends AbstractXmppConnectionAwareMessageHandler { +public class PresenceSendingMessageHandler extends AbstractXmppConnectionAwareMessageHandler { public PresenceSendingMessageHandler() { super(); @@ -50,10 +51,11 @@ public class PresenceSendingMessageHandler extends AbstractXmppConnectionAwareMe @Override protected void handleMessageInternal(Message message) { - Assert.state(isInitialized(), getComponentName() + " must be initialized"); + Assert.state(isInitialized(), () -> getComponentName() + " must be initialized"); Object payload = message.getPayload(); Assert.state(payload instanceof Presence, - "Payload must be of type 'org.jivesoftware.smack.packet.Presence', was: " + payload.getClass().getName()); + () -> "Payload must be of type 'org.jivesoftware.smack.packet.Presence', was: " + + payload.getClass().getName()); try { XMPPConnection xmppConnection = getXmppConnection(); if (!xmppConnection.isConnected() && xmppConnection instanceof AbstractXMPPConnection) { @@ -63,10 +65,10 @@ public class PresenceSendingMessageHandler extends AbstractXmppConnectionAwareMe } catch (InterruptedException e) { Thread.currentThread().interrupt(); - throw new MessageHandlingException(message, "Interrupted", e); + throw new MessageHandlingException(message, "Thread interrupted in the [" + this + ']', e); } catch (Exception e) { - throw new MessageHandlingException(message, "Failed to handle", e); + throw new MessageHandlingException(message, "Failed to handle message in the [" + this + ']', e); } } diff --git a/src/reference/asciidoc/configuration.adoc b/src/reference/asciidoc/configuration.adoc index d71ee8b66e..70dc26c1a0 100644 --- a/src/reference/asciidoc/configuration.adoc +++ b/src/reference/asciidoc/configuration.adoc @@ -217,6 +217,13 @@ The `requestMessage` is stored under `ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KE The `ErrorMessageStrategy` can use that `requestMessage` as the `originalMessage` property of the `ErrorMessage` it creates. The `DefaultErrorMessageStrategy` does exactly that. +Starting with version 5.2, all the `MessageHandlingException` instances thrown by the framework components, includes a component `BeanDefinition` resource and source to determine a configuration point form the exception. +In case of XML configuration, a resource is an XML file path and source an XML tag with its `id` attribute. +With Java & Annotation configuration, a resource is a `@Configuration` class and source is a `@Bean` method. +In most case the target integration flow solution is based on the out-of-the-box components and their configuration options. +When an exception happens at runtime, there is no any end-user code involved in stack trace because an execution is against beans, not their configuration. +Including a resource and source of the bean definition helps to determine possible configuration mistakes and provides better developer experience. + [[global-properties]] === Global Properties diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 44145b0abc..cd682d44de 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -59,6 +59,9 @@ See <<./control-bus.adoc#control-bus,Control Bus>> for more information. The `Function>` strategy has been introduced for the aggregator component to merge and compute headers for output messages. See <<./aggregator.adoc#aggregator-api,Aggregator Programming Model>> for more information. +All the `MessageHandlingException` s thrown in the framework, includes now a bean resource and source for back tracking a configuration part in case no end-user code involved. +See <<./configuration.adoc#namespace-errorhandler,Error Handling>> for more information. + [[x5.2-amqp]] ==== AMQP Changes