diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java index 25377b2b8e..ce084278c6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceHandler.java @@ -53,6 +53,7 @@ public class IntegrationNamespaceHandler extends AbstractIntegrationNamespaceHan registerBeanDefinitionParser("claim-check-in", new ClaimCheckInParser()); registerBeanDefinitionParser("claim-check-out", new ClaimCheckOutParser()); registerBeanDefinitionParser("inbound-channel-adapter", new MethodInvokingInboundChannelAdapterParser()); + registerBeanDefinitionParser("resource-inbound-channel-adapter", new ResourceInboundChannelAdapterParser()); registerBeanDefinitionParser("outbound-channel-adapter", new MethodInvokingOutboundChannelAdapterParser()); registerBeanDefinitionParser("logging-channel-adapter", new LoggingChannelAdapterParser()); registerBeanDefinitionParser("gateway", new GatewayParser()); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java new file mode 100644 index 0000000000..9284805970 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java @@ -0,0 +1,43 @@ +/* + * Copyright 2002-2011 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.config.xml; + +import org.springframework.beans.BeanMetadataElement; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.integration.resource.ResourceMessageSource; +import org.w3c.dom.Element; + +/** + * Parser for 'resource-inbound-channel-adapter' + * + * @author Oleg Zhurakousky + * @since 2.1 + */ +public class ResourceInboundChannelAdapterParser extends AbstractPollingInboundChannelAdapterParser { + + + @Override + protected BeanMetadataElement parseSource(Element element, ParserContext parserContext) { + BeanDefinitionBuilder sourceBuilder = BeanDefinitionBuilder.genericBeanDefinition(ResourceMessageSource.class); + sourceBuilder.addConstructorArgValue(element.getAttribute("pattern")); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(sourceBuilder, element, "pattern-resolver"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(sourceBuilder, element, "filter"); + return sourceBuilder.getBeanDefinition(); + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java index ef72728e6f..31f0424f3d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java @@ -16,9 +16,6 @@ package org.springframework.integration.endpoint; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; - import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; import org.springframework.integration.context.NamedComponent; @@ -36,8 +33,6 @@ import org.springframework.util.Assert; * @author Oleg Zhurakousky */ public class SourcePollingChannelAdapter extends AbstractPollingEndpoint implements TrackableComponent { - - private final Log logger = LogFactory.getLog(this.getClass()); private volatile MessageSource source; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java b/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java new file mode 100644 index 0000000000..7b104b7ca5 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java @@ -0,0 +1,101 @@ +/* + * Copyright 2002-2011 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.resource; + +import java.util.Arrays; +import java.util.Collection; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.MessageSource; +import org.springframework.core.io.Resource; +import org.springframework.core.io.support.ResourcePatternResolver; +import org.springframework.integration.MessagingException; +import org.springframework.integration.endpoint.AbstractMessageSource; +import org.springframework.integration.util.CollectionFilter; +import org.springframework.util.Assert; +import org.springframework.util.CollectionUtils; +import org.springframework.util.ObjectUtils; + +/** + * Implementation of {@link MessageSource} based on {@link ResourcePatternResolver} which will + * attempt to resolve {@link Resource}s based on the pattern specified. + * + * @author Oleg Zhurakousky + * @author Mark Fisher + * @since 2.1 + */ +public class ResourceMessageSource extends AbstractMessageSource implements ApplicationContextAware, InitializingBean { + + private final String pattern; + + private volatile ApplicationContext applicationContext; + + private volatile ResourcePatternResolver patternResolver; + + private volatile CollectionFilter filter; + + + public ResourceMessageSource(String pattern) { + Assert.hasText(pattern, "pattern must not be empty"); + this.pattern = pattern; + } + + + public void setPatternResolver(ResourcePatternResolver patternResolver) { + this.patternResolver = patternResolver; + } + + public void setFilter(CollectionFilter filter) { + this.filter = filter; + } + + public void setApplicationContext(ApplicationContext applicationContext) { + this.applicationContext = applicationContext; + } + + public void afterPropertiesSet() { + if (this.patternResolver == null) { + if (this.applicationContext instanceof ResourcePatternResolver) { + this.patternResolver = this.applicationContext; + } + } + Assert.notNull(this.patternResolver, "no 'patternResolver' available"); + } + + @Override + protected Resource[] doReceive() { + try { + Resource[] resources = this.patternResolver.getResources(this.pattern); + if (this.filter != null && !ObjectUtils.isEmpty(resources)) { + Collection filteredResources = this.filter.filter(Arrays.asList(resources)); + if (CollectionUtils.isEmpty(filteredResources)) { + resources = null; + } + else { + resources = filteredResources.toArray(new Resource[0]); + } + } + return resources; + } + catch (Exception e) { + throw new MessagingException("Attempt to retrieve Resources failed", e); + } + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedCollectionFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedCollectionFilter.java new file mode 100644 index 0000000000..a76280be43 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedCollectionFilter.java @@ -0,0 +1,89 @@ +/* + * Copyright 2002-2011 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.util; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Queue; +import java.util.concurrent.LinkedBlockingQueue; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.util.Assert; + +/** + * An implementation of {@link CollectionFilter} which will queue all items that have been seen until + * the queue reaches its capacity after which one item from the queue will be purged to make room for a + * new item to be added. Note that however unlikely the removed item will now appear as unprocessed + * so it is highly recommended to move/delete resources which correspond to the underlying items once processing + * is done to eliminate duplicate processing. + * + * @author Oleg Zhurakousky + * @author Mark Fisher + * @since 2.1 + */ +public class AcceptOnceUntilPurgedCollectionFilter implements CollectionFilter { + + private final Log logger = LogFactory.getLog(this.getClass()); + + private final Queue seenItems; + + private final Object seenQueueMonitor = new Object(); + + + public AcceptOnceUntilPurgedCollectionFilter() { + this(Integer.MAX_VALUE); + } + + public AcceptOnceUntilPurgedCollectionFilter(int maxCapacity) { + this.seenItems = new LinkedBlockingQueue(maxCapacity); + } + + + public Collection filter(Collection unfilteredElements) { + Assert.notNull(unfilteredElements, "'unfilteredElements' must not be null"); + List filteredElements = new ArrayList(); + if (unfilteredElements.size() > 0) { + for (T element : unfilteredElements) { + if (this.accept(element)) { + filteredElements.add(element); + } + } + } + return filteredElements; + } + + private boolean accept(T item) { + synchronized (this.seenQueueMonitor) { + boolean accepted = false; + + if (!this.seenItems.contains(item)) { + accepted = this.seenItems.offer(item); + if (!accepted){ + logger.warn("'seenQueueMonitor' queue of AcceptOnceUntilPurgedElementFilter is at the capacity, " + + "evicting one item to make room for another"); + this.seenItems.poll(); + accepted = this.seenItems.offer(item); + } + } + + return accepted; + } + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/CollectionFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/CollectionFilter.java new file mode 100644 index 0000000000..b1adc56362 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/CollectionFilter.java @@ -0,0 +1,32 @@ +/* + * Copyright 2002-2011 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.util; + +import java.util.Collection; + +/** + * Base strategy for filtering out a subset of a Collection of elements. + * + * @author Oleg Zhurakousky + * @author Mark Fisher + * @since 2.1 + */ +public interface CollectionFilter { + + Collection filter(Collection unfilteredElements); + +} diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd index 9d1b2f987f..5ebc415514 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.1.xsd @@ -716,6 +716,78 @@ + + + + + Defines a Channel Adapter that receives Resource(s) and sends them to a + MessageChannel identified via 'channel' attribute. + + + + + + + + + + + + Component identifier + + + + + + + Channel where Message will be sent to + + + + + + + Reference to the implementation of org.springframework.integration.util.ElementFilter. + + + + + + + + + + + + Lifecycle attribute signaling if this component should be started during Application Context startup. + + + + + + + Location pattern expression (e.g., "/**/*.txt") + + + + + + + + Reference to a org.springframework.core.io.support.ResourcePatternResolver. + + + + + + + Maximum amount of time in milliseconds to wait when sending a message to the channel if such channel may block. + For example, a Queue Channel can block until space is available if its maximum capacity has been reached. + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourceInboundChannelAdapterParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourceInboundChannelAdapterParserTests.java new file mode 100644 index 0000000000..2af6712e9a --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourceInboundChannelAdapterParserTests.java @@ -0,0 +1,136 @@ +/* + * Copyright 2002-2011 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.resource; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; + +import java.io.File; +import java.util.Collection; +import java.util.Collections; + +import org.junit.Test; + +import org.springframework.beans.factory.BeanCreationException; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.core.io.Resource; +import org.springframework.integration.Message; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.endpoint.SourcePollingChannelAdapter; +import org.springframework.integration.test.util.TestUtils; +import org.springframework.integration.util.CollectionFilter; +import org.springframework.util.CollectionUtils; + +/** + * @author Oleg Zhurakousky + * @since 2.1 + */ +public class ResourceInboundChannelAdapterParserTests { + + @Test + public void testDefaultConfig(){ + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config.xml", this.getClass()); + SourcePollingChannelAdapter resourceAdapter = context.getBean("resourceAdapterDefault", SourcePollingChannelAdapter.class); + ResourceMessageSource source = TestUtils.getPropertyValue(resourceAdapter, "source", ResourceMessageSource.class); + assertNotNull(source); + boolean autoStartup = TestUtils.getPropertyValue(resourceAdapter, "autoStartup", Boolean.class); + assertFalse(autoStartup); + + assertEquals("/**/*", TestUtils.getPropertyValue(source, "pattern")); + assertEquals(context, TestUtils.getPropertyValue(source, "patternResolver")); + } + + @Test(expected=BeanCreationException.class) + public void testDefaultConfigNoLocationPattern(){ + new ClassPathXmlApplicationContext("ResourcePatternResolver-config-fail.xml", this.getClass()); + } + + @Test + public void testCustomPatternResolver(){ + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config-custom.xml", this.getClass()); + SourcePollingChannelAdapter resourceAdapter = context.getBean("resourceAdapterDefault", SourcePollingChannelAdapter.class); + ResourceMessageSource source = TestUtils.getPropertyValue(resourceAdapter, "source", ResourceMessageSource.class); + assertNotNull(source); + assertEquals(context.getBean("customResolver"), TestUtils.getPropertyValue(source, "patternResolver")); + } + + @SuppressWarnings("unchecked") + @Test + public void testUsage() throws Exception{ + + File baseDir = new File(System.getProperty("java.io.tmpdir")); + + for (int i = 0; i < 10; i++) { + File f = new File(baseDir, "testUsage"+i); + f.createNewFile(); + } + + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config-usage.xml", this.getClass()); + QueueChannel resultChannel = context.getBean("resultChannel", QueueChannel.class); + Message message = (Message) resultChannel.receive(3000); + assertNotNull(message); + Resource[] resources = message.getPayload(); + for (Resource resource : resources) { + assertTrue(resource.getURI().toString().contains("testUsage")); + } + } + + @SuppressWarnings("unchecked") + @Test + public void testUsageWithCustomResourceFilter() throws Exception{ + + File baseDir = new File(System.getProperty("java.io.tmpdir")); + for (int i = 0; i < 10; i++) { + File f = new File(baseDir, "testUsageWithRf"+i); + f.createNewFile(); + } + + ApplicationContext context = new ClassPathXmlApplicationContext("ResourcePatternResolver-config-usagerf.xml", this.getClass()); + SourcePollingChannelAdapter resourceAdapter = context.getBean("resourceAdapterDefault", SourcePollingChannelAdapter.class); + ResourceMessageSource source = TestUtils.getPropertyValue(resourceAdapter, "source", ResourceMessageSource.class); + assertNotNull(source); + assertEquals(context.getBean("rlFilter"), TestUtils.getPropertyValue(source, "filter")); + + QueueChannel resultChannel = context.getBean("resultChannel", QueueChannel.class); + + Message message = (Message) resultChannel.receive(1000); + assertNotNull(message); + + message = (Message) resultChannel.receive(1000); + assertNull(message); + + } + + + public static class OneItemAndNeverAgainResourceListFilter implements CollectionFilter { + + private volatile boolean once = false; + + public Collection filter(Collection unfilteredResources) { + if (!once && !CollectionUtils.isEmpty(unfilteredResources)) { + once = true; + return Collections.singletonList(unfilteredResources.iterator().next()); + } + return null; + } + } + +} diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml new file mode 100644 index 0000000000..7c6c3aea2c --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml @@ -0,0 +1,20 @@ + + + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml new file mode 100644 index 0000000000..9671745127 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml @@ -0,0 +1,17 @@ + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml new file mode 100644 index 0000000000..d24ffcfc45 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml @@ -0,0 +1,18 @@ + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml new file mode 100644 index 0000000000..3c71c30647 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml @@ -0,0 +1,21 @@ + + + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml new file mode 100644 index 0000000000..f5eff57ce0 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml @@ -0,0 +1,17 @@ + + + + + + + + + + + + +