From 420e0002c27e3b8cf73b8812641a7a9aaf5dc574 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Sat, 12 Jan 2008 18:07:38 +0000 Subject: [PATCH] Implemented CharacterStreamTargetAdapter and added CharacterStreamTargetAdapterTests. --- .../adapter/AbstractTargetAdapter.java | 5 + .../stream/CharacterStreamSourceAdapter.java | 2 +- .../stream/CharacterStreamTargetAdapter.java | 103 ++++++++ .../CharacterStreamTargetAdapterTests.java | 220 ++++++++++++++++++ 4 files changed, 329 insertions(+), 1 deletion(-) create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapter.java create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapterTests.java diff --git a/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java index 825f8861d4..08c5267c11 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java @@ -16,6 +16,9 @@ package org.springframework.integration.adapter; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + import org.springframework.integration.bus.ConsumerPolicy; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.message.Message; @@ -30,6 +33,8 @@ import org.springframework.util.Assert; */ public abstract class AbstractTargetAdapter implements TargetAdapter { + protected Log logger = LogFactory.getLog(this.getClass()); + private String name; private MessageChannel channel; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapter.java index 4ecdc6d787..69cf8621cf 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapter.java @@ -34,7 +34,7 @@ public class CharacterStreamSourceAdapter extends PollingSourceAdapter { /** - * Factory method for creating an adapter for stdin (System.in). + * Factory method that creates an adapter for stdin (System.in). */ public static CharacterStreamSourceAdapter stdinAdapter(MessageChannel channel) { CharacterStreamSourceAdapter adapter = new CharacterStreamSourceAdapter(System.in); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapter.java new file mode 100644 index 0000000000..f9305d945a --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapter.java @@ -0,0 +1,103 @@ +/* + * Copyright 2002-2007 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.adapter.stream; + +import java.io.BufferedWriter; +import java.io.IOException; +import java.io.OutputStream; +import java.io.OutputStreamWriter; + +import org.springframework.integration.MessageHandlingException; +import org.springframework.integration.adapter.AbstractTargetAdapter; + +/** + * A target adapter that writes to an {@link OutputStream}. 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 true. It is false + * by default. + * + * @author Mark Fisher + */ +public class CharacterStreamTargetAdapter extends AbstractTargetAdapter { + + private BufferedWriter writer; + + private boolean shouldAppendNewLine = false; + + + public CharacterStreamTargetAdapter(OutputStream stream) { + this(stream, -1); + } + + public CharacterStreamTargetAdapter(OutputStream stream, int bufferSize) { + if (bufferSize > 0) { + this.writer = new BufferedWriter(new OutputStreamWriter(stream), bufferSize); + } + else { + this.writer = new BufferedWriter(new OutputStreamWriter(stream)); + } + } + + + /** + * Factory method that creates an adapter for stdout (System.out). + */ + public static CharacterStreamTargetAdapter stdoutAdapter() { + return new CharacterStreamTargetAdapter(System.out); + } + + /** + * Factory method that creates an adapter for stderr (System.err). + */ + public static CharacterStreamTargetAdapter stderrAdapter() { + return new CharacterStreamTargetAdapter(System.err); + } + + + public void setShouldAppendNewLine(boolean shouldAppendNewLine) { + this.shouldAppendNewLine = shouldAppendNewLine; + } + + @Override + protected boolean sendToTarget(Object object) { + if (object == null) { + if (logger.isWarnEnabled()) { + logger.warn("target adapter received null object"); + } + return false; + } + try { + if (object instanceof String) { + writer.write((String) object); + } + else { + writer.write(object.toString()); + } + if (this.shouldAppendNewLine) { + writer.newLine(); + } + writer.flush(); + return true; + } + catch (IOException e) { + throw new MessageHandlingException("IO failure occurred in adapter", e); + } + } + +} diff --git a/spring-integration-core/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapterTests.java b/spring-integration-core/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapterTests.java new file mode 100644 index 0000000000..6aeda1affd --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapterTests.java @@ -0,0 +1,220 @@ +/* + * Copyright 2002-2007 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.adapter.stream; + +import static org.junit.Assert.assertEquals; + +import java.io.ByteArrayOutputStream; + +import org.junit.Test; + +import org.springframework.integration.bus.ChannelPollingMessageRetriever; +import org.springframework.integration.bus.ConsumerPolicy; +import org.springframework.integration.bus.DefaultMessageDispatcher; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.channel.SimpleChannel; +import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.message.StringMessage; + +/** + * @author Mark Fisher + */ +public class CharacterStreamTargetAdapterTests { + + @Test + public void testSingleString() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + channel.send(new StringMessage("foo")); + int count = dispatcher.dispatch(); + assertEquals(1, count); + String result = new String(stream.toByteArray()); + assertEquals("foo", result); + } + + @Test + public void testTwoStringsAndNoNewLinesByDefault() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + channel.send(new StringMessage("foo")); + channel.send(new StringMessage("bar")); + assertEquals(1, dispatcher.dispatch()); + String result1 = new String(stream.toByteArray()); + assertEquals("foo", result1); + assertEquals(1, dispatcher.dispatch()); + String result2 = new String(stream.toByteArray()); + assertEquals("foobar", result2); + } + + @Test + public void testTwoStringsWithNewLines() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + adapter.setShouldAppendNewLine(true); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + channel.send(new StringMessage("foo")); + channel.send(new StringMessage("bar")); + assertEquals(1, dispatcher.dispatch()); + String result1 = new String(stream.toByteArray()); + String newLine = System.getProperty("line.separator"); + assertEquals("foo" + newLine, result1); + assertEquals(1, dispatcher.dispatch()); + String result2 = new String(stream.toByteArray()); + assertEquals("foo" + newLine + "bar" + newLine, result2); + } + + @Test + public void testMaxMessagesPerTaskSameAsMessageCount() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + policy.setMaxMessagesPerTask(2); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + channel.send(new StringMessage("foo")); + channel.send(new StringMessage("bar")); + assertEquals(2, dispatcher.dispatch()); + String result = new String(stream.toByteArray()); + assertEquals("foobar", result); + } + + @Test + public void testMaxMessagesPerTaskExceedsMessageCountWithAppendedNewLines() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + adapter.setShouldAppendNewLine(true); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + policy.setReceiveTimeout(0); + policy.setMaxMessagesPerTask(10); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + channel.send(new StringMessage("foo")); + channel.send(new StringMessage("bar")); + assertEquals(2, dispatcher.dispatch()); + String result = new String(stream.toByteArray()); + String newLine = System.getProperty("line.separator"); + assertEquals("foo" + newLine + "bar" + newLine, result); + } + + @Test + public void testSingleNonStringObject() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + TestObject testObject = new TestObject("foo"); + channel.send(new GenericMessage(testObject)); + int count = dispatcher.dispatch(); + assertEquals(1, count); + String result = new String(stream.toByteArray()); + assertEquals("foo", result); + } + + @Test + public void testTwoNonStringObjectWithOutNewLines() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + policy.setReceiveTimeout(0); + policy.setMaxMessagesPerTask(2); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + TestObject testObject1 = new TestObject("foo"); + TestObject testObject2 = new TestObject("bar"); + channel.send(new GenericMessage(testObject1)); + channel.send(new GenericMessage(testObject2)); + assertEquals(2, dispatcher.dispatch()); + String result = new String(stream.toByteArray()); + assertEquals("foobar", result); + } + + @Test + public void testTwoNonStringObjectWithNewLines() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + MessageChannel channel = new SimpleChannel(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.setChannel(channel); + adapter.setShouldAppendNewLine(true); + ConsumerPolicy policy = ConsumerPolicy.newEventDrivenPolicy(); + policy.setReceiveTimeout(0); + policy.setMaxMessagesPerTask(2); + ChannelPollingMessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); + dispatcher.addHandler(adapter); + dispatcher.start(); + TestObject testObject1 = new TestObject("foo"); + TestObject testObject2 = new TestObject("bar"); + channel.send(new GenericMessage(testObject1)); + channel.send(new GenericMessage(testObject2)); + assertEquals(2, dispatcher.dispatch()); + String result = new String(stream.toByteArray()); + String newLine = System.getProperty("line.separator"); + assertEquals("foo" + newLine + "bar" + newLine, result); + } + + + private static class TestObject { + + private String text; + + TestObject(String text) { + this.text = text; + } + + public String toString() { + return this.text; + } + } + +}