The MessageHeader now has a single Object-typed 'returnAddress' property rather than both 'replyChannel' and 'replyChannelName'. This also removes a package cycle between 'message' and 'channel' (part of INT-112).
This commit is contained in:
@@ -55,7 +55,7 @@ public class MessageBusTests {
|
||||
MessageChannel targetChannel = new SimpleChannel();
|
||||
bus.registerChannel("sourceChannel", sourceChannel);
|
||||
StringMessage message = new StringMessage("test");
|
||||
message.getHeader().setReplyChannelName("targetChannel");
|
||||
message.getHeader().setReturnAddress("targetChannel");
|
||||
sourceChannel.send(message);
|
||||
bus.registerChannel("targetChannel", targetChannel);
|
||||
MessageHandler handler = new MessageHandler() {
|
||||
@@ -105,13 +105,13 @@ public class MessageBusTests {
|
||||
SimpleChannel outputChannel2 = new SimpleChannel();
|
||||
MessageHandler handler1 = new MessageHandler() {
|
||||
public Message<?> handle(Message<?> message) {
|
||||
message.getHeader().setReplyChannelName("output1");
|
||||
message.getHeader().setReturnAddress("output1");
|
||||
return message;
|
||||
}
|
||||
};
|
||||
MessageHandler handler2 = new MessageHandler() {
|
||||
public Message<?> handle(Message<?> message) {
|
||||
message.getHeader().setReplyChannelName("output2");
|
||||
message.getHeader().setReturnAddress("output2");
|
||||
return message;
|
||||
}
|
||||
};
|
||||
@@ -137,13 +137,13 @@ public class MessageBusTests {
|
||||
SimpleChannel outputChannel2 = new SimpleChannel();
|
||||
MessageHandler handler1 = new MessageHandler() {
|
||||
public Message<?> handle(Message<?> message) {
|
||||
message.getHeader().setReplyChannelName("output1");
|
||||
message.getHeader().setReturnAddress("output1");
|
||||
return message;
|
||||
}
|
||||
};
|
||||
MessageHandler handler2 = new MessageHandler() {
|
||||
public Message<?> handle(Message<?> message) {
|
||||
message.getHeader().setReplyChannelName("output2");
|
||||
message.getHeader().setReturnAddress("output2");
|
||||
return message;
|
||||
}
|
||||
};
|
||||
|
||||
@@ -124,7 +124,7 @@ public class EndpointParserTests {
|
||||
((Lifecycle) endpoint).start();
|
||||
Message<?> message = new StringMessage("test");
|
||||
MessageChannel replyChannel = new SimpleChannel();
|
||||
message.getHeader().setReplyChannel(replyChannel);
|
||||
message.getHeader().setReturnAddress(replyChannel);
|
||||
endpoint.handle(message);
|
||||
Message<?> reply = replyChannel.receive(500);
|
||||
assertNotNull(reply);
|
||||
|
||||
@@ -82,7 +82,7 @@ public class DefaultMessageEndpointTests {
|
||||
endpoint.setHandler(handler);
|
||||
endpoint.start();
|
||||
StringMessage testMessage = new StringMessage(1, "test");
|
||||
testMessage.getHeader().setReplyChannel(replyChannel);
|
||||
testMessage.getHeader().setReturnAddress(replyChannel);
|
||||
endpoint.handle(testMessage);
|
||||
endpoint.stop();
|
||||
Message<?> reply = replyChannel.receive(50);
|
||||
@@ -105,7 +105,7 @@ public class DefaultMessageEndpointTests {
|
||||
endpoint.setHandler(handler);
|
||||
endpoint.start();
|
||||
StringMessage testMessage = new StringMessage(1, "test");
|
||||
testMessage.getHeader().setReplyChannelName("replyChannel");
|
||||
testMessage.getHeader().setReturnAddress("replyChannel");
|
||||
endpoint.handle(testMessage);
|
||||
endpoint.stop();
|
||||
Message<?> reply = replyChannel.receive(50);
|
||||
@@ -114,7 +114,7 @@ public class DefaultMessageEndpointTests {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testReplyChannelTakesPrecedenceOverReplyChannelName() throws Exception {
|
||||
public void testDynamicReplyChannel() throws Exception {
|
||||
final MessageChannel replyChannel1 = new SimpleChannel();
|
||||
final MessageChannel replyChannel2 = new SimpleChannel();
|
||||
ChannelRegistry channelRegistry = new DefaultChannelRegistry();
|
||||
@@ -125,18 +125,25 @@ public class DefaultMessageEndpointTests {
|
||||
}
|
||||
};
|
||||
DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint();
|
||||
endpoint.setChannelRegistry(channelRegistry);
|
||||
endpoint.setHandler(handler);
|
||||
endpoint.start();
|
||||
StringMessage testMessage = new StringMessage(1, "test");
|
||||
testMessage.getHeader().setReplyChannel(replyChannel1);
|
||||
testMessage.getHeader().setReplyChannelName("replyChannel2");
|
||||
StringMessage testMessage = new StringMessage("test");
|
||||
testMessage.getHeader().setReturnAddress(replyChannel1);
|
||||
endpoint.handle(testMessage);
|
||||
endpoint.stop();
|
||||
Message<?> reply1 = replyChannel1.receive(50);
|
||||
assertNotNull(reply1);
|
||||
assertEquals("hello test", reply1.getPayload());
|
||||
Message<?> reply2 = replyChannel2.receive(0);
|
||||
assertNull(reply2);
|
||||
testMessage.getHeader().setReturnAddress("replyChannel2");
|
||||
endpoint.handle(testMessage);
|
||||
reply1 = replyChannel1.receive(0);
|
||||
assertNull(reply1);
|
||||
reply2 = replyChannel2.receive(0);
|
||||
assertNotNull(reply2);
|
||||
assertEquals("hello test", reply2.getPayload());
|
||||
endpoint.stop();
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -249,7 +256,7 @@ public class DefaultMessageEndpointTests {
|
||||
endpoint.setHandler(new ConcurrentHandler(handler, createExecutor()));
|
||||
endpoint.start();
|
||||
StringMessage message = new StringMessage(1, "test");
|
||||
message.getHeader().setReplyChannelName("replyChannel");
|
||||
message.getHeader().setReturnAddress("replyChannel");
|
||||
endpoint.handle(message);
|
||||
endpoint.stop();
|
||||
latch.await(500, TimeUnit.MILLISECONDS);
|
||||
@@ -304,7 +311,7 @@ public class DefaultMessageEndpointTests {
|
||||
endpoint.setConcurrencyPolicy(new ConcurrencyPolicy(3, 14));
|
||||
endpoint.start();
|
||||
StringMessage message = new StringMessage(1, "test");
|
||||
message.getHeader().setReplyChannelName("replyChannel");
|
||||
message.getHeader().setReturnAddress("replyChannel");
|
||||
endpoint.handle(message);
|
||||
endpoint.stop();
|
||||
latch.await(500, TimeUnit.MILLISECONDS);
|
||||
|
||||
@@ -138,7 +138,7 @@ public class AggregatingMessageHandlerTests {
|
||||
message.getHeader().setCorrelationId(correlationId);
|
||||
message.getHeader().setSequenceSize(sequenceSize);
|
||||
message.getHeader().setSequenceNumber(sequenceNumber);
|
||||
message.getHeader().setReplyChannel(replyChannel);
|
||||
message.getHeader().setReturnAddress(replyChannel);
|
||||
return message;
|
||||
}
|
||||
|
||||
@@ -182,7 +182,13 @@ public class AggregatingMessageHandlerTests {
|
||||
try {
|
||||
Message<?> result = this.aggregator.handle(message);
|
||||
if (result != null) {
|
||||
message.getHeader().getReplyChannel().send(result);
|
||||
Object returnAddress = message.getHeader().getReturnAddress();
|
||||
if (returnAddress instanceof MessageChannel) {
|
||||
((MessageChannel) returnAddress).send(result);
|
||||
}
|
||||
else {
|
||||
throw new IllegalStateException("'returnAddress' was not a MessageChannel instance");
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (Exception e) {
|
||||
|
||||
Reference in New Issue
Block a user