diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/config/spring-integration-adapters-1.0.xsd b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/config/spring-integration-adapters-1.0.xsd new file mode 100644 index 0000000000..0bd0dce470 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/config/spring-integration-adapters-1.0.xsd @@ -0,0 +1,82 @@ + + + + + + + + + + + + + + + + + Defines a file-based source channel adapter. + + + + + + + + + + + + + + Defines a file-based target channel adapter. + + + + + + + + + + + + + Defines a jms-based source channel adapter. + + + + + + + + + + + + + + + + + + Defines a jms-based target channel adapter. + + + + + + + + + + + + \ No newline at end of file diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/ChannelPublishingJmsListener.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/ChannelPublishingJmsListener.java new file mode 100644 index 0000000000..70b84af3c0 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/ChannelPublishingJmsListener.java @@ -0,0 +1,96 @@ +/* + * 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.jms; + +import javax.jms.JMSException; +import javax.jms.MessageListener; + +import org.springframework.integration.MessagingConfigurationException; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.message.Message; +import org.springframework.integration.message.MessageDeliveryException; +import org.springframework.integration.message.MessageMapper; +import org.springframework.integration.message.SimplePayloadMessageMapper; +import org.springframework.jms.support.converter.MessageConverter; +import org.springframework.jms.support.converter.SimpleMessageConverter; +import org.springframework.util.Assert; + +/** + * JMS {@link MessageListener} implementation that converts the received JMS + * message into a Spring Integration message and then sends that to a channel. + * + * @author Mark Fisher + */ +public class ChannelPublishingJmsListener implements MessageListener { + + private MessageChannel channel; + + private long timeout = -1; + + private MessageConverter converter = new SimpleMessageConverter(); + + private MessageMapper mapper = new SimplePayloadMessageMapper(); + + + public ChannelPublishingJmsListener() { + } + + public ChannelPublishingJmsListener(MessageChannel channel) { + Assert.notNull(channel, "'channel' must not be null"); + this.channel = channel; + } + + + public void setChannel(MessageChannel channel) { + Assert.notNull(channel, "'channel' must not be null"); + this.channel = channel; + } + + public void setTimeout(long timeout) { + this.timeout = timeout; + } + + public void setMessageConverter(MessageConverter messageConverter) { + Assert.notNull(messageConverter, "'messageConverter' must not be null"); + this.converter = messageConverter; + } + + public void setMessageMapper(MessageMapper messageMapper) { + Assert.notNull(messageMapper, "'messageMapper' must not be null"); + this.mapper = messageMapper; + } + + public void onMessage(javax.jms.Message jmsMessage) { + if (this.channel == null) { + throw new MessagingConfigurationException("'channel' must not be null"); + } + try { + Object payload = converter.fromMessage(jmsMessage); + Message messageToSend = mapper.toMessage(payload); + if (this.timeout < 0) { + this.channel.send(messageToSend); + } + else { + this.channel.send(messageToSend, timeout); + } + } + catch (JMSException e) { + throw new MessageDeliveryException("failed to convert JMS Message", e); + } + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/DefaultJmsMessagePostProcessor.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/DefaultJmsMessagePostProcessor.java new file mode 100644 index 0000000000..19151ca83e --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/DefaultJmsMessagePostProcessor.java @@ -0,0 +1,48 @@ +/* + * 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.jms; + +import javax.jms.Destination; +import javax.jms.JMSException; +import javax.jms.Message; + +import org.springframework.integration.message.MessageHeader; + +/** + * A {@link JmsMessagePostProcessor} that passes attributes from a + * {@link MessageHeader} to a JMS message before it is sent to its destination. + * + * @author Mark Fisher + */ +public class DefaultJmsMessagePostProcessor implements JmsMessagePostProcessor { + + public void postProcessJmsMessage(Message jmsMessage, MessageHeader header) throws JMSException { + Object jmsCorrelationId = header.getAttribute(JmsTargetAdapter.JMS_CORRELATION_ID); + if (jmsCorrelationId != null && (jmsCorrelationId instanceof String)) { + jmsMessage.setJMSCorrelationID((String) jmsCorrelationId); + } + Object jmsReplyTo = header.getAttribute(JmsTargetAdapter.JMS_REPLY_TO); + if (jmsReplyTo != null && (jmsReplyTo instanceof Destination)) { + jmsMessage.setJMSReplyTo((Destination) jmsReplyTo); + } + Object jmsType = header.getAttribute(JmsTargetAdapter.JMS_TYPE); + if (jmsType != null && (jmsType instanceof String)) { + jmsMessage.setJMSType((String) jmsType); + } + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsMessageDrivenSourceAdapter.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsMessageDrivenSourceAdapter.java new file mode 100644 index 0000000000..2eeeaf1337 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsMessageDrivenSourceAdapter.java @@ -0,0 +1,149 @@ +/* + * 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.jms; + +import javax.jms.ConnectionFactory; +import javax.jms.Destination; + +import org.springframework.beans.factory.DisposableBean; +import org.springframework.context.Lifecycle; +import org.springframework.core.task.TaskExecutor; +import org.springframework.integration.MessagingConfigurationException; +import org.springframework.integration.adapter.AbstractSourceAdapter; +import org.springframework.jms.listener.AbstractJmsListeningContainer; +import org.springframework.jms.listener.DefaultMessageListenerContainer; +import org.springframework.jms.support.converter.MessageConverter; +import org.springframework.jms.support.converter.SimpleMessageConverter; +import org.springframework.util.Assert; + +/** + * A message-driven adapter for receiving JMS messages and sending to a channel. + * + * @author Mark Fisher + */ +public class JmsMessageDrivenSourceAdapter extends AbstractSourceAdapter implements Lifecycle, DisposableBean { + + private AbstractJmsListeningContainer container; + + private ConnectionFactory connectionFactory; + + private Destination destination; + + private String destinationName; + + private MessageConverter messageConverter = new SimpleMessageConverter(); + + private TaskExecutor taskExecutor; + + private long receiveTimeout = 1000; + + private int concurrentConsumers = 1; + + private int maxConcurrentConsumers = 1; + + private int maxMessagesPerTask = Integer.MIN_VALUE; + + private int idleTaskExecutionLimit = 1; + + private long sendTimeout = -1; + + + public void setContainer(AbstractJmsListeningContainer container) { + this.container = container; + } + + public void setConnectionFactory(ConnectionFactory connectionFactory) { + this.connectionFactory = connectionFactory; + } + + public void setDestination(Destination destination) { + this.destination = destination; + } + + public void setDestinationName(String destinationName) { + this.destinationName = destinationName; + } + + public void setMessageConverter(MessageConverter messageConverter) { + Assert.notNull(messageConverter, "'messageConverter' must not be null"); + this.messageConverter = messageConverter; + } + + public void setSendTimeout(long sendTimeout) { + this.sendTimeout = sendTimeout; + } + + public void setTaskExecutor(TaskExecutor taskExecutor) { + this.taskExecutor = taskExecutor; + } + + @Override + public void initialize() { + if (this.container == null) { + initDefaultContainer(); + } + } + + private void initDefaultContainer() { + if (this.connectionFactory == null || (this.destination == null && this.destinationName == null)) { + throw new MessagingConfigurationException("If a 'container' reference is not provided, then " + + "'connectionFactory' and 'destination' (or 'destinationName') are required."); + } + DefaultMessageListenerContainer dmlc = new DefaultMessageListenerContainer(); + dmlc.setConnectionFactory(this.connectionFactory); + if (this.destination != null) { + dmlc.setDestination(this.destination); + } + if (this.destinationName != null) { + dmlc.setDestinationName(this.destinationName); + } + dmlc.setReceiveTimeout(this.receiveTimeout); + dmlc.setConcurrentConsumers(this.concurrentConsumers); + dmlc.setMaxConcurrentConsumers(this.maxConcurrentConsumers); + dmlc.setMaxMessagesPerTask(this.maxMessagesPerTask); + dmlc.setIdleTaskExecutionLimit(this.idleTaskExecutionLimit); + dmlc.setSessionTransacted(true); + dmlc.setAutoStartup(false); + ChannelPublishingJmsListener listener = new ChannelPublishingJmsListener(this.getChannel()); + listener.setMessageConverter(this.messageConverter); + listener.setMessageMapper(this.getMessageMapper()); + listener.setTimeout(this.sendTimeout); + dmlc.setMessageListener(listener); + if (this.taskExecutor != null) { + dmlc.setTaskExecutor(this.taskExecutor); + } + dmlc.afterPropertiesSet(); + this.container = dmlc; + } + + public boolean isRunning() { + return container.isRunning(); + } + + public void start() { + container.start(); + } + + public void stop() { + container.stop(); + } + + public void destroy() { + container.destroy(); + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsMessagePostProcessor.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsMessagePostProcessor.java new file mode 100644 index 0000000000..2b98a66757 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsMessagePostProcessor.java @@ -0,0 +1,34 @@ +/* + * 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.jms; + +import javax.jms.JMSException; +import javax.jms.Message; + +import org.springframework.integration.message.MessageHeader; + +/** + * Strategy interface for post-processing a JMS Message before it is sent to its + * destination. + * + * @author Mark Fisher + */ +public interface JmsMessagePostProcessor { + + void postProcessJmsMessage(Message jmsMessage, MessageHeader header) throws JMSException; + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsPollableSource.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsPollableSource.java new file mode 100644 index 0000000000..c2cbf0a1bf --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsPollableSource.java @@ -0,0 +1,116 @@ +/* + * 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.jms; + +import java.util.Arrays; +import java.util.Collection; + +import javax.jms.ConnectionFactory; +import javax.jms.Destination; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.integration.MessagingConfigurationException; +import org.springframework.integration.adapter.PollableSource; +import org.springframework.jms.core.JmsTemplate; + +/** + * A source for receiving JMS Messages with a polling listener. This source is + * only recommended for very low message volume. Otherwise, the + * {@link JmsMessageDrivenSourceAdapter} that uses Spring's MessageListener + * container support is highly recommended. + * + * @author Mark Fisher + */ +public class JmsPollableSource implements PollableSource, InitializingBean { + + private ConnectionFactory connectionFactory; + + private Destination destination; + + private String destinationName; + + private JmsTemplate jmsTemplate; + + + public JmsPollableSource(JmsTemplate jmsTemplate) { + this.jmsTemplate = jmsTemplate; + } + + public JmsPollableSource(ConnectionFactory connectionFactory, Destination destination) { + this.connectionFactory = connectionFactory; + this.destination = destination; + this.initJmsTemplate(); + } + + public JmsPollableSource(ConnectionFactory connectionFactory, String destinationName) { + this.connectionFactory = connectionFactory; + this.destinationName = destinationName; + this.initJmsTemplate(); + } + + /** + * No-arg constructor provided for convenience when configuring with + * setters. Note that the initialization callback will validate. + */ + public JmsPollableSource() { + } + + public void setConnectionFactory(ConnectionFactory connectionFactory) { + this.connectionFactory = connectionFactory; + } + + public void setDestination(Destination destination) { + this.destination = destination; + } + + public void setDestinationName(String destinationName) { + this.destinationName = destinationName; + } + + public void setJmsTemplate(JmsTemplate jmsTemplate) { + this.jmsTemplate = jmsTemplate; + } + + public void afterPropertiesSet() { + if (this.jmsTemplate == null) { + if (this.connectionFactory == null || (this.destination == null && this.destinationName == null)) { + throw new MessagingConfigurationException("Either a 'jmsTemplate' or " + + "both 'connectionFactory' and 'destination' (or 'destinationName') are required."); + } + this.initJmsTemplate(); + } + } + + private void initJmsTemplate() { + this.jmsTemplate = new JmsTemplate(); + this.jmsTemplate.setConnectionFactory(this.connectionFactory); + if (this.destination != null) { + this.jmsTemplate.setDefaultDestination(this.destination); + } + else if (this.destinationName != null) { + this.jmsTemplate.setDefaultDestinationName(this.destinationName); + } + else { + throw new MessagingConfigurationException("either 'destination' or 'destinationName' is required"); + } + } + + public Collection poll(int limit) { + return Arrays.asList(this.jmsTemplate.receiveAndConvert()); + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsPollingSourceAdapter.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsPollingSourceAdapter.java new file mode 100644 index 0000000000..32e18ea39b --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsPollingSourceAdapter.java @@ -0,0 +1,44 @@ +/* + * 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.jms; + +import javax.jms.ConnectionFactory; +import javax.jms.Destination; + +import org.springframework.integration.adapter.PollingSourceAdapter; +import org.springframework.jms.core.JmsTemplate; + +/** + * A convenience adapter that wraps a {@link JmsPollableSource}. + * + * @author Mark Fisher + */ +public class JmsPollingSourceAdapter extends PollingSourceAdapter { + + public JmsPollingSourceAdapter(JmsTemplate jmsTemplate) { + super(new JmsPollableSource(jmsTemplate)); + } + + public JmsPollingSourceAdapter(ConnectionFactory connectionFactory, Destination destination) { + super(new JmsPollableSource(connectionFactory, destination)); + } + + public JmsPollingSourceAdapter(ConnectionFactory connectionFactory, String destinationName) { + super(new JmsPollableSource(connectionFactory, destinationName)); + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsTargetAdapter.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsTargetAdapter.java new file mode 100644 index 0000000000..28badc5174 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/JmsTargetAdapter.java @@ -0,0 +1,139 @@ +/* + * 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.jms; + +import javax.jms.ConnectionFactory; +import javax.jms.Destination; +import javax.jms.JMSException; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.integration.MessagingConfigurationException; +import org.springframework.integration.handler.MessageHandler; +import org.springframework.integration.message.Message; +import org.springframework.jms.core.JmsTemplate; +import org.springframework.jms.core.MessagePostProcessor; +import org.springframework.jms.support.converter.SimpleMessageConverter; + +/** + * A target adapter for sending JMS Messages. + * + * @author Mark Fisher + */ +public class JmsTargetAdapter implements MessageHandler, InitializingBean { + + public static final String JMS_CORRELATION_ID = "JMSCorrelationID"; + + public static final String JMS_REPLY_TO = "JMSReplyTo"; + + public static final String JMS_TYPE = "JMSType"; + + + private ConnectionFactory connectionFactory; + + private Destination destination; + + private String destinationName; + + private JmsTemplate jmsTemplate; + + private JmsMessagePostProcessor jmsMessagePostProcessor = new DefaultJmsMessagePostProcessor(); + + + public JmsTargetAdapter(JmsTemplate jmsTemplate) { + this.jmsTemplate = jmsTemplate; + } + + public JmsTargetAdapter(ConnectionFactory connectionFactory, Destination destination) { + this.connectionFactory = connectionFactory; + this.destination = destination; + this.initJmsTemplate(); + } + + public JmsTargetAdapter(ConnectionFactory connectionFactory, String destinationName) { + this.connectionFactory = connectionFactory; + this.destinationName = destinationName; + this.initJmsTemplate(); + } + + /** + * No-arg constructor provided for convenience when configuring with + * setters. Note that the initialization callback will validate. + */ + public JmsTargetAdapter() { + } + + + public void setConnectionFactory(ConnectionFactory connectionFactory) { + this.connectionFactory = connectionFactory; + } + + public void setDestination(Destination destination) { + this.destination = destination; + } + + public void setDestinationName(String destinationName) { + this.destinationName = destinationName; + } + + public void setJmsTemplate(JmsTemplate jmsTemplate) { + this.jmsTemplate = jmsTemplate; + } + + public void setJmsMessagePostProcessor(JmsMessagePostProcessor jmsMessagePostProcessor) { + this.jmsMessagePostProcessor = jmsMessagePostProcessor; + } + + public void afterPropertiesSet() { + if (this.jmsTemplate == null) { + if (this.connectionFactory == null || (this.destination == null && this.destinationName == null)) { + throw new MessagingConfigurationException("Either a 'jmsTemplate' or " + + "*both* 'connectionFactory' and 'destination' (or 'destination-name') are required."); + } + this.initJmsTemplate(); + } + if (this.jmsTemplate.getMessageConverter() == null) { + this.jmsTemplate.setMessageConverter(new SimpleMessageConverter()); + } + } + + private void initJmsTemplate() { + this.jmsTemplate = new JmsTemplate(); + this.jmsTemplate.setConnectionFactory(this.connectionFactory); + if (this.destination != null) { + this.jmsTemplate.setDefaultDestination(this.destination); + } + else { + this.jmsTemplate.setDefaultDestinationName(this.destinationName); + } + } + + public final Message handle(final Message message) { + if (message == null) { + return null; + } + this.jmsTemplate.convertAndSend(message.getPayload(), new MessagePostProcessor() { + public javax.jms.Message postProcessMessage(javax.jms.Message jmsMessage) throws JMSException { + if (jmsMessagePostProcessor != null) { + jmsMessagePostProcessor.postProcessJmsMessage(jmsMessage, message.getHeader()); + } + return jmsMessage; + } + }); + return null; + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/config/JmsSourceAdapterParser.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/config/JmsSourceAdapterParser.java new file mode 100644 index 0000000000..0510175c15 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/config/JmsSourceAdapterParser.java @@ -0,0 +1,158 @@ +/* + * 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.jms.config; + +import org.w3c.dom.Element; + +import org.springframework.beans.factory.BeanCreationException; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser; +import org.springframework.integration.adapter.jms.JmsMessageDrivenSourceAdapter; +import org.springframework.integration.adapter.jms.JmsPollingSourceAdapter; +import org.springframework.util.StringUtils; + +/** + * Parser for the <jms-source/> element. + * + * @author Mark Fisher + */ +public class JmsSourceAdapterParser extends AbstractSingleBeanDefinitionParser { + + private static final String JMS_TEMPLATE_ATTRIBUTE = "jms-template"; + + private static final String CONNECTION_FACTORY_ATTRIBUTE = "connection-factory"; + + private static final String CONNECTION_FACTORY_PROPERTY = "connectionFactory"; + + private static final String DESTINATION_ATTRIBUTE = "destination"; + + private static final String DESTINATION_PROPERTY = "destination"; + + private static final String DESTINATION_NAME_ATTRIBUTE = "destination-name"; + + private static final String DESTINATION_NAME_PROPERTY = "destinationName"; + + private static final String CHANNEL_ATTRIBUTE = "channel"; + + private static final String CHANNEL_PROPERTY = "channel"; + + private static final String POLL_PERIOD_ATTRIBUTE = "poll-period"; + + private static final String POLL_PERIOD_PROPERTY = "period"; + + private static final String MESSAGE_CONVERTER_ATTRIBUTE = "message-converter"; + + private static final String MESSAGE_CONVERTER_PROPERTY = "messageConverter"; + + + protected Class getBeanClass(Element element) { + if (StringUtils.hasText(element.getAttribute(POLL_PERIOD_ATTRIBUTE))) { + return JmsPollingSourceAdapter.class; + } + return JmsMessageDrivenSourceAdapter.class; + } + + protected boolean shouldGenerateId() { + return false; + } + + protected boolean shouldGenerateIdAsFallback() { + return true; + } + + protected void doParse(Element element, BeanDefinitionBuilder builder) { + if (builder.getBeanDefinition().getBeanClass().equals(JmsPollingSourceAdapter.class)) { + parsePollingSourceAdapter(element, builder); + } + else { + parseMessageDrivenSourceAdapter(element, builder); + } + String channel = element.getAttribute(CHANNEL_ATTRIBUTE); + builder.addPropertyReference(CHANNEL_PROPERTY, channel); + } + + private void parsePollingSourceAdapter(Element element, BeanDefinitionBuilder builder) { + String pollPeriod = element.getAttribute(POLL_PERIOD_ATTRIBUTE); + if (!StringUtils.hasText(pollPeriod)) { + throw new BeanCreationException("'" + POLL_PERIOD_ATTRIBUTE + + "' is required for a " + JmsPollingSourceAdapter.class.getSimpleName()); + } + if (StringUtils.hasText(element.getAttribute(MESSAGE_CONVERTER_ATTRIBUTE))) { + throw new BeanCreationException("The '" + MESSAGE_CONVERTER_ATTRIBUTE + "' attribute is not supported for a " + + JmsPollingSourceAdapter.class.getSimpleName() + ". Consider providing a '" + JMS_TEMPLATE_ATTRIBUTE + + "' reference where the template contains a 'messageConverter' property instead."); + } + builder.addPropertyValue(POLL_PERIOD_PROPERTY, pollPeriod); + String jmsTemplate = element.getAttribute(JMS_TEMPLATE_ATTRIBUTE); + String connectionFactory = element.getAttribute(CONNECTION_FACTORY_ATTRIBUTE); + String destination = element.getAttribute(DESTINATION_ATTRIBUTE); + String destinationName = element.getAttribute(DESTINATION_NAME_ATTRIBUTE); + if (StringUtils.hasText(jmsTemplate)) { + if (StringUtils.hasText(connectionFactory) || StringUtils.hasText(destination) || StringUtils.hasText(destinationName)) { + throw new BeanCreationException("when providing '" + JMS_TEMPLATE_ATTRIBUTE + + "', none of '" + CONNECTION_FACTORY_ATTRIBUTE + "', '" + DESTINATION_ATTRIBUTE + + "', or '" + DESTINATION_NAME_ATTRIBUTE + "' should be provided."); + } + builder.addConstructorArgReference(jmsTemplate); + } + else if (StringUtils.hasText(connectionFactory) && (StringUtils.hasText(destination) || StringUtils.hasText(destinationName))) { + builder.addConstructorArgReference(connectionFactory); + if (StringUtils.hasText(destination)) { + builder.addConstructorArgReference(destination); + } + else if (StringUtils.hasText(destinationName)) { + builder.addConstructorArg(destinationName); + } + } + else { + throw new BeanCreationException("either a '" + JMS_TEMPLATE_ATTRIBUTE + "' or both '" + + CONNECTION_FACTORY_ATTRIBUTE + "' and '" + DESTINATION_ATTRIBUTE + "' (or '" + + DESTINATION_NAME_ATTRIBUTE + "') attributes must be provided for a " + + JmsPollingSourceAdapter.class.getSimpleName()); + } + } + + private void parseMessageDrivenSourceAdapter(Element element, BeanDefinitionBuilder builder) { + String connectionFactory = element.getAttribute(CONNECTION_FACTORY_ATTRIBUTE); + String destination = element.getAttribute(DESTINATION_ATTRIBUTE); + String destinationName = element.getAttribute(DESTINATION_NAME_ATTRIBUTE); + String messageConverter = element.getAttribute(MESSAGE_CONVERTER_ATTRIBUTE); + if (StringUtils.hasText(element.getAttribute(JMS_TEMPLATE_ATTRIBUTE))) { + throw new BeanCreationException(JmsMessageDrivenSourceAdapter.class.getSimpleName() + + " does not accept a '" + JMS_TEMPLATE_ATTRIBUTE + "' reference. Both " + + "'" + CONNECTION_FACTORY_ATTRIBUTE + "' and '" + DESTINATION_ATTRIBUTE + + "' (or '" + DESTINATION_NAME_ATTRIBUTE + "') must be provided."); + } + if (StringUtils.hasText(connectionFactory) && (StringUtils.hasText(destination) || StringUtils.hasText(destinationName))) { + builder.addPropertyReference(CONNECTION_FACTORY_PROPERTY, connectionFactory); + if (StringUtils.hasText(destination)) { + builder.addPropertyReference(DESTINATION_PROPERTY, destination); + } + else { + builder.addPropertyValue(DESTINATION_NAME_PROPERTY, destinationName); + } + } + else { + throw new BeanCreationException("Both '" + CONNECTION_FACTORY_ATTRIBUTE + "' and '" + + DESTINATION_ATTRIBUTE + "' (or '" + DESTINATION_NAME_ATTRIBUTE + "') must be provided."); + } + if (StringUtils.hasText(messageConverter)) { + builder.addPropertyReference(MESSAGE_CONVERTER_PROPERTY, messageConverter); + } + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/config/JmsTargetAdapterParser.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/config/JmsTargetAdapterParser.java new file mode 100644 index 0000000000..5403bfa82e --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/jms/config/JmsTargetAdapterParser.java @@ -0,0 +1,109 @@ +/* + * 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.jms.config; + +import org.w3c.dom.Element; + +import org.springframework.beans.factory.BeanCreationException; +import org.springframework.beans.factory.config.RuntimeBeanReference; +import org.springframework.beans.factory.parsing.BeanComponentDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.RootBeanDefinition; +import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser; +import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.integration.adapter.jms.JmsTargetAdapter; +import org.springframework.integration.endpoint.DefaultMessageEndpoint; +import org.springframework.integration.scheduling.Subscription; +import org.springframework.util.StringUtils; + +/** + * Parser for the <jms-target/> element. + * + * @author Mark Fisher + */ +public class JmsTargetAdapterParser extends AbstractSingleBeanDefinitionParser { + + private static final String JMS_TEMPLATE_ATTRIBUTE = "jms-template"; + + private static final String JMS_TEMPLATE_PROPERTY = "jmsTemplate"; + + private static final String CONNECTION_FACTORY_ATTRIBUTE = "connection-factory"; + + private static final String CONNECTION_FACTORY_PROPERTY = "connectionFactory"; + + private static final String DESTINATION_ATTRIBUTE = "destination"; + + private static final String DESTINATION_NAME_ATTRIBUTE = "destination-name"; + + private static final String DESTINATION_PROPERTY = "destination"; + + private static final String DESTINATION_NAME_PROPERTY = "destinationName"; + + private static final String CHANNEL_ATTRIBUTE = "channel"; + + private static final String HANDLER_PROPERTY = "handler"; + + private static final String SUBSCRIPTION_PROPERTY = "subscription"; + + + protected Class getBeanClass(Element element) { + return DefaultMessageEndpoint.class; + } + + protected boolean shouldGenerateId() { + return false; + } + + protected boolean shouldGenerateIdAsFallback() { + return true; + } + + protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { + String jmsTemplate = element.getAttribute(JMS_TEMPLATE_ATTRIBUTE); + String connectionFactory = element.getAttribute(CONNECTION_FACTORY_ATTRIBUTE); + String destination = element.getAttribute(DESTINATION_ATTRIBUTE); + String destinationName = element.getAttribute(DESTINATION_NAME_ATTRIBUTE); + RootBeanDefinition adapterDef = new RootBeanDefinition(JmsTargetAdapter.class); + if (StringUtils.hasText(jmsTemplate)) { + if (StringUtils.hasText(connectionFactory) || StringUtils.hasText(destination) || StringUtils.hasText(destinationName)) { + throw new BeanCreationException("when providing a 'jms-template' reference, none of " + + "'connection-factory', 'destination', or 'destination-name' should be provided."); + } + adapterDef.getPropertyValues().addPropertyValue(JMS_TEMPLATE_PROPERTY, new RuntimeBeanReference(jmsTemplate)); + } + else if (StringUtils.hasText(connectionFactory) && (StringUtils.hasText(destination) ^ StringUtils.hasText(destinationName))) { + adapterDef.getPropertyValues().addPropertyValue(CONNECTION_FACTORY_PROPERTY, new RuntimeBeanReference(connectionFactory)); + if (StringUtils.hasText(destination)) { + adapterDef.getPropertyValues().addPropertyValue(DESTINATION_PROPERTY, new RuntimeBeanReference(destination)); + } + else { + adapterDef.getPropertyValues().addPropertyValue(DESTINATION_NAME_PROPERTY, destinationName); + } + } + else { + throw new BeanCreationException("either a 'jms-template' reference or both " + + "'connection-factory' and 'destination' (or 'destination-name') references must be provided."); + } + String channel = element.getAttribute(CHANNEL_ATTRIBUTE); + Subscription subscription = new Subscription(channel); + String adapterBeanName = parserContext.getReaderContext().generateBeanName(adapterDef); + parserContext.registerBeanComponent(new BeanComponentDefinition(adapterDef, adapterBeanName)); + builder.addPropertyReference(HANDLER_PROPERTY, adapterBeanName); + builder.addPropertyValue(SUBSCRIPTION_PROPERTY, subscription); + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamSource.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamSource.java new file mode 100644 index 0000000000..0637bd6f02 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamSource.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.BufferedInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; + +import org.springframework.integration.adapter.PollableSource; +import org.springframework.integration.message.MessageDeliveryException; + +/** + * A pollable source for receiving bytes from an {@link InputStream}. + * + * @author Mark Fisher + */ +public class ByteStreamSource implements PollableSource { + + private BufferedInputStream stream; + + private Object streamMonitor; + + private int bytesPerMessage = 1024; + + private boolean shouldTruncate = true; + + + public ByteStreamSource(InputStream stream) { + this(stream, -1); + } + + public ByteStreamSource(InputStream stream, int bufferSize) { + this.streamMonitor = stream; + if (stream instanceof BufferedInputStream) { + this.stream = (BufferedInputStream) stream; + } + else if (bufferSize > 0) { + this.stream = new BufferedInputStream(stream, bufferSize); + } + else { + this.stream = new BufferedInputStream(stream); + } + } + + + public void setBytesPerMessage(int bytesPerMessage) { + this.bytesPerMessage = bytesPerMessage; + } + + public void setShouldTruncate(boolean shouldTruncate) { + this.shouldTruncate = shouldTruncate; + } + + public Collection poll(int limit) { + List results = new ArrayList(); + while (results.size() < limit) { + try { + byte[] bytes; + int bytesRead = 0; + synchronized (this.streamMonitor) { + if (stream.available() == 0) { + return results; + } + bytes = new byte[bytesPerMessage]; + bytesRead = stream.read(bytes, 0, bytes.length); + } + if (bytesRead <= 0) { + return results; + } + if (!this.shouldTruncate) { + results.add(bytes); + } + else { + byte[] result = new byte[bytesRead]; + System.arraycopy(bytes, 0, result, 0, result.length); + results.add(result); + } + } + catch (IOException e) { + throw new MessageDeliveryException("IO failure occurred in adapter", e); + } + } + return results; + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamSourceAdapter.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamSourceAdapter.java new file mode 100644 index 0000000000..97fa46fcfe --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamSourceAdapter.java @@ -0,0 +1,43 @@ +/* + * 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.InputStream; + +import org.springframework.integration.adapter.PollingSourceAdapter; + +/** + * A polling source adapter that wraps a {@link ByteStreamSource}. + * + * @author Mark Fisher + */ +public class ByteStreamSourceAdapter extends PollingSourceAdapter { + + public ByteStreamSourceAdapter(InputStream stream) { + super(new ByteStreamSource(stream)); + } + + + public void setBytesPerMessage(int bytesPerMessage) { + ((ByteStreamSource) this.getSource()).setBytesPerMessage(bytesPerMessage); + } + + public void setShouldTruncate(boolean shouldTruncate) { + ((ByteStreamSource) this.getSource()).setShouldTruncate(shouldTruncate); + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamTargetAdapter.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamTargetAdapter.java new file mode 100644 index 0000000000..8e87911e07 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/ByteStreamTargetAdapter.java @@ -0,0 +1,77 @@ +/* + * 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.BufferedOutputStream; +import java.io.IOException; +import java.io.OutputStream; + +import org.springframework.integration.adapter.AbstractTargetAdapter; +import org.springframework.integration.message.MessageHandlingException; + +/** + * A target adapter that writes a byte array to an {@link OutputStream}. + * + * @author Mark Fisher + */ +public class ByteStreamTargetAdapter extends AbstractTargetAdapter { + + private BufferedOutputStream stream; + + + public ByteStreamTargetAdapter(OutputStream stream) { + this(stream, -1); + } + + public ByteStreamTargetAdapter(OutputStream stream, int bufferSize) { + if (bufferSize > 0) { + this.stream = new BufferedOutputStream(stream, bufferSize); + } + else { + this.stream = new BufferedOutputStream(stream); + } + } + + + @Override + protected boolean sendToTarget(Object object) { + if (object == null) { + if (logger.isWarnEnabled()) { + logger.warn(this.getClass().getSimpleName() + " received null object"); + } + return false; + } + try { + if (object instanceof String) { + this.stream.write(((String) object).getBytes()); + } + else if (object instanceof byte[]){ + this.stream.write((byte[]) object); + } + else { + throw new MessageHandlingException(this.getClass().getSimpleName() + + " only supports byte array and String-based messages"); + } + this.stream.flush(); + return true; + } + catch (IOException e) { + throw new MessageHandlingException("IO failure occurred in adapter", e); + } + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSource.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSource.java new file mode 100644 index 0000000000..0e71df1ce4 --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSource.java @@ -0,0 +1,81 @@ +/* + * 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.BufferedReader; +import java.io.IOException; +import java.io.InputStream; +import java.io.InputStreamReader; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; + +import org.springframework.integration.adapter.PollableSource; +import org.springframework.integration.message.MessageDeliveryException; + +/** + * A pollable source for text-based {@link InputStream InputStreams}. + * + * @author Mark Fisher + */ +public class CharacterStreamSource implements PollableSource { + + private BufferedReader reader; + + private Object streamMonitor; + + + public CharacterStreamSource(InputStream stream) { + this(stream, -1); + } + + public CharacterStreamSource(InputStream stream, int bufferSize) { + this.streamMonitor = stream; + if (bufferSize > 0) { + this.reader = new BufferedReader(new InputStreamReader(stream), bufferSize); + } + else { + this.reader = new BufferedReader(new InputStreamReader(stream)); + } + } + + + public Collection poll(int limit) { + List results = new ArrayList(); + while (results.size() < limit) { + try { + String line = null; + synchronized (this.streamMonitor) { + boolean isReady = reader.ready(); + if (!isReady) { + return results; + } + line = reader.readLine(); + } + if (line == null) { + return results; + } + results.add(line); + } + catch (IOException e) { + throw new MessageDeliveryException("IO failure occurred in adapter", e); + } + } + return results; + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapter.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapter.java new file mode 100644 index 0000000000..69cf8621cf --- /dev/null +++ b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapter.java @@ -0,0 +1,45 @@ +/* + * 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.InputStream; + +import org.springframework.integration.adapter.PollingSourceAdapter; +import org.springframework.integration.channel.MessageChannel; + +/** + * A polling source adapter that wraps a {@link CharacterStreamSource}. + * + * @author Mark Fisher + */ +public class CharacterStreamSourceAdapter extends PollingSourceAdapter { + + public CharacterStreamSourceAdapter(InputStream stream) { + super(new CharacterStreamSource(stream)); + } + + + /** + * Factory method that creates an adapter for stdin (System.in). + */ + public static CharacterStreamSourceAdapter stdinAdapter(MessageChannel channel) { + CharacterStreamSourceAdapter adapter = new CharacterStreamSourceAdapter(System.in); + adapter.setChannel(channel); + return adapter; + } + +} diff --git a/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapter.java b/spring-integration-adapters/src/main/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapter.java new file mode 100644 index 0000000000..e5cace6b03 --- /dev/null +++ b/spring-integration-adapters/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.adapter.AbstractTargetAdapter; +import org.springframework.integration.message.MessageHandlingException; + +/** + * 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-adapters/src/test/java/org/springframework/integration/adapter/jms/DefaultJmsMessagePostProcessorTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/DefaultJmsMessagePostProcessorTests.java new file mode 100644 index 0000000000..cc947bf0cb --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/DefaultJmsMessagePostProcessorTests.java @@ -0,0 +1,102 @@ +/* + * 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.jms; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; + +import javax.jms.Destination; +import javax.jms.JMSException; + +import org.junit.Test; + +import org.springframework.integration.message.StringMessage; + +/** + * @author Mark Fisher + */ +public class DefaultJmsMessagePostProcessorTests { + + @Test + public void testJmsReplyTo() throws JMSException { + StringMessage message = new StringMessage("test1"); + Destination replyTo = new Destination() {}; + message.getHeader().setAttribute(JmsTargetAdapter.JMS_REPLY_TO, replyTo); + DefaultJmsMessagePostProcessor postProcessor = new DefaultJmsMessagePostProcessor(); + javax.jms.Message jmsMessage = new StubTextMessage(); + postProcessor.postProcessJmsMessage(jmsMessage, message.getHeader()); + assertNotNull(jmsMessage.getJMSReplyTo()); + assertSame(replyTo, jmsMessage.getJMSReplyTo()); + } + + @Test + public void testJmsReplyToIgnoredIfIncorrectType() throws JMSException { + StringMessage message = new StringMessage("test1"); + message.getHeader().setAttribute(JmsTargetAdapter.JMS_REPLY_TO, "not-a-destination"); + DefaultJmsMessagePostProcessor postProcessor = new DefaultJmsMessagePostProcessor(); + javax.jms.Message jmsMessage = new StubTextMessage(); + postProcessor.postProcessJmsMessage(jmsMessage, message.getHeader()); + assertNull(jmsMessage.getJMSReplyTo()); + } + + @Test + public void testJmsCorrelationId() throws JMSException { + StringMessage message = new StringMessage("test1"); + String jmsCorrelationId = "ABC-123"; + message.getHeader().setAttribute(JmsTargetAdapter.JMS_CORRELATION_ID, jmsCorrelationId); + DefaultJmsMessagePostProcessor postProcessor = new DefaultJmsMessagePostProcessor(); + javax.jms.Message jmsMessage = new StubTextMessage(); + postProcessor.postProcessJmsMessage(jmsMessage, message.getHeader()); + assertNotNull(jmsMessage.getJMSCorrelationID()); + assertEquals(jmsCorrelationId, jmsMessage.getJMSCorrelationID()); + } + + @Test + public void testJmsCorrelationIdIgnoredIfIncorrectType() throws JMSException { + StringMessage message = new StringMessage("test1"); + message.getHeader().setAttribute(JmsTargetAdapter.JMS_CORRELATION_ID, new Integer(123)); + DefaultJmsMessagePostProcessor postProcessor = new DefaultJmsMessagePostProcessor(); + javax.jms.Message jmsMessage = new StubTextMessage(); + postProcessor.postProcessJmsMessage(jmsMessage, message.getHeader()); + assertNull(jmsMessage.getJMSCorrelationID()); + } + + @Test + public void testJmsType() throws JMSException { + StringMessage message = new StringMessage("test1"); + String jmsType = "testing"; + message.getHeader().setAttribute(JmsTargetAdapter.JMS_TYPE, jmsType); + DefaultJmsMessagePostProcessor postProcessor = new DefaultJmsMessagePostProcessor(); + javax.jms.Message jmsMessage = new StubTextMessage(); + postProcessor.postProcessJmsMessage(jmsMessage, message.getHeader()); + assertNotNull(jmsMessage.getJMSType()); + assertEquals(jmsType, jmsMessage.getJMSType()); + } + + @Test + public void testJmsTypeIgnoredIfIncorrectType() throws JMSException { + StringMessage message = new StringMessage("test1"); + message.getHeader().setAttribute(JmsTargetAdapter.JMS_TYPE, new Integer(123)); + DefaultJmsMessagePostProcessor postProcessor = new DefaultJmsMessagePostProcessor(); + javax.jms.Message jmsMessage = new StubTextMessage(); + postProcessor.postProcessJmsMessage(jmsMessage, message.getHeader()); + assertNull(jmsMessage.getJMSType()); + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubConnection.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubConnection.java new file mode 100644 index 0000000000..abec6a3bef --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubConnection.java @@ -0,0 +1,83 @@ +/* + * 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.jms; + +import javax.jms.Connection; +import javax.jms.ConnectionConsumer; +import javax.jms.ConnectionMetaData; +import javax.jms.Destination; +import javax.jms.ExceptionListener; +import javax.jms.JMSException; +import javax.jms.ServerSessionPool; +import javax.jms.Session; +import javax.jms.Topic; + +/** + * @author Mark Fisher + */ +public class StubConnection implements Connection { + + private String messageText; + + + public StubConnection(String messageText) { + this.messageText = messageText; + } + + + public void close() throws JMSException { + } + + public ConnectionConsumer createConnectionConsumer(Destination destination, String messageSelector, + ServerSessionPool sessionPool, int maxMessages) throws JMSException { + return null; + } + + public ConnectionConsumer createDurableConnectionConsumer(Topic topic, String subscriptionName, + String messageSelector, ServerSessionPool sessionPool, int maxMessages) throws JMSException { + return null; + } + + public Session createSession(boolean transacted, int acknowledgeMode) throws JMSException { + return new StubSession(this.messageText); + } + + public String getClientID() throws JMSException { + return null; + } + + public ExceptionListener getExceptionListener() throws JMSException { + return null; + } + + public ConnectionMetaData getMetaData() throws JMSException { + return null; + } + + public void setClientID(String clientID) throws JMSException { + } + + public void setExceptionListener(ExceptionListener listener) throws JMSException { + } + + public void start() throws JMSException { + } + + public void stop() throws JMSException { + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubConsumer.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubConsumer.java new file mode 100644 index 0000000000..e659ef9f14 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubConsumer.java @@ -0,0 +1,65 @@ +/* + * 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.jms; + +import javax.jms.JMSException; +import javax.jms.Message; +import javax.jms.MessageConsumer; +import javax.jms.MessageListener; + +/** + * @author Mark Fisher + */ +public class StubConsumer implements MessageConsumer { + + private String messageText; + + + public StubConsumer(String messageText) { + this.messageText = messageText; + } + + + public void close() throws JMSException { + } + + public MessageListener getMessageListener() throws JMSException { + return null; + } + + public String getMessageSelector() throws JMSException { + return null; + } + + public Message receive() throws JMSException { + StubTextMessage message = new StubTextMessage(); + message.setText(this.messageText); + return message; + } + + public Message receive(long timeout) throws JMSException { + return this.receive(); + } + + public Message receiveNoWait() throws JMSException { + return this.receive(); + } + + public void setMessageListener(MessageListener listener) throws JMSException { + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubDestination.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubDestination.java new file mode 100644 index 0000000000..e23331e733 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubDestination.java @@ -0,0 +1,26 @@ +/* + * 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.jms; + +import javax.jms.Destination; + +/** + * @author Mark Fisher + */ +public class StubDestination implements Destination { + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubSession.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubSession.java new file mode 100644 index 0000000000..55f8a16c09 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubSession.java @@ -0,0 +1,168 @@ +/* + * 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.jms; + +import java.io.Serializable; + +import javax.jms.BytesMessage; +import javax.jms.Destination; +import javax.jms.JMSException; +import javax.jms.MapMessage; +import javax.jms.Message; +import javax.jms.MessageConsumer; +import javax.jms.MessageListener; +import javax.jms.MessageProducer; +import javax.jms.ObjectMessage; +import javax.jms.Queue; +import javax.jms.QueueBrowser; +import javax.jms.Session; +import javax.jms.StreamMessage; +import javax.jms.TemporaryQueue; +import javax.jms.TemporaryTopic; +import javax.jms.TextMessage; +import javax.jms.Topic; +import javax.jms.TopicSubscriber; + +/** + * @author Mark Fisher + */ +public class StubSession implements Session { + + private String messageText; + + + public StubSession(String messageText) { + this.messageText = messageText; + } + + + public void close() throws JMSException { + } + + public void commit() throws JMSException { + } + + public QueueBrowser createBrowser(Queue queue) throws JMSException { + return null; + } + + public QueueBrowser createBrowser(Queue queue, String messageSelector) throws JMSException { + return null; + } + + public BytesMessage createBytesMessage() throws JMSException { + return null; + } + + public MessageConsumer createConsumer(Destination destination) throws JMSException { + return new StubConsumer(this.messageText); + } + + public MessageConsumer createConsumer(Destination destination, String messageSelector) throws JMSException { + return new StubConsumer(this.messageText); + } + + public MessageConsumer createConsumer(Destination destination, String messageSelector, boolean NoLocal) + throws JMSException { + return new StubConsumer(this.messageText); + } + + public TopicSubscriber createDurableSubscriber(Topic topic, String name) throws JMSException { + return null; + } + + public TopicSubscriber createDurableSubscriber(Topic topic, String name, String messageSelector, boolean noLocal) + throws JMSException { + return null; + } + + public MapMessage createMapMessage() throws JMSException { + return null; + } + + public Message createMessage() throws JMSException { + return null; + } + + public ObjectMessage createObjectMessage() throws JMSException { + return null; + } + + public ObjectMessage createObjectMessage(Serializable object) throws JMSException { + return null; + } + + public MessageProducer createProducer(Destination destination) throws JMSException { + return null; + } + + public Queue createQueue(String queueName) throws JMSException { + return null; + } + + public StreamMessage createStreamMessage() throws JMSException { + return null; + } + + public TemporaryQueue createTemporaryQueue() throws JMSException { + return null; + } + + public TemporaryTopic createTemporaryTopic() throws JMSException { + return null; + } + + public TextMessage createTextMessage() throws JMSException { + return null; + } + + public TextMessage createTextMessage(String text) throws JMSException { + return null; + } + + public Topic createTopic(String topicName) throws JMSException { + return null; + } + + public int getAcknowledgeMode() throws JMSException { + return 0; + } + + public MessageListener getMessageListener() throws JMSException { + return null; + } + + public boolean getTransacted() throws JMSException { + return false; + } + + public void recover() throws JMSException { + } + + public void rollback() throws JMSException { + } + + public void run() { + } + + public void setMessageListener(MessageListener listener) throws JMSException { + } + + public void unsubscribe(String name) throws JMSException { + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubTextMessage.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubTextMessage.java new file mode 100644 index 0000000000..08322c33eb --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/StubTextMessage.java @@ -0,0 +1,188 @@ +package org.springframework.integration.adapter.jms; + +import java.util.Enumeration; + +import javax.jms.Destination; +import javax.jms.JMSException; +import javax.jms.TextMessage; + +public class StubTextMessage implements TextMessage { + + private String text; + + private Destination replyTo; + + private String correlationID; + + private String type; + + + public String getText() throws JMSException { + return this.text; + } + + public void setText(String text) throws JMSException { + this.text = text; + } + + public void acknowledge() throws JMSException { + } + + public void clearBody() throws JMSException { + } + + public void clearProperties() throws JMSException { + } + + public boolean getBooleanProperty(String name) throws JMSException { + return false; + } + + public byte getByteProperty(String name) throws JMSException { + return 0; + } + + public double getDoubleProperty(String name) throws JMSException { + return 0; + } + + public float getFloatProperty(String name) throws JMSException { + return 0; + } + + public int getIntProperty(String name) throws JMSException { + return 0; + } + + public String getJMSCorrelationID() throws JMSException { + return this.correlationID; + } + + public byte[] getJMSCorrelationIDAsBytes() throws JMSException { + return null; + } + + public int getJMSDeliveryMode() throws JMSException { + return 0; + } + + public Destination getJMSDestination() throws JMSException { + return null; + } + + public long getJMSExpiration() throws JMSException { + return 0; + } + + public String getJMSMessageID() throws JMSException { + return null; + } + + public int getJMSPriority() throws JMSException { + return 0; + } + + public boolean getJMSRedelivered() throws JMSException { + return false; + } + + public Destination getJMSReplyTo() throws JMSException { + return this.replyTo; + } + + public long getJMSTimestamp() throws JMSException { + return 0; + } + + public String getJMSType() throws JMSException { + return this.type; + } + + public long getLongProperty(String name) throws JMSException { + return 0; + } + + public Object getObjectProperty(String name) throws JMSException { + return null; + } + + public Enumeration getPropertyNames() throws JMSException { + return null; + } + + public short getShortProperty(String name) throws JMSException { + return 0; + } + + public String getStringProperty(String name) throws JMSException { + return null; + } + + public boolean propertyExists(String name) throws JMSException { + return false; + } + + public void setBooleanProperty(String name, boolean value) throws JMSException { + } + + public void setByteProperty(String name, byte value) throws JMSException { + } + + public void setDoubleProperty(String name, double value) throws JMSException { + } + + public void setFloatProperty(String name, float value) throws JMSException { + } + + public void setIntProperty(String name, int value) throws JMSException { + } + + public void setJMSCorrelationID(String correlationID) throws JMSException { + this.correlationID = correlationID; + } + + public void setJMSCorrelationIDAsBytes(byte[] correlationID) throws JMSException { + } + + public void setJMSDeliveryMode(int deliveryMode) throws JMSException { + } + + public void setJMSDestination(Destination destination) throws JMSException { + } + + public void setJMSExpiration(long expiration) throws JMSException { + } + + public void setJMSMessageID(String id) throws JMSException { + } + + public void setJMSPriority(int priority) throws JMSException { + } + + public void setJMSRedelivered(boolean redelivered) throws JMSException { + } + + public void setJMSReplyTo(Destination replyTo) throws JMSException { + this.replyTo = replyTo; + } + + public void setJMSTimestamp(long timestamp) throws JMSException { + } + + public void setJMSType(String type) throws JMSException { + this.type = type; + } + + public void setLongProperty(String name, long value) throws JMSException { + } + + public void setObjectProperty(String name, Object value) throws JMSException { + } + + public void setShortProperty(String name, short value) throws JMSException { + } + + public void setStringProperty(String name, String value) throws JMSException { + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsSourceAdapterParserTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsSourceAdapterParserTests.java new file mode 100644 index 0000000000..ca3608b01c --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsSourceAdapterParserTests.java @@ -0,0 +1,196 @@ +/* + * 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.jms.config; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +import javax.jms.JMSException; +import javax.jms.Session; +import javax.jms.TextMessage; + +import org.junit.Test; + +import org.springframework.beans.factory.BeanCreationException; +import org.springframework.beans.factory.BeanDefinitionStoreException; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.integration.adapter.jms.JmsMessageDrivenSourceAdapter; +import org.springframework.integration.adapter.jms.JmsPollingSourceAdapter; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.message.Message; +import org.springframework.jms.support.converter.MessageConversionException; +import org.springframework.jms.support.converter.MessageConverter; + +/** + * @author Mark Fisher + */ +public class JmsSourceAdapterParserTests { + + @Test + public void testPollingAdapterWithJmsTemplate() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "pollingAdapterWithJmsTemplate.xml", this.getClass()); + context.start(); + JmsPollingSourceAdapter adapter = (JmsPollingSourceAdapter) context.getBean("adapter"); + adapter.processMessages(); + MessageChannel channel = (MessageChannel) context.getBean("channel"); + Message message = channel.receive(500); + assertNotNull("message should not be null", message); + assertEquals("polling-test", message.getPayload()); + context.stop(); + } + + @Test + public void testPollingAdapterWithConnectionFactoryAndDestination() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "pollingAdapterWithConnectionFactoryAndDestination.xml", this.getClass()); + context.start(); + JmsPollingSourceAdapter adapter = (JmsPollingSourceAdapter) context.getBean("adapter"); + adapter.processMessages(); + MessageChannel channel = (MessageChannel) context.getBean("channel"); + Message message = channel.receive(500); + assertNotNull("message should not be null", message); + assertEquals("polling-test", message.getPayload()); + context.stop(); + } + + @Test + public void testPollingAdapterWithConnectionFactoryAndDestinationName() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "pollingAdapterWithConnectionFactoryAndDestinationName.xml", this.getClass()); + context.start(); + JmsPollingSourceAdapter adapter = (JmsPollingSourceAdapter) context.getBean("adapter"); + adapter.processMessages(); + MessageChannel channel = (MessageChannel) context.getBean("channel"); + Message message = channel.receive(500); + assertNotNull("message should not be null", message); + assertEquals("polling-test", message.getPayload()); + context.stop(); + } + + @Test + public void testMessageDrivenAdapterWithConnectionFactoryAndDestination() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "messageDrivenAdapterWithConnectionFactoryAndDestination.xml", this.getClass()); + context.start(); + JmsMessageDrivenSourceAdapter adapter = (JmsMessageDrivenSourceAdapter) context.getBean("adapter"); + assertEquals(JmsMessageDrivenSourceAdapter.class, adapter.getClass()); + MessageChannel channel = (MessageChannel) context.getBean("channel"); + Message message = channel.receive(3000); + assertNotNull("message should not be null", message); + assertEquals("message-driven-test", message.getPayload()); + context.stop(); + } + + @Test + public void testMessageDrivenAdapterWithConnectionFactoryAndDestinationName() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "messageDrivenAdapterWithConnectionFactoryAndDestinationName.xml", this.getClass()); + context.start(); + JmsMessageDrivenSourceAdapter adapter = (JmsMessageDrivenSourceAdapter) context.getBean("adapter"); + assertEquals(JmsMessageDrivenSourceAdapter.class, adapter.getClass()); + MessageChannel channel = (MessageChannel) context.getBean("channel"); + Message message = channel.receive(3000); + assertNotNull("message should not be null", message); + assertEquals("message-driven-test", message.getPayload()); + context.stop(); + } + + @Test + public void testMessageDrivenAdapterWithMessageConverter() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "messageDrivenAdapterWithMessageConverter.xml", this.getClass()); + JmsMessageDrivenSourceAdapter adapter = (JmsMessageDrivenSourceAdapter) context.getBean("adapter"); + assertEquals(JmsMessageDrivenSourceAdapter.class, adapter.getClass()); + MessageChannel channel = (MessageChannel) context.getBean("channel"); + Message message = channel.receive(3000); + assertNotNull("message should not be null", message); + assertEquals("converted-test-message", message.getPayload()); + context.stop(); + } + + @Test(expected=BeanDefinitionStoreException.class) + public void testPollingAdapterWithConnectionFactoryOnly() { + try { + new ClassPathXmlApplicationContext("pollingAdapterWithConnectionFactoryOnly.xml", this.getClass()); + } + catch (RuntimeException e) { + assertEquals(BeanCreationException.class, e.getCause().getClass()); + throw e; + } + } + + @Test(expected=BeanDefinitionStoreException.class) + public void testPollingAdapterWithDestinationOnly() { + try { + new ClassPathXmlApplicationContext("pollingAdapterWithDestinationOnly.xml", this.getClass()); + } + catch (RuntimeException e) { + assertEquals(BeanCreationException.class, e.getCause().getClass()); + throw e; + } + } + + @Test(expected=BeanDefinitionStoreException.class) + public void testPollingAdapterWithDestinationNameOnly() { + try { + new ClassPathXmlApplicationContext("pollingAdapterWithDestinationNameOnly.xml", this.getClass()); + } + catch (RuntimeException e) { + assertEquals(BeanCreationException.class, e.getCause().getClass()); + throw e; + } + } + + @Test(expected=BeanDefinitionStoreException.class) + public void testPollingAdapterWithoutPollPeriod() { + try { + new ClassPathXmlApplicationContext("pollingAdapterWithoutPollPeriod.xml", this.getClass()); + } + catch (RuntimeException e) { + assertEquals(BeanCreationException.class, e.getCause().getClass()); + throw e; + } + } + + @Test(expected=BeanDefinitionStoreException.class) + public void testMessageDrivenAdapterWithConnectionFactoryOnly() { + try { + new ClassPathXmlApplicationContext("messageDrivenAdapterWithConnectionFactoryOnly.xml", this.getClass()); + } + catch (RuntimeException e) { + assertEquals(BeanCreationException.class, e.getCause().getClass()); + throw e; + } + } + + + public static class TestMessageConverter implements MessageConverter { + + public Object fromMessage(javax.jms.Message message) throws JMSException, MessageConversionException { + String original = ((TextMessage) message).getText(); + return "converted-" + original; + } + + public javax.jms.Message toMessage(Object object, Session session) throws JMSException, + MessageConversionException { + return null; + } + + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsTargetAdapterParserTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsTargetAdapterParserTests.java new file mode 100644 index 0000000000..57ef368a85 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsTargetAdapterParserTests.java @@ -0,0 +1,52 @@ +/* + * 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.jms.config; + +import static org.junit.Assert.assertEquals; + +import org.junit.Test; + +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.integration.adapter.jms.JmsTargetAdapter; +import org.springframework.integration.endpoint.DefaultMessageEndpoint; + +/** + * @author Mark Fisher + */ +public class JmsTargetAdapterParserTests { + + @Test + public void testTargetAdapterWithConnectionFactoryAndDestination() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "targetAdapterWithConnectionFactoryAndDestination.xml", this.getClass()); + DefaultMessageEndpoint endpoint = (DefaultMessageEndpoint) context.getBean("adapter"); + assertEquals(JmsTargetAdapter.class, endpoint.getHandler().getClass()); + assertEquals("adapter", endpoint.getName()); + assertEquals("testChannel", endpoint.getSubscription().getChannelName()); + } + + @Test + public void testTargetAdapterWithConnectionFactoryAndDestinationName() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "targetAdapterWithConnectionFactoryAndDestinationName.xml", this.getClass()); + DefaultMessageEndpoint endpoint = (DefaultMessageEndpoint) context.getBean("adapter"); + assertEquals(JmsTargetAdapter.class, endpoint.getHandler().getClass()); + assertEquals("adapter", endpoint.getName()); + assertEquals("testChannel", endpoint.getSubscription().getChannelName()); + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryAndDestination.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryAndDestination.xml new file mode 100644 index 0000000000..5b59551afc --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryAndDestination.xml @@ -0,0 +1,29 @@ + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryAndDestinationName.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryAndDestinationName.xml new file mode 100644 index 0000000000..83757edcb3 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryAndDestinationName.xml @@ -0,0 +1,27 @@ + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryOnly.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryOnly.xml new file mode 100644 index 0000000000..7a70079776 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithConnectionFactoryOnly.xml @@ -0,0 +1,24 @@ + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithMessageConverter.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithMessageConverter.xml new file mode 100644 index 0000000000..b00256fca4 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/messageDrivenAdapterWithMessageConverter.xml @@ -0,0 +1,32 @@ + + + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryAndDestination.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryAndDestination.xml new file mode 100644 index 0000000000..a7fb9715ac --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryAndDestination.xml @@ -0,0 +1,30 @@ + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryAndDestinationName.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryAndDestinationName.xml new file mode 100644 index 0000000000..b8b256a706 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryAndDestinationName.xml @@ -0,0 +1,28 @@ + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryOnly.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryOnly.xml new file mode 100644 index 0000000000..c6d68cec1b --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithConnectionFactoryOnly.xml @@ -0,0 +1,24 @@ + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithDestinationNameOnly.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithDestinationNameOnly.xml new file mode 100644 index 0000000000..5d185d0ca0 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithDestinationNameOnly.xml @@ -0,0 +1,16 @@ + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithDestinationOnly.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithDestinationOnly.xml new file mode 100644 index 0000000000..508b17d074 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithDestinationOnly.xml @@ -0,0 +1,18 @@ + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithJmsTemplate.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithJmsTemplate.xml new file mode 100644 index 0000000000..cdcff85d11 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithJmsTemplate.xml @@ -0,0 +1,29 @@ + + + + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithoutPollPeriod.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithoutPollPeriod.xml new file mode 100644 index 0000000000..cb48b1f27d --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/pollingAdapterWithoutPollPeriod.xml @@ -0,0 +1,29 @@ + + + + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/targetAdapterWithConnectionFactoryAndDestination.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/targetAdapterWithConnectionFactoryAndDestination.xml new file mode 100644 index 0000000000..6dd4ebb62b --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/targetAdapterWithConnectionFactoryAndDestination.xml @@ -0,0 +1,29 @@ + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/targetAdapterWithConnectionFactoryAndDestinationName.xml b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/targetAdapterWithConnectionFactoryAndDestinationName.xml new file mode 100644 index 0000000000..cec70369ae --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/targetAdapterWithConnectionFactoryAndDestinationName.xml @@ -0,0 +1,27 @@ + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceAdapterTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceAdapterTests.java new file mode 100644 index 0000000000..46cf2ef0b4 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceAdapterTests.java @@ -0,0 +1,173 @@ +/* + * 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 static org.junit.Assert.assertNull; + +import java.io.ByteArrayInputStream; + +import org.junit.Test; + +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.channel.SimpleChannel; +import org.springframework.integration.message.Message; + +/** + * @author Mark Fisher + */ +public class ByteStreamSourceAdapterTests { + + @Test + public void testEndOfStream() { + byte[] bytes = new byte[] {1,2,3}; + ByteArrayInputStream stream = new ByteArrayInputStream(bytes); + MessageChannel channel = new SimpleChannel(); + ByteStreamSourceAdapter adapter = new ByteStreamSourceAdapter(stream); + adapter.setChannel(channel); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + Message message1 = channel.receive(500); + byte[] payload = (byte[]) message1.getPayload(); + assertEquals(3, payload.length); + assertEquals(1, payload[0]); + assertEquals(2, payload[1]); + assertEquals(3, payload[2]); + Message message2 = channel.receive(0); + assertNull(message2); + adapter.processMessages(); + 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 SimpleChannel(); + ByteStreamSourceAdapter adapter = new ByteStreamSourceAdapter(stream); + adapter.setChannel(channel); + adapter.setBytesPerMessage(8); + adapter.setMaxMessagesPerTask(5); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + 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 SimpleChannel(); + ByteStreamSourceAdapter adapter = new ByteStreamSourceAdapter(stream); + adapter.setBytesPerMessage(4); + adapter.setInitialDelay(10000); + adapter.setMaxMessagesPerTask(1); + adapter.setChannel(channel); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + 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); + adapter.processMessages(); + 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 SimpleChannel(); + ByteStreamSourceAdapter adapter = new ByteStreamSourceAdapter(stream); + adapter.setInitialDelay(10000); + adapter.setChannel(channel); + adapter.setBytesPerMessage(4); + adapter.setMaxMessagesPerTask(5); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(2, count); + 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]); + Message message3 = channel.receive(0); + assertNull(message3); + } + + @Test + public void testByteArrayIsTruncated() { + byte[] bytes = new byte[] {0,1,2,3,4,5}; + ByteArrayInputStream stream = new ByteArrayInputStream(bytes); + MessageChannel channel = new SimpleChannel(); + ByteStreamSourceAdapter adapter = new ByteStreamSourceAdapter(stream); + adapter.setInitialDelay(10000); + adapter.setBytesPerMessage(4); + adapter.setMaxMessagesPerTask(1); + adapter.setChannel(channel); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + Message message1 = channel.receive(0); + assertEquals(4, ((byte[]) message1.getPayload()).length); + Message message2 = channel.receive(0); + assertNull(message2); + adapter.processMessages(); + Message message3 = channel.receive(0); + assertEquals(2, ((byte[]) message3.getPayload()).length); + } + + @Test + public void testByteArrayIsNotTruncated() { + byte[] bytes = new byte[] {0,1,2,3,4,5}; + ByteArrayInputStream stream = new ByteArrayInputStream(bytes); + MessageChannel channel = new SimpleChannel(); + ByteStreamSourceAdapter adapter = new ByteStreamSourceAdapter(stream); + adapter.setInitialDelay(10000); + adapter.setBytesPerMessage(4); + adapter.setShouldTruncate(false); + adapter.setMaxMessagesPerTask(1); + adapter.setChannel(channel); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + Message message1 = channel.receive(0); + assertEquals(4, ((byte[]) message1.getPayload()).length); + Message message2 = channel.receive(0); + assertNull(message2); + adapter.processMessages(); + Message message3 = channel.receive(0); + assertEquals(4, ((byte[]) message3.getPayload()).length); + assertEquals(0, ((byte[]) message3.getPayload())[3]); + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamTargetAdapterTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamTargetAdapterTests.java new file mode 100644 index 0000000000..4f5b47d0a8 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamTargetAdapterTests.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 java.io.IOException; + +import org.junit.Test; + +import org.springframework.integration.channel.DispatcherPolicy; +import org.springframework.integration.channel.SimpleChannel; +import org.springframework.integration.dispatcher.DefaultMessageDispatcher; +import org.springframework.integration.dispatcher.MessageDispatcher; +import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.message.StringMessage; +import org.springframework.integration.scheduling.SimpleMessagingTaskScheduler; + +/** + * @author Mark Fisher + */ +public class ByteStreamTargetAdapterTests { + + @Test + public void testSingleByteArray() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + adapter.handle(new GenericMessage(new byte[] {1,2,3})); + byte[] result = stream.toByteArray(); + assertEquals(3, result.length); + assertEquals(1, result[0]); + assertEquals(2, result[1]); + assertEquals(3, result[2]); + } + + @Test + public void testSingleString() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + adapter.handle(new StringMessage("foo")); + byte[] result = stream.toByteArray(); + assertEquals(3, result.length); + assertEquals("foo", new String(result)); + } + + @Test + public void testMaxMessagesPerTaskSameAsMessageCount() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(3); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new GenericMessage(new byte[] {1,2,3}), 0); + channel.send(new GenericMessage(new byte[] {4,5,6}), 0); + channel.send(new GenericMessage(new byte[] {7,8,9}), 0); + assertEquals(3, dispatcher.dispatch()); + byte[] result = stream.toByteArray(); + assertEquals(9, result.length); + assertEquals(1, result[0]); + assertEquals(9, result[8]); + } + + @Test + public void testMaxMessagesPerTaskLessThanMessageCount() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(2); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new GenericMessage(new byte[] {1,2,3}), 0); + channel.send(new GenericMessage(new byte[] {4,5,6}), 0); + channel.send(new GenericMessage(new byte[] {7,8,9}), 0); + assertEquals(2, dispatcher.dispatch()); + byte[] result = stream.toByteArray(); + assertEquals(6, result.length); + assertEquals(1, result[0]); + } + + @Test + public void testMaxMessagesPerTaskExceedsMessageCount() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(5); + dispatcherPolicy.setReceiveTimeout(0); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new GenericMessage(new byte[] {1,2,3}), 0); + channel.send(new GenericMessage(new byte[] {4,5,6}), 0); + channel.send(new GenericMessage(new byte[] {7,8,9}), 0); + assertEquals(3, dispatcher.dispatch()); + byte[] result = stream.toByteArray(); + assertEquals(9, result.length); + assertEquals(1, result[0]); + } + + @Test + public void testMaxMessagesLessThanMessageCountWithMultipleDispatches() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(2); + dispatcherPolicy.setReceiveTimeout(0); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new GenericMessage(new byte[] {1,2,3}), 0); + channel.send(new GenericMessage(new byte[] {4,5,6}), 0); + channel.send(new GenericMessage(new byte[] {7,8,9}), 0); + assertEquals(2, dispatcher.dispatch()); + byte[] result1 = stream.toByteArray(); + assertEquals(6, result1.length); + assertEquals(1, result1[0]); + assertEquals(1, dispatcher.dispatch()); + byte[] result2 = stream.toByteArray(); + assertEquals(9, result2.length); + assertEquals(1, result2[0]); + assertEquals(7, result2[6]); + } + + @Test + public void testMaxMessagesExceedsMessageCountWithMultipleDispatches() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(5); + dispatcherPolicy.setReceiveTimeout(0); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new GenericMessage(new byte[] {1,2,3}), 0); + channel.send(new GenericMessage(new byte[] {4,5,6}), 0); + channel.send(new GenericMessage(new byte[] {7,8,9}), 0); + assertEquals(3, dispatcher.dispatch()); + byte[] result1 = stream.toByteArray(); + assertEquals(9, result1.length); + assertEquals(1, result1[0]); + assertEquals(0, dispatcher.dispatch()); + byte[] result2 = stream.toByteArray(); + assertEquals(9, result2.length); + assertEquals(1, result2[0]); + } + + @Test + public void testStreamResetBetweenDispatches() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(2); + dispatcherPolicy.setReceiveTimeout(0); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new GenericMessage(new byte[] {1,2,3}), 0); + channel.send(new GenericMessage(new byte[] {4,5,6}), 0); + channel.send(new GenericMessage(new byte[] {7,8,9}), 0); + assertEquals(2, dispatcher.dispatch()); + byte[] result1 = stream.toByteArray(); + assertEquals(6, result1.length); + stream.reset(); + assertEquals(1, dispatcher.dispatch()); + byte[] result2 = stream.toByteArray(); + assertEquals(3, result2.length); + assertEquals(7, result2[0]); + } + + @Test + public void testStreamWriteBetweenDispatches() throws IOException { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + ByteStreamTargetAdapter adapter = new ByteStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(2); + dispatcherPolicy.setReceiveTimeout(0); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new GenericMessage(new byte[] {1,2,3}), 0); + channel.send(new GenericMessage(new byte[] {4,5,6}), 0); + channel.send(new GenericMessage(new byte[] {7,8,9}), 0); + assertEquals(2, dispatcher.dispatch()); + byte[] result1 = stream.toByteArray(); + assertEquals(6, result1.length); + stream.write(new byte[] {123}); + stream.flush(); + assertEquals(1, dispatcher.dispatch()); + byte[] result2 = stream.toByteArray(); + assertEquals(10, result2.length); + assertEquals(1, result2[0]); + assertEquals(123, result2[6]); + assertEquals(7, result2[7]); + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapterTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapterTests.java new file mode 100644 index 0000000000..52d32815b5 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceAdapterTests.java @@ -0,0 +1,111 @@ +/* + * 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 static org.junit.Assert.assertNull; + +import java.io.ByteArrayInputStream; + +import org.junit.Test; + +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.channel.SimpleChannel; +import org.springframework.integration.message.Message; + +/** + * @author Mark Fisher + */ +public class CharacterStreamSourceAdapterTests { + + @Test + public void testEndOfStream() { + byte[] bytes = "test".getBytes(); + ByteArrayInputStream stream = new ByteArrayInputStream(bytes); + MessageChannel channel = new SimpleChannel(); + CharacterStreamSourceAdapter adapter = new CharacterStreamSourceAdapter(stream); + adapter.setChannel(channel); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + Message message1 = channel.receive(0); + assertEquals("test", message1.getPayload()); + Message message2 = channel.receive(0); + assertNull(message2); + adapter.processMessages(); + Message message3 = channel.receive(0); + assertNull(message3); + } + + @Test + public void testEndOfStreamWithMaxMessagesPerTask() { + byte[] bytes = "test".getBytes(); + ByteArrayInputStream stream = new ByteArrayInputStream(bytes); + MessageChannel channel = new SimpleChannel(); + CharacterStreamSourceAdapter adapter = new CharacterStreamSourceAdapter(stream); + adapter.setChannel(channel); + adapter.setMaxMessagesPerTask(5); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + 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"; + ByteArrayInputStream stream = new ByteArrayInputStream(s.getBytes()); + MessageChannel channel = new SimpleChannel(); + CharacterStreamSourceAdapter adapter = new CharacterStreamSourceAdapter(stream); + adapter.setInitialDelay(10000); + adapter.setMaxMessagesPerTask(1); + adapter.setChannel(channel); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(1, count); + Message message1 = channel.receive(0); + assertEquals("test1", message1.getPayload()); + Message message2 = channel.receive(0); + assertNull(message2); + adapter.processMessages(); + Message message3 = channel.receive(0); + assertEquals("test2", message3.getPayload()); + } + + @Test + public void testLessThanMaxMessagesAvailable() { + String s = "test1" + System.getProperty("line.separator") + "test2"; + ByteArrayInputStream stream = new ByteArrayInputStream(s.getBytes()); + MessageChannel channel = new SimpleChannel(); + CharacterStreamSourceAdapter adapter = new CharacterStreamSourceAdapter(stream); + adapter.setChannel(channel); + adapter.setMaxMessagesPerTask(5); + adapter.start(); + int count = adapter.processMessages(); + assertEquals(2, count); + Message message1 = channel.receive(0); + assertEquals("test1", message1.getPayload()); + Message message2 = channel.receive(0); + assertEquals("test2", message2.getPayload()); + Message message3 = channel.receive(0); + assertNull(message3); + } + +} diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapterTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapterTests.java new file mode 100644 index 0000000000..cbd9365956 --- /dev/null +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamTargetAdapterTests.java @@ -0,0 +1,191 @@ +/* + * 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.channel.DispatcherPolicy; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.channel.SimpleChannel; +import org.springframework.integration.dispatcher.DefaultMessageDispatcher; +import org.springframework.integration.dispatcher.MessageDispatcher; +import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.message.StringMessage; +import org.springframework.integration.scheduling.SimpleMessagingTaskScheduler; + +/** + * @author Mark Fisher + */ +public class CharacterStreamTargetAdapterTests { + + private SimpleMessagingTaskScheduler scheduler = new SimpleMessagingTaskScheduler(1); + + + @Test + public void testSingleString() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + adapter.handle(new StringMessage("foo")); + 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); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new StringMessage("foo"), 0); + channel.send(new StringMessage("bar"), 0); + 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.setShouldAppendNewLine(true); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new StringMessage("foo"), 0); + channel.send(new StringMessage("bar"), 0); + 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(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(2); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + channel.send(new StringMessage("foo"), 0); + channel.send(new StringMessage("bar"), 0); + assertEquals(2, dispatcher.dispatch()); + String result = new String(stream.toByteArray()); + assertEquals("foobar", result); + } + + @Test + public void testMaxMessagesPerTaskExceedsMessageCountWithAppendedNewLines() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setMaxMessagesPerTask(10); + dispatcherPolicy.setReceiveTimeout(0); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + adapter.setShouldAppendNewLine(true); + dispatcher.addHandler(adapter); + channel.send(new StringMessage("foo"), 0); + channel.send(new StringMessage("bar"), 0); + 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); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + 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(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setReceiveTimeout(0); + dispatcherPolicy.setMaxMessagesPerTask(2); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + MessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + dispatcher.addHandler(adapter); + TestObject testObject1 = new TestObject("foo"); + TestObject testObject2 = new TestObject("bar"); + channel.send(new GenericMessage(testObject1), 0); + channel.send(new GenericMessage(testObject2), 0); + assertEquals(2, dispatcher.dispatch()); + String result = new String(stream.toByteArray()); + assertEquals("foobar", result); + } + + @Test + public void testTwoNonStringObjectWithNewLines() { + ByteArrayOutputStream stream = new ByteArrayOutputStream(); + CharacterStreamTargetAdapter adapter = new CharacterStreamTargetAdapter(stream); + DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + dispatcherPolicy.setReceiveTimeout(0); + dispatcherPolicy.setMaxMessagesPerTask(2); + SimpleChannel channel = new SimpleChannel(dispatcherPolicy); + DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); + adapter.setShouldAppendNewLine(true); + dispatcher.addHandler(adapter); + TestObject testObject1 = new TestObject("foo"); + TestObject testObject2 = new TestObject("bar"); + channel.send(new GenericMessage(testObject1), 0); + channel.send(new GenericMessage(testObject2), 0); + 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; + } + } + +}