INT-351 - nextTarget is now reset after a message has been sent.
This commit is contained in:
@@ -153,7 +153,9 @@ public class DefaultEndpoint<T extends MessageHandler> extends AbstractRequestRe
|
||||
}
|
||||
replyMessage = MessageBuilder.fromMessage(replyMessage)
|
||||
.copyHeadersIfAbsent(requestMessage.getHeaders())
|
||||
.setHeaderIfAbsent(MessageHeaders.CORRELATION_ID, requestMessage.getHeaders().getId()).build();
|
||||
.setHeaderIfAbsent(MessageHeaders.CORRELATION_ID, requestMessage.getHeaders().getId())
|
||||
.setNextTarget((String)null)
|
||||
.build();
|
||||
if (!this.getMessageExchangeTemplate().send(replyMessage, replyTarget)) {
|
||||
throw new MessageEndpointReplyException(replyMessage, requestMessage,
|
||||
"failed to send reply to '" + replyTarget + "'");
|
||||
|
||||
@@ -16,16 +16,16 @@
|
||||
|
||||
package org.springframework.integration.endpoint;
|
||||
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertFalse;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.bus.DefaultMessageBus;
|
||||
@@ -45,6 +45,7 @@ import org.springframework.integration.message.selector.MessageSelectorChain;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Marius Bogoevici
|
||||
*/
|
||||
public class DefaultEndpointTests {
|
||||
|
||||
@@ -359,6 +360,32 @@ public class DefaultEndpointTests {
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void nextTargetNotPropagatedPastCurrentEndpoint() {
|
||||
final QueueChannel intermediateItemChannel = new QueueChannel(1);
|
||||
final QueueChannel finalChannel = new QueueChannel(1);
|
||||
DefaultEndpoint<MessageHandler> primaryEndpoint = new DefaultEndpoint<MessageHandler>(new MessageHandler() {
|
||||
public Message<?> handle(Message<?> message) {
|
||||
return MessageBuilder.fromMessage(message).setNextTarget(intermediateItemChannel).build();
|
||||
}
|
||||
});
|
||||
DefaultEndpoint<MessageHandler> secondaryEndpoint = new DefaultEndpoint<MessageHandler>(new MessageHandler() {
|
||||
public Message<?> handle(Message<?> message) {
|
||||
return message;
|
||||
}
|
||||
});
|
||||
secondaryEndpoint.setOutputChannel(finalChannel);
|
||||
Message<String> message = MessageBuilder.fromPayload("test").build();
|
||||
primaryEndpoint.send(message);
|
||||
Message<?> reply = intermediateItemChannel.receive(500);
|
||||
secondaryEndpoint.send(reply);
|
||||
Message<?> replyOnIntermediateChannel = intermediateItemChannel.receive(500);
|
||||
assertNull(replyOnIntermediateChannel);
|
||||
Message<?> replyOnFinalChannel = finalChannel.receive(500);
|
||||
assertNotNull(replyOnFinalChannel);
|
||||
}
|
||||
|
||||
|
||||
private static class TestHandler implements MessageHandler {
|
||||
|
||||
public Message<?> handle(Message<?> message) {
|
||||
|
||||
Reference in New Issue
Block a user