SimpleChannel is now QueueChannel. Also added RendezvousChannel for the zero-capacity SynchronousQueue-based version.
This commit is contained in:
@@ -32,7 +32,7 @@ import org.springframework.context.event.ContextStartedEvent;
|
||||
import org.springframework.context.event.ContextStoppedEvent;
|
||||
import org.springframework.context.support.ClassPathXmlApplicationContext;
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.message.Message;
|
||||
|
||||
/**
|
||||
@@ -42,7 +42,7 @@ public class ApplicationEventSourceAdapterTests {
|
||||
|
||||
@Test
|
||||
public void testAnyApplicationEventSentByDefault() {
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ApplicationEventSourceAdapter adapter = new ApplicationEventSourceAdapter(channel);
|
||||
Message<?> message1 = channel.receive(0);
|
||||
assertNull(message1);
|
||||
@@ -58,7 +58,7 @@ public class ApplicationEventSourceAdapterTests {
|
||||
|
||||
@Test
|
||||
public void testOnlyConfiguredEventTypesAreSent() {
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ApplicationEventSourceAdapter adapter = new ApplicationEventSourceAdapter(channel);
|
||||
List<Class<? extends ApplicationEvent>> eventTypes = new ArrayList<Class<? extends ApplicationEvent>>();
|
||||
eventTypes.add(TestApplicationEvent1.class);
|
||||
|
||||
@@ -27,7 +27,7 @@ import org.springframework.context.ApplicationEvent;
|
||||
import org.springframework.context.ApplicationEventPublisher;
|
||||
import org.springframework.integration.bus.MessageBus;
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
import org.springframework.integration.scheduling.Subscription;
|
||||
@@ -45,7 +45,7 @@ public class ApplicationEventTargetAdapterTests {
|
||||
latch.countDown();
|
||||
}
|
||||
};
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ApplicationEventTargetAdapter adapter = new ApplicationEventTargetAdapter();
|
||||
adapter.setApplicationEventPublisher(publisher);
|
||||
MessageBus bus = new MessageBus();
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
|
||||
<bean id="bus" class="org.springframework.integration.bus.MessageBus"/>
|
||||
|
||||
<bean id="channel" class="org.springframework.integration.channel.SimpleChannel"/>
|
||||
<bean id="channel" class="org.springframework.integration.channel.QueueChannel"/>
|
||||
|
||||
<bean id="adapter" class="org.springframework.integration.adapter.event.ApplicationEventSourceAdapter">
|
||||
<constructor-arg ref="channel"/>
|
||||
|
||||
@@ -30,7 +30,7 @@ import java.util.concurrent.Executors;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.message.Message;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
import org.springframework.mock.web.MockHttpServletRequest;
|
||||
@@ -45,7 +45,7 @@ public class HttpInvokerSourceAdapterTests {
|
||||
|
||||
@Test
|
||||
public void testRequestOnly() throws Exception {
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
HttpInvokerSourceAdapter adapter = new HttpInvokerSourceAdapter(channel);
|
||||
adapter.setExpectReply(false);
|
||||
adapter.afterPropertiesSet();
|
||||
@@ -60,7 +60,7 @@ public class HttpInvokerSourceAdapterTests {
|
||||
|
||||
@Test
|
||||
public void testRequestExpectingReply() throws Exception {
|
||||
final MessageChannel channel = new SimpleChannel();
|
||||
final MessageChannel channel = new QueueChannel();
|
||||
Executors.newSingleThreadExecutor().execute(new Runnable() {
|
||||
public void run() {
|
||||
Message<?> message = channel.receive();
|
||||
|
||||
@@ -26,7 +26,7 @@ import org.springframework.context.ApplicationContext;
|
||||
import org.springframework.context.support.ClassPathXmlApplicationContext;
|
||||
import org.springframework.integration.adapter.rmi.RmiSourceAdapter;
|
||||
import org.springframework.integration.adapter.rmi.RmiTargetAdapter;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.endpoint.HandlerEndpoint;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
|
||||
@@ -35,7 +35,7 @@ import org.springframework.integration.message.StringMessage;
|
||||
*/
|
||||
public class RmiTargetAdapterParserTests {
|
||||
|
||||
private final SimpleChannel testChannel = new SimpleChannel();
|
||||
private final QueueChannel testChannel = new QueueChannel();
|
||||
|
||||
|
||||
@Before
|
||||
|
||||
@@ -25,7 +25,7 @@ import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.adapter.PollingSourceAdapter;
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.message.Message;
|
||||
import org.springframework.integration.scheduling.PollingSchedule;
|
||||
|
||||
@@ -38,7 +38,7 @@ public class ByteStreamSourceAdapterTests {
|
||||
public void testEndOfStream() {
|
||||
byte[] bytes = new byte[] {1,2,3};
|
||||
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ByteStreamSource source = new ByteStreamSource(stream);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
schedule.setInitialDelay(10000);
|
||||
@@ -61,7 +61,7 @@ public class ByteStreamSourceAdapterTests {
|
||||
public void testEndOfStreamWithMaxMessagesPerTask() throws Exception {
|
||||
byte[] bytes = new byte[] {0,1,2,3,4,5,6,7};
|
||||
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ByteStreamSource source = new ByteStreamSource(stream);
|
||||
source.setBytesPerMessage(8);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
@@ -79,7 +79,7 @@ public class ByteStreamSourceAdapterTests {
|
||||
public void testMultipleMessagesWithSingleMessagePerTask() {
|
||||
byte[] bytes = new byte[] {0,1,2,3,4,5,6,7};
|
||||
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ByteStreamSource source = new ByteStreamSource(stream);
|
||||
source.setBytesPerMessage(4);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
@@ -104,7 +104,7 @@ public class ByteStreamSourceAdapterTests {
|
||||
public void testLessThanMaxMessagesAvailable() {
|
||||
byte[] bytes = new byte[] {0,1,2,3,4,5,6,7};
|
||||
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ByteStreamSource source = new ByteStreamSource(stream);
|
||||
source.setBytesPerMessage(4);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
@@ -128,7 +128,7 @@ public class ByteStreamSourceAdapterTests {
|
||||
public void testByteArrayIsTruncated() {
|
||||
byte[] bytes = new byte[] {0,1,2,3,4,5};
|
||||
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ByteStreamSource source = new ByteStreamSource(stream);
|
||||
source.setBytesPerMessage(4);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
@@ -149,7 +149,7 @@ public class ByteStreamSourceAdapterTests {
|
||||
public void testByteArrayIsNotTruncated() {
|
||||
byte[] bytes = new byte[] {0,1,2,3,4,5};
|
||||
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
ByteStreamSource source = new ByteStreamSource(stream);
|
||||
source.setBytesPerMessage(4);
|
||||
source.setShouldTruncate(false);
|
||||
|
||||
@@ -24,7 +24,7 @@ import java.io.IOException;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.channel.DispatcherPolicy;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.dispatcher.PollingDispatcher;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
@@ -62,7 +62,7 @@ public class ByteStreamTargetAdapterTests {
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(3);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
@@ -81,7 +81,7 @@ public class ByteStreamTargetAdapterTests {
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
@@ -100,7 +100,7 @@ public class ByteStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(5);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
@@ -119,7 +119,7 @@ public class ByteStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
@@ -143,7 +143,7 @@ public class ByteStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(5);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
@@ -166,7 +166,7 @@ public class ByteStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
@@ -189,7 +189,7 @@ public class ByteStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
|
||||
|
||||
@@ -25,7 +25,7 @@ import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.adapter.PollingSourceAdapter;
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.message.Message;
|
||||
import org.springframework.integration.scheduling.PollingSchedule;
|
||||
|
||||
@@ -37,7 +37,7 @@ public class CharacterStreamSourceAdapterTests {
|
||||
@Test
|
||||
public void testEndOfStream() {
|
||||
StringReader reader = new StringReader("test");
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
CharacterStreamSource source = new CharacterStreamSource(reader);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
schedule.setInitialDelay(10000);
|
||||
@@ -55,7 +55,7 @@ public class CharacterStreamSourceAdapterTests {
|
||||
@Test
|
||||
public void testEndOfStreamWithMaxMessagesPerTask() {
|
||||
StringReader reader = new StringReader("test");
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
CharacterStreamSource source = new CharacterStreamSource(reader);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
schedule.setInitialDelay(10000);
|
||||
@@ -72,7 +72,7 @@ public class CharacterStreamSourceAdapterTests {
|
||||
public void testMultipleLinesWithSingleMessagePerTask() {
|
||||
String s = "test1" + System.getProperty("line.separator") + "test2";
|
||||
StringReader reader = new StringReader(s);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
CharacterStreamSource source = new CharacterStreamSource(reader);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
schedule.setInitialDelay(10000);
|
||||
@@ -92,7 +92,7 @@ public class CharacterStreamSourceAdapterTests {
|
||||
public void testLessThanMaxMessagesAvailable() {
|
||||
String s = "test1" + System.getProperty("line.separator") + "test2";
|
||||
StringReader reader = new StringReader(s);
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
CharacterStreamSource source = new CharacterStreamSource(reader);
|
||||
PollingSchedule schedule = new PollingSchedule(1000);
|
||||
schedule.setInitialDelay(5000);
|
||||
|
||||
@@ -24,7 +24,7 @@ import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.channel.DispatcherPolicy;
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.dispatcher.PollingDispatcher;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
@@ -44,7 +44,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
|
||||
@Test
|
||||
public void testTwoStringsAndNoNewLinesByDefault() {
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
@@ -59,7 +59,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
|
||||
@Test
|
||||
public void testTwoStringsWithNewLines() {
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
@@ -80,7 +80,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
@@ -96,7 +96,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(10);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
dispatcher.subscribe(adapter);
|
||||
@@ -109,7 +109,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
|
||||
@Test
|
||||
public void testSingleNonStringObject() {
|
||||
MessageChannel channel = new SimpleChannel();
|
||||
MessageChannel channel = new QueueChannel();
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
@@ -127,7 +127,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
TestObject testObject1 = new TestObject("foo");
|
||||
@@ -145,7 +145,7 @@ public class CharacterStreamTargetAdapterTests {
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
SimpleChannel channel = new SimpleChannel(5, dispatcherPolicy);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
PollingDispatcher dispatcher = new PollingDispatcher(channel, null);
|
||||
dispatcher.subscribe(adapter);
|
||||
|
||||
Reference in New Issue
Block a user