ByteStreamTargetAdapter is now ByteStreamTarget. CharacterStreamTargetAdapter is now CharacterStreamTarget.
This commit is contained in:
@@ -28,22 +28,22 @@ import org.springframework.integration.message.MessagingException;
|
||||
import org.springframework.integration.message.Target;
|
||||
|
||||
/**
|
||||
* A target adapter that writes a byte array to an {@link OutputStream}.
|
||||
* A target that writes a byte array to an {@link OutputStream}.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class ByteStreamTargetAdapter implements Target {
|
||||
public class ByteStreamTarget implements Target {
|
||||
|
||||
private final Log logger = LogFactory.getLog(this.getClass());
|
||||
|
||||
private BufferedOutputStream stream;
|
||||
private final BufferedOutputStream stream;
|
||||
|
||||
|
||||
public ByteStreamTargetAdapter(OutputStream stream) {
|
||||
public ByteStreamTarget(OutputStream stream) {
|
||||
this(stream, -1);
|
||||
}
|
||||
|
||||
public ByteStreamTargetAdapter(OutputStream stream, int bufferSize) {
|
||||
public ByteStreamTarget(OutputStream stream, int bufferSize) {
|
||||
if (bufferSize > 0) {
|
||||
this.stream = new BufferedOutputStream(stream, bufferSize);
|
||||
}
|
||||
@@ -52,6 +52,7 @@ public class ByteStreamTargetAdapter implements Target {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public boolean send(Message message) {
|
||||
Object payload = message.getPayload();
|
||||
if (payload == null) {
|
||||
@@ -75,7 +76,7 @@ public class ByteStreamTargetAdapter implements Target {
|
||||
return true;
|
||||
}
|
||||
catch (IOException e) {
|
||||
throw new MessagingException("IO failure occurred in adapter", e);
|
||||
throw new MessagingException("IO failure occurred in target", e);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -33,16 +33,15 @@ import org.springframework.integration.message.Target;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* A target adapter that writes to a {@link Writer}. String-based
|
||||
* objects will be written directly, but if the object is not itself a
|
||||
* {@link String}, the adapter will write the result of the object's
|
||||
* {@link #toString()} method. To append a new-line after each write, set the
|
||||
* {@link #shouldAppendNewLine} flag to <em>true</em>. It is <em>false</em>
|
||||
* by default.
|
||||
* A target that writes to a {@link Writer}. String-based objects will be
|
||||
* written directly, but if the object is not itself a {@link String}, the
|
||||
* target will write the result of the object's {@link #toString()} method.
|
||||
* To append a new-line after each write, set the {@link #shouldAppendNewLine}
|
||||
* flag to <em>true</em>. It is <em>false</em> by default.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class CharacterStreamTargetAdapter implements Target {
|
||||
public class CharacterStreamTarget implements Target {
|
||||
|
||||
private final Log logger = LogFactory.getLog(this.getClass());
|
||||
|
||||
@@ -51,11 +50,11 @@ public class CharacterStreamTargetAdapter implements Target {
|
||||
private volatile boolean shouldAppendNewLine = false;
|
||||
|
||||
|
||||
public CharacterStreamTargetAdapter(Writer writer) {
|
||||
public CharacterStreamTarget(Writer writer) {
|
||||
this(writer, -1);
|
||||
}
|
||||
|
||||
public CharacterStreamTargetAdapter(Writer writer, int bufferSize) {
|
||||
public CharacterStreamTarget(Writer writer, int bufferSize) {
|
||||
Assert.notNull(writer, "writer must not be null");
|
||||
if (writer instanceof BufferedWriter) {
|
||||
this.writer = (BufferedWriter) writer;
|
||||
@@ -70,43 +69,43 @@ public class CharacterStreamTargetAdapter implements Target {
|
||||
|
||||
|
||||
/**
|
||||
* Factory method that creates an adapter for stdout (System.out) with the
|
||||
* Factory method that creates a target for stdout (System.out) with the
|
||||
* default charset encoding.
|
||||
*/
|
||||
public static CharacterStreamTargetAdapter stdoutAdapter() {
|
||||
return stdoutAdapter(null);
|
||||
public static CharacterStreamTarget stdout() {
|
||||
return stdout(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Factory method that creates an adapter for stdout (System.out) with the
|
||||
* Factory method that creates a target for stdout (System.out) with the
|
||||
* specified charset encoding.
|
||||
*/
|
||||
public static CharacterStreamTargetAdapter stdoutAdapter(String charsetName) {
|
||||
return createAdapterForStream(System.out, charsetName);
|
||||
public static CharacterStreamTarget stdout(String charsetName) {
|
||||
return createTargetForStream(System.out, charsetName);
|
||||
}
|
||||
|
||||
/**
|
||||
* Factory method that creates an adapter for stderr (System.err) with the
|
||||
* Factory method that creates a target for stderr (System.err) with the
|
||||
* default charset encoding.
|
||||
*/
|
||||
public static CharacterStreamTargetAdapter stderrAdapter() {
|
||||
return stderrAdapter(null);
|
||||
public static CharacterStreamTarget stderr() {
|
||||
return stderr(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Factory method that creates an adapter for stderr (System.err) with the
|
||||
* Factory method that creates a target for stderr (System.err) with the
|
||||
* specified charset encoding.
|
||||
*/
|
||||
public static CharacterStreamTargetAdapter stderrAdapter(String charsetName) {
|
||||
return createAdapterForStream(System.err, charsetName);
|
||||
public static CharacterStreamTarget stderr(String charsetName) {
|
||||
return createTargetForStream(System.err, charsetName);
|
||||
}
|
||||
|
||||
private static CharacterStreamTargetAdapter createAdapterForStream(OutputStream stream, String charsetName) {
|
||||
private static CharacterStreamTarget createTargetForStream(OutputStream stream, String charsetName) {
|
||||
if (charsetName == null) {
|
||||
return new CharacterStreamTargetAdapter(new OutputStreamWriter(stream));
|
||||
return new CharacterStreamTarget(new OutputStreamWriter(stream));
|
||||
}
|
||||
try {
|
||||
return new CharacterStreamTargetAdapter(new OutputStreamWriter(stream, charsetName));
|
||||
return new CharacterStreamTarget(new OutputStreamWriter(stream, charsetName));
|
||||
}
|
||||
catch (UnsupportedEncodingException e) {
|
||||
throw new ConfigurationException("unsupported encoding: " + charsetName, e);
|
||||
@@ -122,7 +121,7 @@ public class CharacterStreamTargetAdapter implements Target {
|
||||
Object payload = message.getPayload();
|
||||
if (payload == null) {
|
||||
if (logger.isWarnEnabled()) {
|
||||
logger.warn("target adapter received null payload");
|
||||
logger.warn("target received null payload");
|
||||
}
|
||||
return false;
|
||||
}
|
||||
@@ -146,7 +145,7 @@ public class CharacterStreamTargetAdapter implements Target {
|
||||
return true;
|
||||
}
|
||||
catch (IOException e) {
|
||||
throw new MessagingException("IO failure occurred in adapter", e);
|
||||
throw new MessagingException("IO failure occurred in target", e);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -20,7 +20,7 @@ import org.w3c.dom.Element;
|
||||
|
||||
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
|
||||
import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser;
|
||||
import org.springframework.integration.adapter.stream.CharacterStreamTargetAdapter;
|
||||
import org.springframework.integration.adapter.stream.CharacterStreamTarget;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
@@ -32,7 +32,7 @@ public class ConsoleTargetParser extends AbstractSingleBeanDefinitionParser {
|
||||
|
||||
@Override
|
||||
protected Class<?> getBeanClass(Element element) {
|
||||
return CharacterStreamTargetAdapter.class;
|
||||
return CharacterStreamTarget.class;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -48,10 +48,10 @@ public class ConsoleTargetParser extends AbstractSingleBeanDefinitionParser {
|
||||
@Override
|
||||
protected void doParse(Element element, BeanDefinitionBuilder builder) {
|
||||
if ("true".equals(element.getAttribute("error"))) {
|
||||
builder.setFactoryMethod("stderrAdapter");
|
||||
builder.setFactoryMethod("stderr");
|
||||
}
|
||||
else {
|
||||
builder.setFactoryMethod("stdoutAdapter");
|
||||
builder.setFactoryMethod("stdout");
|
||||
}
|
||||
String charsetName = element.getAttribute("charset");
|
||||
if (StringUtils.hasText(charsetName)) {
|
||||
|
||||
@@ -32,7 +32,7 @@ import org.springframework.integration.scheduling.PollingSchedule;
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class ByteStreamSourceAdapterTests {
|
||||
public class ByteStreamSourceTests {
|
||||
|
||||
@Test
|
||||
public void testEndOfStream() {
|
||||
@@ -33,13 +33,13 @@ import org.springframework.integration.message.StringMessage;
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class ByteStreamTargetAdapterTests {
|
||||
public class ByteStreamTargetTests {
|
||||
|
||||
@Test
|
||||
public void testSingleByteArray() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
adapter.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
target.send(new GenericMessage<byte[]>(new byte[] {1,2,3}));
|
||||
byte[] result = stream.toByteArray();
|
||||
assertEquals(3, result.length);
|
||||
assertEquals(1, result[0]);
|
||||
@@ -50,8 +50,8 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testSingleString() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
adapter.send(new StringMessage("foo"));
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
target.send(new StringMessage("foo"));
|
||||
byte[] result = stream.toByteArray();
|
||||
assertEquals(3, result.length);
|
||||
assertEquals("foo", new String(result));
|
||||
@@ -60,12 +60,12 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testMaxMessagesPerTaskSameAsMessageCount() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(3);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().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);
|
||||
@@ -79,12 +79,12 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testMaxMessagesPerTaskLessThanMessageCount() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().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);
|
||||
@@ -97,13 +97,13 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testMaxMessagesPerTaskExceedsMessageCount() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(5);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().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);
|
||||
@@ -116,13 +116,13 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testMaxMessagesLessThanMessageCountWithMultipleDispatches() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().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);
|
||||
@@ -140,13 +140,13 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testMaxMessagesExceedsMessageCountWithMultipleDispatches() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(5);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().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);
|
||||
@@ -163,13 +163,13 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testStreamResetBetweenDispatches() {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().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);
|
||||
@@ -186,13 +186,13 @@ public class ByteStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testStreamWriteBetweenDispatches() throws IOException {
|
||||
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||
ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream);
|
||||
ByteStreamTarget target = new ByteStreamTarget(stream);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().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);
|
||||
@@ -32,7 +32,7 @@ import org.springframework.integration.scheduling.PollingSchedule;
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class CharacterStreamSourceAdapterTests {
|
||||
public class CharacterStreamSourceTests {
|
||||
|
||||
@Test
|
||||
public void testEndOfStream() {
|
||||
@@ -33,13 +33,13 @@ import org.springframework.integration.message.StringMessage;
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class CharacterStreamTargetAdapterTests {
|
||||
public class CharacterStreamTargetTests {
|
||||
|
||||
@Test
|
||||
public void testSingleString() {
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
adapter.send(new StringMessage("foo"));
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
target.send(new StringMessage("foo"));
|
||||
assertEquals("foo", writer.toString());
|
||||
}
|
||||
|
||||
@@ -47,9 +47,9 @@ public class CharacterStreamTargetAdapterTests {
|
||||
public void testTwoStringsAndNoNewLinesByDefault() {
|
||||
MessageChannel channel = new QueueChannel();
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().subscribe(target);
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
task.run();
|
||||
@@ -62,10 +62,10 @@ public class CharacterStreamTargetAdapterTests {
|
||||
public void testTwoStringsWithNewLines() {
|
||||
MessageChannel channel = new QueueChannel();
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
target.setShouldAppendNewLine(true);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().subscribe(target);
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
task.run();
|
||||
@@ -78,12 +78,12 @@ public class CharacterStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testMaxMessagesPerTaskSameAsMessageCount() {
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().subscribe(target);
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
task.run();
|
||||
@@ -93,14 +93,14 @@ public class CharacterStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testMaxMessagesPerTaskExceedsMessageCountWithAppendedNewLines() {
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setMaxMessagesPerTask(10);
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
task.getDispatcher().subscribe(target);
|
||||
target.setShouldAppendNewLine(true);
|
||||
channel.send(new StringMessage("foo"), 0);
|
||||
channel.send(new StringMessage("bar"), 0);
|
||||
task.run();
|
||||
@@ -112,9 +112,9 @@ public class CharacterStreamTargetAdapterTests {
|
||||
public void testSingleNonStringObject() {
|
||||
MessageChannel channel = new QueueChannel();
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().subscribe(target);
|
||||
TestObject testObject = new TestObject("foo");
|
||||
channel.send(new GenericMessage<TestObject>(testObject));
|
||||
task.run();
|
||||
@@ -124,13 +124,13 @@ public class CharacterStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testTwoNonStringObjectWithOutNewLines() {
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().subscribe(target);
|
||||
TestObject testObject1 = new TestObject("foo");
|
||||
TestObject testObject2 = new TestObject("bar");
|
||||
channel.send(new GenericMessage<TestObject>(testObject1), 0);
|
||||
@@ -142,14 +142,14 @@ public class CharacterStreamTargetAdapterTests {
|
||||
@Test
|
||||
public void testTwoNonStringObjectWithNewLines() {
|
||||
StringWriter writer = new StringWriter();
|
||||
CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(writer);
|
||||
CharacterStreamTarget target = new CharacterStreamTarget(writer);
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setReceiveTimeout(0);
|
||||
dispatcherPolicy.setMaxMessagesPerTask(2);
|
||||
QueueChannel channel = new QueueChannel(5, dispatcherPolicy);
|
||||
adapter.setShouldAppendNewLine(true);
|
||||
target.setShouldAppendNewLine(true);
|
||||
PollingDispatcherTask task = new PollingDispatcherTask(new DefaultPollingDispatcher(channel), null);
|
||||
task.getDispatcher().subscribe(adapter);
|
||||
task.getDispatcher().subscribe(target);
|
||||
TestObject testObject1 = new TestObject("foo");
|
||||
TestObject testObject2 = new TestObject("bar");
|
||||
channel.send(new GenericMessage<TestObject>(testObject1), 0);
|
||||
@@ -33,7 +33,7 @@ import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.springframework.beans.factory.BeanCreationException;
|
||||
import org.springframework.context.support.ClassPathXmlApplicationContext;
|
||||
import org.springframework.integration.ConfigurationException;
|
||||
import org.springframework.integration.adapter.stream.CharacterStreamTargetAdapter;
|
||||
import org.springframework.integration.adapter.stream.CharacterStreamTarget;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
|
||||
/**
|
||||
@@ -61,8 +61,8 @@ public class ConsoleTargetParserTests {
|
||||
public void testConsoleTargetWithDefaultCharset() {
|
||||
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
|
||||
"consoleTargetParserTests.xml", ConsoleTargetParserTests.class);
|
||||
CharacterStreamTargetAdapter target =
|
||||
(CharacterStreamTargetAdapter) context.getBean("targetWithDefaultCharset");
|
||||
CharacterStreamTarget target =
|
||||
(CharacterStreamTarget) context.getBean("targetWithDefaultCharset");
|
||||
DirectFieldAccessor targetAccessor = new DirectFieldAccessor(target);
|
||||
Writer bufferedWriter = (Writer) targetAccessor.getPropertyValue("writer");
|
||||
assertEquals(BufferedWriter.class, bufferedWriter.getClass());
|
||||
@@ -81,8 +81,8 @@ public class ConsoleTargetParserTests {
|
||||
public void testConsoleTargetWithProvidedCharset() {
|
||||
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
|
||||
"consoleTargetParserTests.xml", ConsoleTargetParserTests.class);
|
||||
CharacterStreamTargetAdapter target =
|
||||
(CharacterStreamTargetAdapter) context.getBean("targetWithProvidedCharset");
|
||||
CharacterStreamTarget target =
|
||||
(CharacterStreamTarget) context.getBean("targetWithProvidedCharset");
|
||||
DirectFieldAccessor targetAccessor = new DirectFieldAccessor(target);
|
||||
Writer bufferedWriter = (Writer) targetAccessor.getPropertyValue("writer");
|
||||
assertEquals(BufferedWriter.class, bufferedWriter.getClass());
|
||||
@@ -117,8 +117,8 @@ public class ConsoleTargetParserTests {
|
||||
public void testErrorTarget() {
|
||||
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
|
||||
"consoleTargetParserTests.xml", ConsoleTargetParserTests.class);
|
||||
CharacterStreamTargetAdapter target =
|
||||
(CharacterStreamTargetAdapter) context.getBean("stderrTarget");
|
||||
CharacterStreamTarget target =
|
||||
(CharacterStreamTarget) context.getBean("stderrTarget");
|
||||
DirectFieldAccessor targetAccessor = new DirectFieldAccessor(target);
|
||||
Writer bufferedWriter = (Writer) targetAccessor.getPropertyValue("writer");
|
||||
assertEquals(BufferedWriter.class, bufferedWriter.getClass());
|
||||
|
||||
Reference in New Issue
Block a user