INT-1129 modified the argument order for MessagingTemplate send(..) methods so that MessageChannel always comes first - consistent with Destination/destinationName being the first argument in all JmsTemplate send(..) methods
This commit is contained in:
@@ -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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -51,7 +51,7 @@ public class ResequencingMessageGroupProcessor implements MessageGroupProcessor
|
||||
List<Message<?>> sorted = new ArrayList<Message<?>>(messages);
|
||||
Collections.sort(sorted, comparator);
|
||||
for (Message<?> message : sorted) {
|
||||
messagingTemplate.send(message, outputChannel);
|
||||
messagingTemplate.send(outputChannel, message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<Boolean>() {
|
||||
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<Message<?>>() {
|
||||
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);
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String> message1 = MessageBuilder.withPayload("test1").setReplyChannel(replyChannel).build();
|
||||
Message<String> message2 = MessageBuilder.withPayload("test2").setReplyChannel(replyChannel).build();
|
||||
Message<String> 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"));
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user