diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java index 4db7f20745..b9b610dfaa 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java @@ -50,7 +50,7 @@ public abstract class AbstractAggregatingMessageGroupProcessor implements Messag MessageBuilder builder = (payload instanceof Message) ? MessageBuilder.fromMessage((Message) payload) : MessageBuilder.withPayload(payload); Message message = builder.copyHeadersIfAbsent(headers).build(); - channelTemplate.send(message, outputChannel); + channelTemplate.send(outputChannel, message); } /** diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/PassThroughMessageGroupProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/PassThroughMessageGroupProcessor.java index ba7b98e973..ebfc95228c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/PassThroughMessageGroupProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/PassThroughMessageGroupProcessor.java @@ -30,7 +30,7 @@ public class PassThroughMessageGroupProcessor implements MessageGroupProcessor { public void processAndSend(MessageGroup group, MessagingTemplate messagingTemplate, MessageChannel outputChannel) { for (Message message : group.getUnmarked()) { - messagingTemplate.send(message, outputChannel); + messagingTemplate.send(outputChannel, message); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/ResequencingMessageGroupProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/ResequencingMessageGroupProcessor.java index 4e71a725e7..1e3940c3b6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/ResequencingMessageGroupProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/ResequencingMessageGroupProcessor.java @@ -51,7 +51,7 @@ public class ResequencingMessageGroupProcessor implements MessageGroupProcessor List> sorted = new ArrayList>(messages); Collections.sort(sorted, comparator); for (Message message : sorted) { - messagingTemplate.send(message, outputChannel); + messagingTemplate.send(outputChannel, message); } } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aop/MessagePublishingInterceptor.java b/spring-integration-core/src/main/java/org/springframework/integration/aop/MessagePublishingInterceptor.java index 93838d062d..ecdcea49f1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aop/MessagePublishingInterceptor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aop/MessagePublishingInterceptor.java @@ -136,7 +136,7 @@ public class MessagePublishingInterceptor implements MethodInterceptor { channel = this.channelResolver.resolveChannelName(channelName); } if (channel != null) { - this.messagingTemplate.send(message, channel); + this.messagingTemplate.send(channel, message); } else { this.messagingTemplate.send(message); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/MessagingTemplate.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/MessagingTemplate.java index 97d9a4e01c..d163ad1da4 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/MessagingTemplate.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/MessagingTemplate.java @@ -76,14 +76,14 @@ public class MessagingTemplate implements InitializingBean { /** - * Create a MessageChannelTemplate with no default channel. Note, that one + * Create a MessagingTemplate with no default channel. Note, that one * may be provided by invoking {@link #setDefaultChannel(MessageChannel)}. */ public MessagingTemplate() { } /** - * Create a MessageChannelTemplate with the given default channel. + * Create a MessagingTemplate with the given default channel. */ public MessagingTemplate(MessageChannel defaultChannel) { this.defaultChannel = defaultChannel; @@ -166,19 +166,19 @@ public class MessagingTemplate implements InitializingBean { } public boolean send(final Message message) { - return this.send(message, this.getRequiredDefaultChannel()); + return this.send(this.getRequiredDefaultChannel(), message); } - public boolean send(final Message message, final MessageChannel channel) { + public boolean send(final MessageChannel channel, final Message message) { TransactionTemplate txTemplate = this.getTransactionTemplate(); if (txTemplate != null) { return txTemplate.execute(new TransactionCallback() { public Boolean doInTransaction(TransactionStatus status) { - return doSend(message, channel); + return doSend(channel, message); } }); } - return this.doSend(message, channel); + return this.doSend(channel, message); } public Message receive() { @@ -201,22 +201,22 @@ public class MessagingTemplate implements InitializingBean { } public Message sendAndReceive(final Message request) { - return this.sendAndReceive(request, this.getRequiredDefaultChannel()); + return this.sendAndReceive(this.getRequiredDefaultChannel(), request); } - public Message sendAndReceive(final Message request, final MessageChannel channel) { + public Message sendAndReceive(final MessageChannel channel, final Message request) { TransactionTemplate txTemplate = this.getTransactionTemplate(); if (txTemplate != null) { return txTemplate.execute(new TransactionCallback>() { public Message doInTransaction(TransactionStatus status) { - return doSendAndReceive(request, channel); + return doSendAndReceive(channel, request); } }); } - return this.doSendAndReceive(request, channel); + return this.doSendAndReceive(channel, request); } - private boolean doSend(Message message, MessageChannel channel) { + private boolean doSend(MessageChannel channel, Message message) { Assert.notNull(channel, "channel must not be null"); long timeout = this.sendTimeout; boolean sent = (timeout >= 0) @@ -240,7 +240,7 @@ public class MessagingTemplate implements InitializingBean { return message; } - private Message doSendAndReceive(Message request, MessageChannel channel) { + private Message doSendAndReceive(MessageChannel channel, Message request) { Object originalReplyChannelHeader = request.getHeaders().getReplyChannel(); Object originalErrorChannelHeader = request.getHeaders().getErrorChannel(); TemporaryReplyChannel replyChannel = new TemporaryReplyChannel(this.receiveTimeout); @@ -248,7 +248,7 @@ public class MessagingTemplate implements InitializingBean { .setReplyChannel(replyChannel) .setErrorChannel(replyChannel) .build(); - if (!this.doSend(request, channel)) { + if (!this.doSend(channel, request)) { throw new MessageDeliveryException(request, "failed to send message to channel"); } Message reply = this.doReceive(replyChannel); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java index e6b3968bae..82a3639c55 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java @@ -52,7 +52,7 @@ public abstract class MessageProducerSupport extends AbstractEndpoint implements if (message != null) { message.getHeaders().getHistory().addEvent(this); } - return this.messagingTemplate.send(message, this.outputChannel); + return this.messagingTemplate.send(this.outputChannel, message); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java index 9e4893fbb2..0987fb951e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java @@ -75,7 +75,7 @@ public class SourcePollingChannelAdapter extends AbstractPollingEndpoint { protected boolean doPoll() { Message message = this.source.receive(); if (message != null) { - return this.messagingTemplate.send(message, this.outputChannel); + return this.messagingTemplate.send(this.outputChannel, message); } return false; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java b/spring-integration-core/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java index 2c48fb0f25..92bd8d5159 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java @@ -144,7 +144,7 @@ public abstract class AbstractMessagingGateway extends AbstractEndpoint { "send is not supported, because no request channel has been configured"); Message message = this.toMessage(object); Assert.notNull(message, "message must not be null"); - if (!this.messagingTemplate.send(message, this.requestChannel)) { + if (!this.messagingTemplate.send(this.requestChannel, message)) { throw new MessageDeliveryException(message, "failed to send Message to channel"); } } @@ -195,7 +195,7 @@ public abstract class AbstractMessagingGateway extends AbstractEndpoint { Message reply = null; Throwable error = null; try { - reply = this.messagingTemplate.sendAndReceive(message, this.requestChannel); + reply = this.messagingTemplate.sendAndReceive(this.requestChannel, message); if (reply instanceof ErrorMessage) { error = ((ErrorMessage) reply).getPayload(); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java index 3d5a6f52d4..3c07fdb575 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java @@ -128,7 +128,7 @@ public abstract class AbstractReplyProducingMessageHandler extends AbstractMessa if (logger.isDebugEnabled()) { logger.debug("handler '" + this + "' sending reply Message: " + replyMessage); } - return this.messagingTemplate.send(replyMessage, replyChannel); + return this.messagingTemplate.send(replyChannel, replyMessage); } /** 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 88e36f6460..8581899652 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 @@ -241,7 +241,7 @@ public class DelayHandler extends IntegrationObjectSupport implements MessageHan MessageChannel errorChannel = resolveErrorChannelIfPossible(message); if (errorChannel != null) { ErrorMessage errorMessage = new ErrorMessage(exception); - boolean sent = messagingTemplate.send(errorMessage, errorChannel); + boolean sent = messagingTemplate.send(errorChannel, errorMessage); if (!sent && logger.isWarnEnabled()) { logger.warn("Failed to send MessageDeliveryException to error channel.", exception); } @@ -263,7 +263,7 @@ public class DelayHandler extends IntegrationObjectSupport implements MessageHan private void sendMessageToReplyChannel(Message message) { MessageChannel replyChannel = this.resolveReplyChannel(message); - this.messagingTemplate.send(message, replyChannel); + this.messagingTemplate.send(replyChannel, message); } private MessageChannel resolveReplyChannel(Message message) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMessageRouter.java b/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMessageRouter.java index 178be3da1e..d56e968e55 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMessageRouter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMessageRouter.java @@ -116,7 +116,7 @@ public abstract class AbstractMessageRouter extends AbstractMessageHandler { .setHeader(MessageHeaders.ID, UUID.randomUUID()) .build(); if (channel != null) { - if (this.messagingTemplate.send(messageToSend, channel)) { + if (this.messagingTemplate.send(channel, messageToSend)) { sent = true; } else if (!this.ignoreSendFailures) { @@ -128,7 +128,7 @@ public abstract class AbstractMessageRouter extends AbstractMessageHandler { } if (!sent) { if (this.defaultOutputChannel != null) { - sent = this.messagingTemplate.send(message, this.defaultOutputChannel); + sent = this.messagingTemplate.send(this.defaultOutputChannel, message); } else if (this.resolutionRequired) { throw new MessageDeliveryException(message, diff --git a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java index f89b2478a8..0ff10b3d47 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java @@ -237,7 +237,7 @@ public class AggregatorTests { for (Message message : group.getUnmarked()) { product *= (Integer) message.getPayload(); } - messagingTemplate.send(MessageBuilder.withPayload(product).build(), outputChannel); + messagingTemplate.send(outputChannel, MessageBuilder.withPayload(product).build()); } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java index 45a3f38381..ad6ced6f28 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java @@ -350,7 +350,7 @@ public class ConcurrentAggregatorTests { for (Message message : group.getUnmarked()) { product *= (Integer) message.getPayload(); } - messagingTemplate.send(MessageBuilder.withPayload(product).build(), outputChannel); + messagingTemplate.send(outputChannel, MessageBuilder.withPayload(product).build()); } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/MethodInvokingMessageGroupProcessorTests.java b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/MethodInvokingMessageGroupProcessorTests.java index ac691e5237..880b378ea4 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/MethodInvokingMessageGroupProcessorTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/MethodInvokingMessageGroupProcessorTests.java @@ -96,7 +96,7 @@ public class MethodInvokingMessageGroupProcessorTests { when(messageGroupMock.getUnmarked()).thenReturn(messagesUpForProcessing); processor.processAndSend(messageGroupMock, messagingTemplate, outputChannel); // verify - verify(messagingTemplate).send(messageCaptor.capture(), eq(outputChannel)); + verify(messagingTemplate).send(eq(outputChannel), messageCaptor.capture()); assertThat((Integer) messageCaptor.getValue().getPayload(), is(7)); } @@ -120,7 +120,7 @@ public class MethodInvokingMessageGroupProcessorTests { when(messageGroupMock.getUnmarked()).thenReturn(messagesUpForProcessing); processor.processAndSend(messageGroupMock, messagingTemplate, outputChannel); // verify - verify(messagingTemplate).send(messageCaptor.capture(), eq(outputChannel)); + verify(messagingTemplate).send(eq(outputChannel), messageCaptor.capture()); assertThat((Integer) messageCaptor.getValue().getPayload(), is(7)); } @@ -155,7 +155,7 @@ public class MethodInvokingMessageGroupProcessorTests { when(messageGroupMock.getUnmarked()).thenReturn(messagesUpForProcessing); processor.processAndSend(messageGroupMock, messagingTemplate, outputChannel); // verify - verify(messagingTemplate).send(messageCaptor.capture(), eq(outputChannel)); + verify(messagingTemplate).send(eq(outputChannel), messageCaptor.capture()); assertThat((Integer) messageCaptor.getValue().getPayload(), is(7)); } @@ -186,7 +186,7 @@ public class MethodInvokingMessageGroupProcessorTests { when(messageGroupMock.getUnmarked()).thenReturn(messagesUpForProcessing); processor.processAndSend(messageGroupMock, messagingTemplate, outputChannel); // verify - verify(messagingTemplate).send(messageCaptor.capture(), eq(outputChannel)); + verify(messagingTemplate).send(eq(outputChannel), messageCaptor.capture()); assertThat((Integer) messageCaptor.getValue().getPayload(), is(7)); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/MessagingTemplateTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/MessagingTemplateTests.java index 1106369e2e..0e9de84ffa 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/MessagingTemplateTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/MessagingTemplateTests.java @@ -80,7 +80,7 @@ public class MessagingTemplateTests { public void send() { MessagingTemplate template = new MessagingTemplate(); QueueChannel channel = new QueueChannel(); - template.send(new StringMessage("test"), channel); + template.send(channel, new StringMessage("test")); Message reply = channel.receive(0); assertNotNull(reply); assertEquals("test", reply.getPayload()); @@ -112,7 +112,7 @@ public class MessagingTemplateTests { QueueChannel explicitChannel = new QueueChannel(); QueueChannel defaultChannel = new QueueChannel(); MessagingTemplate template = new MessagingTemplate(defaultChannel); - template.send(new StringMessage("test"), explicitChannel); + template.send(explicitChannel, new StringMessage("test")); Message reply = explicitChannel.receive(0); assertNotNull(reply); assertEquals("test", reply.getPayload()); @@ -182,7 +182,7 @@ public class MessagingTemplateTests { public void sendAndReceive() { MessagingTemplate template = new MessagingTemplate(); template.setReceiveTimeout(3000); - Message reply = template.sendAndReceive(new StringMessage("test"), this.requestChannel); + Message reply = template.sendAndReceive(this.requestChannel, new StringMessage("test")); assertEquals("TEST", reply.getPayload()); } @@ -201,7 +201,7 @@ public class MessagingTemplateTests { MessagingTemplate template = new MessagingTemplate(defaultChannel); template.setReceiveTimeout(3000); Message message = new StringMessage("test"); - Message reply = template.sendAndReceive(message, this.requestChannel); + Message reply = template.sendAndReceive(this.requestChannel, message); assertEquals("TEST", reply.getPayload()); assertNull(defaultChannel.receive(0)); } @@ -228,9 +228,9 @@ public class MessagingTemplateTests { Message message1 = MessageBuilder.withPayload("test1").setReplyChannel(replyChannel).build(); Message message2 = MessageBuilder.withPayload("test2").setReplyChannel(replyChannel).build(); Message message3 = MessageBuilder.withPayload("test3").setReplyChannel(replyChannel).build(); - template.send(message1, this.requestChannel); - template.send(message2, this.requestChannel); - template.send(message3, this.requestChannel); + template.send(this.requestChannel, message1); + template.send(this.requestChannel, message2); + template.send(this.requestChannel, message3); latch.await(2000, TimeUnit.MILLISECONDS); assertEquals(0, latch.getCount()); assertTrue(replies.contains("TEST1")); diff --git a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/messages/XmppMessageDrivenEndpoint.java b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/messages/XmppMessageDrivenEndpoint.java index 724e01be59..6b44c6242b 100644 --- a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/messages/XmppMessageDrivenEndpoint.java +++ b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/messages/XmppMessageDrivenEndpoint.java @@ -138,7 +138,7 @@ public class XmppMessageDrivenEndpoint extends AbstractEndpoint implements Lifec MessageBuilder messageBuilder = MessageBuilder.withPayload(payload) .setHeader(XmppHeaders.TYPE, xmppMessage.getType()) .setHeader(XmppHeaders.CHAT, chat); - messagingTemplate.send(messageBuilder.build(), requestChannel); + messagingTemplate.send(requestChannel, messageBuilder.build()); } } diff --git a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/presence/XmppRosterEventMessageDrivenEndpoint.java b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/presence/XmppRosterEventMessageDrivenEndpoint.java index 6047b499aa..1240af2db8 100644 --- a/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/presence/XmppRosterEventMessageDrivenEndpoint.java +++ b/spring-integration-xmpp/src/main/java/org/springframework/integration/xmpp/presence/XmppRosterEventMessageDrivenEndpoint.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.xmpp.presence; import org.apache.commons.lang.StringUtils; @@ -107,8 +108,9 @@ public class XmppRosterEventMessageDrivenEndpoint extends AbstractEndpoint imple protected void forwardRosterEventMessage(Presence presence) { try { Message msg = this.messageMapper.toMessage(presence); - messagingTemplate.send(msg, requestChannel); - } catch (Exception e) { + messagingTemplate.send(requestChannel, msg); + } + catch (Exception e) { logger.error("Failed to map packet to message ", e); } } diff --git a/spring-integration-xmpp/src/test/java/org/springframework/integration/xmpp/config/XmppHeaderEnricherParserTests.java b/spring-integration-xmpp/src/test/java/org/springframework/integration/xmpp/config/XmppHeaderEnricherParserTests.java index 5ff0f5d743..d2af444579 100644 --- a/spring-integration-xmpp/src/test/java/org/springframework/integration/xmpp/config/XmppHeaderEnricherParserTests.java +++ b/spring-integration-xmpp/src/test/java/org/springframework/integration/xmpp/config/XmppHeaderEnricherParserTests.java @@ -61,7 +61,7 @@ public class XmppHeaderEnricherParserTests { logger.debug(String.format("%s=%s (class: %s)", h, message.getHeaders().get(h), message.getHeaders().get(h).getClass().toString())); } }); - messagingTemplate.send(MessageBuilder.withPayload("foo").build(), input); + messagingTemplate.send(input, MessageBuilder.withPayload("foo").build()); } }