From 193f3d2439a8d530a01751e3595629fb7033c281 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 11 Nov 2011 12:14:16 -0500 Subject: [PATCH] INT-2235 added support for resource-inbound-channel-adapter removed pre-fetching logic changed ResourcePatternResolvingMessageSource to return multiple Resources made filter optional --- .../xml/IntegrationNamespaceHandler.java | 1 + .../ResourceInboundChannelAdapterParser.java | 42 ++++++ .../endpoint/SourcePollingChannelAdapter.java | 5 - .../resource/ResourceMessageSource.java | 104 ++++++++++++++ .../AcceptOnceUntilPurgedElementFilter.java | 81 +++++++++++ .../integration/util/ElementFilter.java | 28 ++++ .../config/xml/spring-integration-2.1.xsd | 72 ++++++++++ .../ResourcePatternResolver-config-custom.xml | 20 +++ .../ResourcePatternResolver-config-fail.xml | 17 +++ .../ResourcePatternResolver-config-usage.xml | 18 +++ ...ResourcePatternResolver-config-usagerf.xml | 21 +++ .../ResourcePatternResolver-config.xml | 17 +++ .../ResourcePatternResolverParserTests.java | 132 ++++++++++++++++++ 13 files changed, 553 insertions(+), 5 deletions(-) create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-custom.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-fail.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usage.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config-usagerf.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolver-config.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java 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..ad02416e98 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResourceInboundChannelAdapterParser.java @@ -0,0 +1,42 @@ +/* + * 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); + IntegrationNamespaceUtils.setValueIfAttributeDefined(sourceBuilder, element, "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..6a9714a187 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/resource/ResourceMessageSource.java @@ -0,0 +1,104 @@ +/* + * 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.ArrayList; +import java.util.List; + +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.ElementFilter; +import org.springframework.util.Assert; +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 + * @since 2.1 + */ +public class ResourceMessageSource extends AbstractMessageSource implements ApplicationContextAware, InitializingBean { + + private volatile String pattern; + + private volatile ApplicationContext applicationContext; + + private volatile ResourcePatternResolver patternResolver; + + private volatile ElementFilter filter; + + + public void setPatternResolver(ResourcePatternResolver patternResolver) { + this.patternResolver = patternResolver; + } + + public void setPattern(String pattern) { + this.pattern = pattern; + } + + public void setFilter(ElementFilter 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' is specified"); + Assert.hasText(this.pattern, "'pattern' must be specified"); + } + + @Override + protected Resource[] doReceive() { + try { + Resource[] resources = this.patternResolver.getResources(this.pattern); + if (this.filter != null && !ObjectUtils.isEmpty(resources)) { + List filteredResources = new ArrayList(); + for (Resource resource : resources) { + Resource filteredResource = this.filter.filter(resource); + if (filteredResource != null) { + filteredResources.add(filteredResource); + } + } + if (filteredResources.size() == 0) { + 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/AcceptOnceUntilPurgedElementFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java new file mode 100644 index 0000000000..e34ceeee97 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/AcceptOnceUntilPurgedElementFilter.java @@ -0,0 +1,81 @@ +/* + * 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.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 ElementFilter} which will queue all items that's 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 corresponds to the underlying items once processing + * is done to eliminate duplicate processing. + * + * @author Oleg Zhurakousky + * @since 2.1 + */ +public class AcceptOnceUntilPurgedElementFilter implements ElementFilter { + + private final Log logger = LogFactory.getLog(this.getClass()); + + private final Queue seenItems; + + private final Object seenQueueMonitor = new Object(); + + public AcceptOnceUntilPurgedElementFilter(){ + this(Integer.MAX_VALUE); + } + + public AcceptOnceUntilPurgedElementFilter(int maxCapacity){ + seenItems = new LinkedBlockingQueue(maxCapacity); + } + + 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; + } + } + + public T filter(T unfilteredElement) { + Assert.notNull(unfilteredElement, "'unfilteredElement' must not be null"); + + if (this.accept(unfilteredElement)){ + return unfilteredElement; + } + else { + return null; + } + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java new file mode 100644 index 0000000000..d0aaac9277 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/ElementFilter.java @@ -0,0 +1,28 @@ +/* + * 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; + + +/** + * Base strategy for filtering out an element + * + * @author Oleg Zhurakousky + * @since 2.1 + */ +public interface ElementFilter { + + T filter(T unfilteredElement); +} 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/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..5ff664ebf9 --- /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 @@ + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java new file mode 100644 index 0000000000..68da765a4d --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/resource/ResourcePatternResolverParserTests.java @@ -0,0 +1,132 @@ +/* + * 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.io.File; + +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.ElementFilter; + +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; + +/** + * @author Oleg Zhurakousky + * + */ +public class ResourcePatternResolverParserTests { + + @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 ElementFilter { + + private volatile boolean once = false; + + public Resource filter(Resource unfilteredElement) { + + if (!once){ + once = true; + return unfilteredElement; + } + return null; + } + + } +}