From 019ec2123d3f2de5381942465a920a3835379c30 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 19 May 2008 17:09:32 +0000 Subject: [PATCH] PollingSourceEndpoint is now just SourceEndpoint. --- .../jms/config/JmsSourceParserTests.java | 3 - .../adapter/stream/ByteStreamSourceTests.java | 14 +- .../stream/CharacterStreamSourceTests.java | 10 +- ...essageEndpointAnnotationPostProcessor.java | 4 +- .../config/SourceEndpointParser.java | 4 +- ...rceEndpoint.java => AbstractEndpoint.java} | 40 ++--- .../endpoint/PollingSourceEndpoint.java | 158 ------------------ .../integration/endpoint/SourceEndpoint.java | 135 ++++++++++++++- .../integration/endpoint/TargetEndpoint.java | 27 +-- .../integration/bus/MessageBusTests.java | 4 +- ...intTests.java => SourceEndpointTests.java} | 24 +-- .../integration/handler/adapterTests.xml | 2 +- .../integration/samples/cafe/Barista.java | 4 +- 13 files changed, 179 insertions(+), 250 deletions(-) rename spring-integration-core/src/main/java/org/springframework/integration/endpoint/{AbstractSourceEndpoint.java => AbstractEndpoint.java} (56%) delete mode 100644 spring-integration-core/src/main/java/org/springframework/integration/endpoint/PollingSourceEndpoint.java rename spring-integration-core/src/test/java/org/springframework/integration/endpoint/{PollingSourceEndpointTests.java => SourceEndpointTests.java} (93%) diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsSourceParserTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsSourceParserTests.java index ff57a0be18..0969bafbd3 100644 --- a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsSourceParserTests.java +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/jms/config/JmsSourceParserTests.java @@ -27,7 +27,6 @@ import org.springframework.beans.factory.NoSuchBeanDefinitionException; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.adapter.jms.JmsSource; import org.springframework.integration.channel.MessageChannel; -import org.springframework.integration.endpoint.PollingSourceEndpoint; import org.springframework.integration.message.Message; /** @@ -120,8 +119,6 @@ public class JmsSourceParserTests { ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( "jmsSourceEndpoint.xml", this.getClass()); context.start(); - PollingSourceEndpoint endpoint = (PollingSourceEndpoint) context.getBean("endpoint"); - assertEquals(JmsSource.class, endpoint.getSource().getClass()); MessageChannel channel = (MessageChannel) context.getBean("channel"); Message message = channel.receive(3000); assertNotNull("message should not be null", message); diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java index f1b91e7e90..2c4be97ee2 100644 --- a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java @@ -25,7 +25,7 @@ import org.junit.Test; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.endpoint.PollingSourceEndpoint; +import org.springframework.integration.endpoint.SourceEndpoint; import org.springframework.integration.message.Message; import org.springframework.integration.scheduling.PollingSchedule; @@ -42,7 +42,7 @@ public class ByteStreamSourceTests { ByteStreamSource source = new ByteStreamSource(stream); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.run(); Message message1 = channel.receive(500); byte[] payload = (byte[]) message1.getPayload(); @@ -66,7 +66,7 @@ public class ByteStreamSourceTests { source.setBytesPerMessage(8); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(5); endpoint.run(); Message message1 = channel.receive(500); @@ -84,7 +84,7 @@ public class ByteStreamSourceTests { source.setBytesPerMessage(4); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(1); endpoint.run(); Message message1 = channel.receive(0); @@ -109,7 +109,7 @@ public class ByteStreamSourceTests { source.setBytesPerMessage(4); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(5); endpoint.run(); Message message1 = channel.receive(0); @@ -133,7 +133,7 @@ public class ByteStreamSourceTests { source.setBytesPerMessage(4); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(1); endpoint.run(); Message message1 = channel.receive(0); @@ -155,7 +155,7 @@ public class ByteStreamSourceTests { source.setShouldTruncate(false); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(1); endpoint.run(); Message message1 = channel.receive(0); diff --git a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java index 758de209cb..825128fef2 100644 --- a/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java +++ b/spring-integration-adapters/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java @@ -25,7 +25,7 @@ import org.junit.Test; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.endpoint.PollingSourceEndpoint; +import org.springframework.integration.endpoint.SourceEndpoint; import org.springframework.integration.message.Message; import org.springframework.integration.scheduling.PollingSchedule; @@ -41,7 +41,7 @@ public class CharacterStreamSourceTests { CharacterStreamSource source = new CharacterStreamSource(reader); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.run(); Message message1 = channel.receive(0); assertEquals("test", message1.getPayload()); @@ -59,7 +59,7 @@ public class CharacterStreamSourceTests { CharacterStreamSource source = new CharacterStreamSource(reader); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(5); endpoint.run(); Message message1 = channel.receive(0); @@ -76,7 +76,7 @@ public class CharacterStreamSourceTests { CharacterStreamSource source = new CharacterStreamSource(reader); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(1); endpoint.run(); Message message1 = channel.receive(0); @@ -96,7 +96,7 @@ public class CharacterStreamSourceTests { CharacterStreamSource source = new CharacterStreamSource(reader); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(5000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(5); endpoint.run(); Message message1 = channel.receive(500); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java index cc92aac092..c70a648659 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java @@ -48,7 +48,7 @@ import org.springframework.integration.channel.ChannelRegistryAware; import org.springframework.integration.dispatcher.SynchronousChannel; import org.springframework.integration.endpoint.ConcurrencyPolicy; import org.springframework.integration.endpoint.HandlerEndpoint; -import org.springframework.integration.endpoint.PollingSourceEndpoint; +import org.springframework.integration.endpoint.SourceEndpoint; import org.springframework.integration.handler.AbstractMessageHandlerAdapter; import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.handler.MessageHandlerChain; @@ -163,7 +163,7 @@ public class MessageEndpointAnnotationPostProcessor implements BeanPostProcessor PollingSchedule schedule = new PollingSchedule(period); schedule.setInitialDelay(initialDelay); schedule.setFixedRate(fixedRate); - PollingSourceEndpoint sourceEndpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint sourceEndpoint = new SourceEndpoint(source, channel, schedule); String channelName = beanName + "-inputChannel"; messageBus.registerChannel(channelName, channel); messageBus.registerEndpoint(beanName + "-sourceEndpoint", sourceEndpoint); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/SourceEndpointParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/SourceEndpointParser.java index b26edc00e9..3854351a28 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/SourceEndpointParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/SourceEndpointParser.java @@ -22,7 +22,7 @@ import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.AbstractSimpleBeanDefinitionParser; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.ConfigurationException; -import org.springframework.integration.endpoint.PollingSourceEndpoint; +import org.springframework.integration.endpoint.SourceEndpoint; import org.springframework.integration.scheduling.PollingSchedule; import org.springframework.integration.scheduling.Schedule; import org.springframework.util.StringUtils; @@ -36,7 +36,7 @@ import org.springframework.util.xml.DomUtils; public class SourceEndpointParser extends AbstractSimpleBeanDefinitionParser { protected final Class getBeanClass(Element element) { - return PollingSourceEndpoint.class; + return SourceEndpoint.class; } protected boolean shouldGenerateId() { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractSourceEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java similarity index 56% rename from spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractSourceEndpoint.java rename to spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java index 54cf1c9865..a34d7a9aaf 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractSourceEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java @@ -19,48 +19,32 @@ package org.springframework.integration.endpoint; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.springframework.integration.channel.MessageChannel; -import org.springframework.integration.message.Source; -import org.springframework.util.Assert; - /** - * Base class for {@link SourceEndpoint} implementations. + * Base class for {@link MessageEndpoint} implementations. * * @author Mark Fisher */ -public abstract class AbstractSourceEndpoint implements SourceEndpoint { +public abstract class AbstractEndpoint implements MessageEndpoint { protected final Log logger = LogFactory.getLog(this.getClass()); - private final Source source; - - private final MessageChannel channel; - private volatile String name; - public AbstractSourceEndpoint(Source source, MessageChannel channel) { - Assert.notNull(source, "source must not be null"); - Assert.notNull(channel, "channel must not be null"); - this.source = source; - this.channel = channel; - } - - - public Source getSource() { - return this.source; - } - - public void setBeanName(String name) { - this.name = name; - } - public String getName() { return this.name; } - protected MessageChannel getChannel() { - return this.channel; + public void setName(String name) { + this.name = name; + } + + public void setBeanName(String beanName) { + this.setName(beanName); + } + + public String toString() { + return (this.name != null) ? this.name : super.toString(); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/PollingSourceEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/PollingSourceEndpoint.java deleted file mode 100644 index f987a62797..0000000000 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/PollingSourceEndpoint.java +++ /dev/null @@ -1,158 +0,0 @@ -/* - * Copyright 2002-2008 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.endpoint; - -import java.lang.reflect.Method; -import java.util.List; - -import org.aopalliance.aop.Advice; - -import org.springframework.aop.framework.ProxyFactory; -import org.springframework.aop.support.StaticMethodMatcherPointcutAdvisor; -import org.springframework.beans.factory.InitializingBean; -import org.springframework.integration.channel.DispatcherPolicy; -import org.springframework.integration.channel.MessageChannel; -import org.springframework.integration.dispatcher.DefaultPollingDispatcher; -import org.springframework.integration.dispatcher.PollingDispatcher; -import org.springframework.integration.dispatcher.PollingDispatcherTask; -import org.springframework.integration.message.Source; -import org.springframework.integration.scheduling.MessagingTask; -import org.springframework.integration.scheduling.Schedule; -import org.springframework.util.Assert; - -/** - * A channel adapter that retrieves messages from a {@link Source} - * and then sends the resulting messages to the provided {@link MessageChannel}. - * - * @author Mark Fisher - */ -public class PollingSourceEndpoint extends AbstractSourceEndpoint implements MessagingTask, InitializingBean { - - private final Schedule schedule; - - private final DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); - - private volatile PollingDispatcher dispatcher; - - private volatile List dispatchAdviceChain; - - private volatile MessagingTask task; - - private volatile List taskAdviceChain; - - private volatile boolean taskInitialized; - - private final Object taskMonitor = new Object(); - - - public PollingSourceEndpoint(Source source, MessageChannel channel, Schedule schedule) { - super(source, channel); - Assert.notNull(schedule, "schedule must not be null"); - this.dispatcher = new DefaultPollingDispatcher(source, this.dispatcherPolicy); - this.dispatcher.subscribe(this.getChannel()); - this.schedule = schedule; - } - - - public void setMaxMessagesPerTask(int maxMessagesPerTask) { - this.dispatcherPolicy.setMaxMessagesPerTask(maxMessagesPerTask); - } - - public void setSendTimeout(long sendTimeout) { - this.dispatcher.setSendTimeout(sendTimeout); - } - - public void setTaskAdviceChain(List taskAdviceChain) { - this.taskAdviceChain = taskAdviceChain; - } - - public void setDispatchAdviceChain(List dispatchAdviceChain) { - this.dispatchAdviceChain = dispatchAdviceChain; - } - - public Schedule getSchedule() { - return this.schedule; - } - - public void afterPropertiesSet() { - this.initializeTask(); - } - - public void initializeTask() { - synchronized (this.taskMonitor) { - if (this.taskInitialized) { - return; - } - this.refreshTask(); - this.taskInitialized = true; - } - } - - public void refreshTask() { - synchronized (this.taskMonitor) { - PollingDispatcher dispatcherProxy = null; - if (this.dispatchAdviceChain != null && this.dispatchAdviceChain.size() > 0) { - ProxyFactory proxyFactory = new ProxyFactory(this.dispatcher); - proxyFactory.setInterfaces(new Class[] { PollingDispatcher.class }); - for (Advice advice : this.dispatchAdviceChain) { - proxyFactory.addAdvisor(new MethodNameAdvisor(advice, "dispatch")); - } - dispatcherProxy = (PollingDispatcher) proxyFactory.getProxy(); - } - this.task = new PollingDispatcherTask((dispatcherProxy != null) ? dispatcherProxy : this.dispatcher, this.schedule); - if (this.taskAdviceChain != null && this.taskAdviceChain.size() > 0) { - ProxyFactory proxyFactory = new ProxyFactory(this.task); - proxyFactory.setInterfaces(new Class[] { MessagingTask.class }); - for (Advice advice : this.taskAdviceChain) { - proxyFactory.addAdvisor(new MethodNameAdvisor(advice, "run")); - } - this.task = (MessagingTask) proxyFactory.getProxy(); - } - } - } - - private MessagingTask getTask() { - synchronized (this.taskMonitor) { - if (!this.taskInitialized) { - this.initializeTask(); - } - return this.task; - } - } - - public void run() { - this.getTask().run(); - } - - - @SuppressWarnings("serial") - private static class MethodNameAdvisor extends StaticMethodMatcherPointcutAdvisor { - - private final String methodName; - - MethodNameAdvisor(Advice advice, String methodName) { - super(advice); - this.methodName = methodName; - } - - @SuppressWarnings("unchecked") - public boolean matches(Method method, Class targetClass) { - return method.getName().equals(methodName); - } - } - -} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourceEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourceEndpoint.java index 7ce7c276d5..2a1f069fbc 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourceEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourceEndpoint.java @@ -16,15 +16,144 @@ package org.springframework.integration.endpoint; +import java.lang.reflect.Method; +import java.util.List; + +import org.aopalliance.aop.Advice; + +import org.springframework.aop.framework.ProxyFactory; +import org.springframework.aop.support.StaticMethodMatcherPointcutAdvisor; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.integration.channel.DispatcherPolicy; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.dispatcher.DefaultPollingDispatcher; +import org.springframework.integration.dispatcher.PollingDispatcher; +import org.springframework.integration.dispatcher.PollingDispatcherTask; import org.springframework.integration.message.Source; +import org.springframework.integration.scheduling.MessagingTask; +import org.springframework.integration.scheduling.Schedule; +import org.springframework.util.Assert; /** - * Base interface for source endpoints. + * A channel adapter that retrieves messages from a {@link Source} + * and then sends the resulting messages to the provided {@link MessageChannel}. * * @author Mark Fisher */ -public interface SourceEndpoint extends MessageEndpoint { +public class SourceEndpoint extends AbstractEndpoint implements MessagingTask, InitializingBean { - Source getSource(); + private final Schedule schedule; + + private final DispatcherPolicy dispatcherPolicy = new DispatcherPolicy(); + + private volatile PollingDispatcher dispatcher; + + private volatile List dispatchAdviceChain; + + private volatile MessagingTask task; + + private volatile List taskAdviceChain; + + private volatile boolean taskInitialized; + + private final Object taskMonitor = new Object(); + + + public SourceEndpoint(Source source, MessageChannel channel, Schedule schedule) { + Assert.notNull(source, "source must not be null"); + Assert.notNull(channel, "channel must not be null"); + Assert.notNull(schedule, "schedule must not be null"); + this.dispatcher = new DefaultPollingDispatcher(source, this.dispatcherPolicy); + this.dispatcher.subscribe(channel); + this.schedule = schedule; + } + + + public void setMaxMessagesPerTask(int maxMessagesPerTask) { + this.dispatcherPolicy.setMaxMessagesPerTask(maxMessagesPerTask); + } + + public void setSendTimeout(long sendTimeout) { + this.dispatcher.setSendTimeout(sendTimeout); + } + + public void setTaskAdviceChain(List taskAdviceChain) { + this.taskAdviceChain = taskAdviceChain; + } + + public void setDispatchAdviceChain(List dispatchAdviceChain) { + this.dispatchAdviceChain = dispatchAdviceChain; + } + + public Schedule getSchedule() { + return this.schedule; + } + + public void afterPropertiesSet() { + this.initializeTask(); + } + + public void initializeTask() { + synchronized (this.taskMonitor) { + if (this.taskInitialized) { + return; + } + this.refreshTask(); + this.taskInitialized = true; + } + } + + public void refreshTask() { + synchronized (this.taskMonitor) { + PollingDispatcher dispatcherProxy = null; + if (this.dispatchAdviceChain != null && this.dispatchAdviceChain.size() > 0) { + ProxyFactory proxyFactory = new ProxyFactory(this.dispatcher); + proxyFactory.setInterfaces(new Class[] { PollingDispatcher.class }); + for (Advice advice : this.dispatchAdviceChain) { + proxyFactory.addAdvisor(new MethodNameAdvisor(advice, "dispatch")); + } + dispatcherProxy = (PollingDispatcher) proxyFactory.getProxy(); + } + this.task = new PollingDispatcherTask((dispatcherProxy != null) ? dispatcherProxy : this.dispatcher, this.schedule); + if (this.taskAdviceChain != null && this.taskAdviceChain.size() > 0) { + ProxyFactory proxyFactory = new ProxyFactory(this.task); + proxyFactory.setInterfaces(new Class[] { MessagingTask.class }); + for (Advice advice : this.taskAdviceChain) { + proxyFactory.addAdvisor(new MethodNameAdvisor(advice, "run")); + } + this.task = (MessagingTask) proxyFactory.getProxy(); + } + } + } + + private MessagingTask getTask() { + synchronized (this.taskMonitor) { + if (!this.taskInitialized) { + this.initializeTask(); + } + return this.task; + } + } + + public void run() { + this.getTask().run(); + } + + + @SuppressWarnings("serial") + private static class MethodNameAdvisor extends StaticMethodMatcherPointcutAdvisor { + + private final String methodName; + + MethodNameAdvisor(Advice advice, String methodName) { + super(advice); + this.methodName = methodName; + } + + @SuppressWarnings("unchecked") + public boolean matches(Method method, Class targetClass) { + return method.getName().equals(methodName); + } + } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/TargetEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/TargetEndpoint.java index 3bd9be284c..b69753cd27 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/TargetEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/TargetEndpoint.java @@ -25,9 +25,6 @@ import java.util.concurrent.SynchronousQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; - import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.context.Lifecycle; @@ -44,15 +41,11 @@ import org.springframework.integration.util.ErrorHandler; import org.springframework.util.Assert; /** - * Base class for {@link MessageEndpoint} implementations. + * Base class for {@link MessageEndpoint} implementations to which Messages may be sent. * * @author Mark Fisher */ -public class TargetEndpoint implements Target, MessageEndpoint, ChannelRegistryAware, InitializingBean, Lifecycle { - - protected final Log logger = LogFactory.getLog(this.getClass()); - - private volatile String name; +public class TargetEndpoint extends AbstractEndpoint implements Target, ChannelRegistryAware, InitializingBean, Lifecycle { private volatile Target target; @@ -80,18 +73,6 @@ public class TargetEndpoint implements Target, MessageEndpoint, ChannelRegistryA } - public String getName() { - return this.name; - } - - public void setName(String name) { - this.name = name; - } - - public void setBeanName(String beanName) { - this.setName(beanName); - } - public Target getTarget() { return this.target; } @@ -227,8 +208,4 @@ public class TargetEndpoint implements Target, MessageEndpoint, ChannelRegistryA } } - public String toString() { - return (this.name != null) ? this.name : super.toString(); - } - } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java index 35a912a58a..04c46d5427 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java @@ -33,7 +33,7 @@ import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.channel.RendezvousChannel; import org.springframework.integration.endpoint.ConcurrencyPolicy; -import org.springframework.integration.endpoint.PollingSourceEndpoint; +import org.springframework.integration.endpoint.SourceEndpoint; import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.ErrorMessage; import org.springframework.integration.message.GenericMessage; @@ -172,7 +172,7 @@ public class MessageBusTests { public void testErrorChannelWithFailedDispatch() throws InterruptedException { MessageBus bus = new MessageBus(); CountDownLatch latch = new CountDownLatch(1); - PollingSourceEndpoint sourceEndpoint = new PollingSourceEndpoint(new FailingSource(latch), new QueueChannel(), new PollingSchedule(1000)); + SourceEndpoint sourceEndpoint = new SourceEndpoint(new FailingSource(latch), new QueueChannel(), new PollingSchedule(1000)); bus.registerEndpoint("testEndpoint", sourceEndpoint); bus.start(); latch.await(1000, TimeUnit.MILLISECONDS); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingSourceEndpointTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/SourceEndpointTests.java similarity index 93% rename from spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingSourceEndpointTests.java rename to spring-integration-core/src/test/java/org/springframework/integration/endpoint/SourceEndpointTests.java index 3e25192f8e..52f088c09f 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingSourceEndpointTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/SourceEndpointTests.java @@ -41,14 +41,14 @@ import org.springframework.integration.scheduling.PollingSchedule; /** * @author Mark Fisher */ -public class PollingSourceEndpointTests { +public class SourceEndpointTests { @Test public void testPolledSourceSendsToChannel() { TestSource source = new TestSource("testing", 1); QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(100); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.run(); Message message = channel.receive(1000); assertNotNull("message should not be null", message); @@ -61,7 +61,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(1); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setSendTimeout(10); endpoint.run(); Message message1 = channel.receive(1000); @@ -82,7 +82,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); endpoint.setMaxMessagesPerTask(5); endpoint.run(); Message message1 = channel.receive(0); @@ -104,7 +104,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); final StringBuffer buffer = new StringBuffer(); List taskAdviceChain = new ArrayList(); taskAdviceChain.add(new MethodBeforeAdvice() { @@ -138,7 +138,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); final StringBuffer buffer = new StringBuffer(); List dispatchAdviceChain = new ArrayList(); dispatchAdviceChain.add(new MethodBeforeAdvice() { @@ -172,7 +172,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); List dispatchAdviceChain = new ArrayList(); List taskAdviceChain = new ArrayList(); final StringBuffer buffer = new StringBuffer(); @@ -216,7 +216,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); List dispatchAdviceChain = new ArrayList(); List taskAdviceChain = new ArrayList(); final StringBuffer buffer = new StringBuffer(); @@ -278,7 +278,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); List dispatchAdviceChain = new ArrayList(); List taskAdviceChain = new ArrayList(); final StringBuffer buffer = new StringBuffer(); @@ -317,7 +317,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); List dispatchAdviceChain = new ArrayList(); List taskAdviceChain = new ArrayList(); final StringBuffer buffer = new StringBuffer(); @@ -356,7 +356,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); List dispatchAdviceChain = new ArrayList(); List taskAdviceChain = new ArrayList(); final StringBuffer buffer = new StringBuffer(); @@ -396,7 +396,7 @@ public class PollingSourceEndpointTests { QueueChannel channel = new QueueChannel(); PollingSchedule schedule = new PollingSchedule(1000); schedule.setInitialDelay(10000); - PollingSourceEndpoint endpoint = new PollingSourceEndpoint(source, channel, schedule); + SourceEndpoint endpoint = new SourceEndpoint(source, channel, schedule); List dispatchAdviceChain = new ArrayList(); List taskAdviceChain = new ArrayList(); final StringBuffer buffer = new StringBuffer(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/handler/adapterTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/handler/adapterTests.xml index 7f84348b71..26d98c8001 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/handler/adapterTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/handler/adapterTests.xml @@ -10,7 +10,7 @@ - + diff --git a/spring-integration-samples/src/main/java/org/springframework/integration/samples/cafe/Barista.java b/spring-integration-samples/src/main/java/org/springframework/integration/samples/cafe/Barista.java index 89919ece3a..81bcbde014 100644 --- a/spring-integration-samples/src/main/java/org/springframework/integration/samples/cafe/Barista.java +++ b/spring-integration-samples/src/main/java/org/springframework/integration/samples/cafe/Barista.java @@ -26,9 +26,9 @@ import org.springframework.stereotype.Component; @Component public class Barista { - private long hotDrinkDelay = 1000; + private long hotDrinkDelay = 2000; - private long coldDrinkDelay = 700; + private long coldDrinkDelay = 1000; private AtomicInteger hotDrinkCounter = new AtomicInteger();