Added support for interceptors on MessageEndpoints.

This commit is contained in:
Mark Fisher
2008-06-30 22:26:04 +00:00
parent f9ad394850
commit 07b2179baa
35 changed files with 836 additions and 908 deletions

View File

@@ -26,7 +26,9 @@ import org.junit.Test;
import org.springframework.integration.channel.MessageChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.endpoint.SourceEndpoint;
import org.springframework.integration.message.CommandMessage;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.PollCommand;
import org.springframework.integration.scheduling.PollingSchedule;
/**
@@ -40,10 +42,8 @@ public class ByteStreamSourceTests {
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
MessageChannel channel = new QueueChannel();
ByteStreamSource source = new ByteStreamSource(stream);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.run();
SourceEndpoint endpoint = new SourceEndpoint(source, channel);
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message1 = channel.receive(500);
byte[] payload = (byte[]) message1.getPayload();
assertEquals(3, payload.length);
@@ -52,74 +52,7 @@ public class ByteStreamSourceTests {
assertEquals(3, payload[2]);
Message<?> message2 = channel.receive(0);
assertNull(message2);
endpoint.run();
Message<?> message3 = channel.receive(0);
assertNull(message3);
}
@Test
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 QueueChannel();
ByteStreamSource source = new ByteStreamSource(stream);
source.setBytesPerMessage(8);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(5);
endpoint.run();
Message<?> message1 = channel.receive(500);
assertEquals(8, ((byte[]) message1.getPayload()).length);
Message<?> message2 = channel.receive(0);
assertNull(message2);
}
@Test
public void testMultipleMessagesWithSingleMessagePerTask() {
byte[] bytes = new byte[] {0,1,2,3,4,5,6,7};
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
MessageChannel channel = new QueueChannel();
ByteStreamSource source = new ByteStreamSource(stream);
source.setBytesPerMessage(4);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(1);
endpoint.run();
Message<?> message1 = channel.receive(0);
byte[] bytes1 = (byte[]) message1.getPayload();
assertEquals(4, bytes1.length);
assertEquals(0, bytes1[0]);
Message<?> message2 = channel.receive(0);
assertNull(message2);
endpoint.run();
Message<?> message3 = channel.receive(0);
byte[] bytes3 = (byte[]) message3.getPayload();
assertEquals(4, bytes3.length);
assertEquals(4, bytes3[0]);
}
@Test
public void testLessThanMaxMessagesAvailable() {
byte[] bytes = new byte[] {0,1,2,3,4,5,6,7};
ByteArrayInputStream stream = new ByteArrayInputStream(bytes);
MessageChannel channel = new QueueChannel();
ByteStreamSource source = new ByteStreamSource(stream);
source.setBytesPerMessage(4);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(5);
endpoint.run();
Message<?> message1 = channel.receive(0);
byte[] bytes1 = (byte[]) message1.getPayload();
assertEquals(4, bytes1.length);
assertEquals(0, bytes1[0]);
Message<?> message2 = channel.receive(0);
byte[] bytes2 = (byte[]) message2.getPayload();
assertEquals(4, bytes2.length);
assertEquals(4, bytes2[0]);
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message3 = channel.receive(0);
assertNull(message3);
}
@@ -133,14 +66,13 @@ public class ByteStreamSourceTests {
source.setBytesPerMessage(4);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(1);
endpoint.run();
SourceEndpoint endpoint = new SourceEndpoint(source, channel);
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message1 = channel.receive(0);
assertEquals(4, ((byte[]) message1.getPayload()).length);
Message<?> message2 = channel.receive(0);
assertNull(message2);
endpoint.run();
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message3 = channel.receive(0);
assertEquals(2, ((byte[]) message3.getPayload()).length);
}
@@ -155,14 +87,13 @@ public class ByteStreamSourceTests {
source.setShouldTruncate(false);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(1);
endpoint.run();
SourceEndpoint endpoint = new SourceEndpoint(source, channel);
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message1 = channel.receive(0);
assertEquals(4, ((byte[]) message1.getPayload()).length);
Message<?> message2 = channel.receive(0);
assertNull(message2);
endpoint.run();
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message3 = channel.receive(0);
assertEquals(4, ((byte[]) message3.getPayload()).length);
assertEquals(0, ((byte[]) message3.getPayload())[3]);

View File

@@ -25,7 +25,6 @@ import org.junit.Test;
import org.springframework.integration.channel.DispatcherPolicy;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.dispatcher.DefaultPollingDispatcher;
import org.springframework.integration.dispatcher.PollingDispatcherTask;
import org.springframework.integration.message.GenericMessage;
import org.springframework.integration.message.StringMessage;
@@ -64,8 +63,8 @@ public class ByteStreamTargetTests {
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
dispatcherPolicy.setMaxMessagesPerTask(3);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
@@ -83,8 +82,8 @@ public class ByteStreamTargetTests {
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
dispatcherPolicy.setMaxMessagesPerTask(2);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
@@ -102,8 +101,8 @@ public class ByteStreamTargetTests {
dispatcherPolicy.setMaxMessagesPerTask(5);
dispatcherPolicy.setReceiveTimeout(0);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
@@ -121,8 +120,8 @@ public class ByteStreamTargetTests {
dispatcherPolicy.setMaxMessagesPerTask(2);
dispatcherPolicy.setReceiveTimeout(0);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
@@ -145,8 +144,8 @@ public class ByteStreamTargetTests {
dispatcherPolicy.setMaxMessagesPerTask(5);
dispatcherPolicy.setReceiveTimeout(0);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
@@ -168,8 +167,8 @@ public class ByteStreamTargetTests {
dispatcherPolicy.setMaxMessagesPerTask(2);
dispatcherPolicy.setReceiveTimeout(0);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);
@@ -191,8 +190,8 @@ public class ByteStreamTargetTests {
dispatcherPolicy.setMaxMessagesPerTask(2);
dispatcherPolicy.setReceiveTimeout(0);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new GenericMessage<byte[]>(new byte[] {1,2,3}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {4,5,6}), 0);
channel.send(new GenericMessage<byte[]>(new byte[] {7,8,9}), 0);

View File

@@ -26,7 +26,9 @@ import org.junit.Test;
import org.springframework.integration.channel.MessageChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.endpoint.SourceEndpoint;
import org.springframework.integration.message.CommandMessage;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.PollCommand;
import org.springframework.integration.scheduling.PollingSchedule;
/**
@@ -41,68 +43,13 @@ public class CharacterStreamSourceTests {
CharacterStreamSource source = new CharacterStreamSource(reader);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.run();
SourceEndpoint endpoint = new SourceEndpoint(source, channel);
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message1 = channel.receive(0);
assertEquals("test", message1.getPayload());
Message<?> message2 = channel.receive(0);
assertNull(message2);
endpoint.run();
Message<?> message3 = channel.receive(0);
assertNull(message3);
}
@Test
public void testEndOfStreamWithMaxMessagesPerTask() {
StringReader reader = new StringReader("test");
MessageChannel channel = new QueueChannel();
CharacterStreamSource source = new CharacterStreamSource(reader);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(5);
endpoint.run();
Message<?> message1 = channel.receive(0);
assertEquals("test", message1.getPayload());
Message<?> message2 = channel.receive(0);
assertNull(message2);
}
@Test
public void testMultipleLinesWithSingleMessagePerTask() {
String s = "test1" + System.getProperty("line.separator") + "test2";
StringReader reader = new StringReader(s);
MessageChannel channel = new QueueChannel();
CharacterStreamSource source = new CharacterStreamSource(reader);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(10000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(1);
endpoint.run();
Message<?> message1 = channel.receive(0);
assertEquals("test1", message1.getPayload());
Message<?> message2 = channel.receive(0);
assertNull(message2);
endpoint.run();
Message<?> message3 = channel.receive(0);
assertEquals("test2", message3.getPayload());
}
@Test
public void testLessThanMaxMessagesAvailable() {
String s = "test1" + System.getProperty("line.separator") + "test2";
StringReader reader = new StringReader(s);
MessageChannel channel = new QueueChannel();
CharacterStreamSource source = new CharacterStreamSource(reader);
PollingSchedule schedule = new PollingSchedule(1000);
schedule.setInitialDelay(5000);
SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule);
endpoint.setMaxMessagesPerTask(5);
endpoint.run();
Message<?> message1 = channel.receive(500);
assertEquals("test1", message1.getPayload());
Message<?> message2 = channel.receive(500);
assertEquals("test2", message2.getPayload());
endpoint.invoke(new CommandMessage(new PollCommand()));
Message<?> message3 = channel.receive(0);
assertNull(message3);
}

View File

@@ -25,7 +25,6 @@ import org.junit.Test;
import org.springframework.integration.channel.DispatcherPolicy;
import org.springframework.integration.channel.MessageChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.dispatcher.DefaultPollingDispatcher;
import org.springframework.integration.dispatcher.PollingDispatcherTask;
import org.springframework.integration.message.GenericMessage;
import org.springframework.integration.message.StringMessage;
@@ -48,8 +47,8 @@ public class CharacterStreamTargetTests {
MessageChannel channel = new QueueChannel();
StringWriter writer = new StringWriter();
CharacterStreamTarget target = new CharacterStreamTarget(writer);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
task.run();
@@ -64,8 +63,8 @@ public class CharacterStreamTargetTests {
StringWriter writer = new StringWriter();
CharacterStreamTarget target = new CharacterStreamTarget(writer);
target.setShouldAppendNewLine(true);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
task.run();
@@ -82,8 +81,8 @@ public class CharacterStreamTargetTests {
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
dispatcherPolicy.setMaxMessagesPerTask(2);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
task.run();
@@ -98,8 +97,8 @@ public class CharacterStreamTargetTests {
dispatcherPolicy.setMaxMessagesPerTask(10);
dispatcherPolicy.setReceiveTimeout(0);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
target.setShouldAppendNewLine(true);
channel.send(new StringMessage("foo"), 0);
channel.send(new StringMessage("bar"), 0);
@@ -113,8 +112,8 @@ public class CharacterStreamTargetTests {
MessageChannel channel = new QueueChannel();
StringWriter writer = new StringWriter();
CharacterStreamTarget target = new CharacterStreamTarget(writer);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
TestObject testObject = new TestObject("foo");
channel.send(new GenericMessage<TestObject>(testObject));
task.run();
@@ -129,8 +128,8 @@ public class CharacterStreamTargetTests {
dispatcherPolicy.setReceiveTimeout(0);
dispatcherPolicy.setMaxMessagesPerTask(2);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
TestObject testObject1 = new TestObject("foo");
TestObject testObject2 = new TestObject("bar");
channel.send(new GenericMessage<TestObject>(testObject1), 0);
@@ -148,8 +147,8 @@ public class CharacterStreamTargetTests {
dispatcherPolicy.setMaxMessagesPerTask(2);
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
target.setShouldAppendNewLine(true);
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
task.getDispatcher().subscribe(target);
PollingDispatcherTask task = new PollingDispatcherTask(channel, null);
task.subscribe(target);
TestObject testObject1 = new TestObject("foo");
TestObject testObject2 = new TestObject("bar");
channel.send(new GenericMessage<TestObject>(testObject1), 0);