From ec98fa35c4a2358c43d4eba1c22c7ebc6a133b4a Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Tue, 11 Dec 2007 07:33:06 +0000 Subject: [PATCH] Added basic support for annotated endpoints. --- .../{annotation => aop}/Publisher.java | 2 +- .../aop/PublisherAnnotationAdvisor.java | 1 - .../aop/PublisherAnnotationPostProcessor.java | 1 - .../integration/bus/MessageBus.java | 4 + .../endpoint/MessageHandlerAdapter.java | 13 +- .../AnnotationAwareMessageEndpoint.java | 155 ++++++++++++++++++ .../annotation/DefaultOutput.java} | 4 +- .../{ => endpoint}/annotation/Handler.java | 2 +- .../annotation/MessageEndpoint.java | 8 +- .../annotation/Polled.java} | 6 +- .../aop/PublisherAnnotationAdvisorTests.java | 1 - .../aop/PublisherAnnotationTestBean.java | 1 - .../AnnotationAwareMessageEndpointTests.java | 51 ++++++ 13 files changed, 235 insertions(+), 14 deletions(-) rename spring-eai-core/src/main/java/org/springframework/integration/{annotation => aop}/Publisher.java (95%) create mode 100644 spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpoint.java rename spring-eai-core/src/main/java/org/springframework/integration/{annotation/Target.java => endpoint/annotation/DefaultOutput.java} (93%) rename spring-eai-core/src/main/java/org/springframework/integration/{ => endpoint}/annotation/Handler.java (95%) rename spring-eai-core/src/main/java/org/springframework/integration/{ => endpoint}/annotation/MessageEndpoint.java (88%) rename spring-eai-core/src/main/java/org/springframework/integration/{annotation/Source.java => endpoint/annotation/Polled.java} (91%) create mode 100644 spring-eai-core/src/test/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpointTests.java diff --git a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Publisher.java b/spring-eai-core/src/main/java/org/springframework/integration/aop/Publisher.java similarity index 95% rename from spring-eai-core/src/main/java/org/springframework/integration/annotation/Publisher.java rename to spring-eai-core/src/main/java/org/springframework/integration/aop/Publisher.java index 804f3da6e8..b5096b43c9 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Publisher.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/aop/Publisher.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.integration.annotation; +package org.springframework.integration.aop; import java.lang.annotation.Documented; import java.lang.annotation.ElementType; diff --git a/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationAdvisor.java b/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationAdvisor.java index 8d2323bee9..3b12453573 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationAdvisor.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationAdvisor.java @@ -23,7 +23,6 @@ import org.aopalliance.aop.Advice; import org.springframework.aop.Pointcut; import org.springframework.aop.support.AbstractPointcutAdvisor; import org.springframework.aop.support.annotation.AnnotationMatchingPointcut; -import org.springframework.integration.annotation.Publisher; import org.springframework.integration.channel.ChannelMapping; import org.springframework.util.Assert; diff --git a/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationPostProcessor.java b/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationPostProcessor.java index b2130efcbe..672dc3219f 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationPostProcessor.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/aop/PublisherAnnotationPostProcessor.java @@ -25,7 +25,6 @@ import org.springframework.aop.support.AopUtils; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanClassLoaderAware; import org.springframework.beans.factory.config.BeanPostProcessor; -import org.springframework.integration.annotation.Publisher; import org.springframework.integration.channel.ChannelMapping; import org.springframework.util.Assert; diff --git a/spring-eai-core/src/main/java/org/springframework/integration/bus/MessageBus.java b/spring-eai-core/src/main/java/org/springframework/integration/bus/MessageBus.java index f5da8bce59..6d7631ea2d 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/bus/MessageBus.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/bus/MessageBus.java @@ -72,6 +72,10 @@ public class MessageBus implements ChannelMapping, ApplicationContextAware, Life this.activateSubscriptions(applicationContext); } + public void setAutoCreateChannels(boolean autoCreateChannels) { + this.autoCreateChannels = autoCreateChannels; + } + @SuppressWarnings("unchecked") private void registerChannels(ApplicationContext context) { Map channelBeans = diff --git a/spring-eai-core/src/main/java/org/springframework/integration/endpoint/MessageHandlerAdapter.java b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/MessageHandlerAdapter.java index 692f3a8a64..944612c8dc 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/endpoint/MessageHandlerAdapter.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/MessageHandlerAdapter.java @@ -17,6 +17,7 @@ package org.springframework.integration.endpoint; import org.springframework.beans.factory.InitializingBean; +import org.springframework.core.Ordered; import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageMapper; @@ -31,7 +32,7 @@ import org.springframework.util.Assert; * * @author Mark Fisher */ -public class MessageHandlerAdapter implements MessageHandler, InitializingBean { +public class MessageHandlerAdapter implements MessageHandler, Ordered, InitializingBean { private T object; @@ -41,6 +42,8 @@ public class MessageHandlerAdapter implements MessageHandler, InitializingBea private SimpleMethodInvoker invoker; + private int order = Integer.MAX_VALUE; + public void setObject(T object) { Assert.notNull(object, "'object' must not be null"); @@ -52,6 +55,14 @@ public class MessageHandlerAdapter implements MessageHandler, InitializingBea this.method = method; } + public void setOrder(int order) { + this.order = order; + } + + public int getOrder() { + return this.order; + } + public void afterPropertiesSet() { this.invoker = new SimpleMethodInvoker(this.object, this.method); } diff --git a/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpoint.java b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpoint.java new file mode 100644 index 0000000000..a04e69004b --- /dev/null +++ b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpoint.java @@ -0,0 +1,155 @@ +/* + * 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.endpoint.annotation; + +import java.lang.annotation.Annotation; +import java.lang.reflect.Method; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.core.OrderComparator; +import org.springframework.core.annotation.AnnotationUtils; +import org.springframework.core.annotation.Order; +import org.springframework.integration.bus.ConsumerPolicy; +import org.springframework.integration.bus.MessageBus; +import org.springframework.integration.endpoint.GenericMessageEndpoint; +import org.springframework.integration.endpoint.InboundMethodInvokingChannelAdapter; +import org.springframework.integration.endpoint.MessageHandlerAdapter; +import org.springframework.integration.endpoint.OutboundMethodInvokingChannelAdapter; +import org.springframework.integration.handler.MessageHandler; +import org.springframework.integration.handler.MessageHandlerChain; +import org.springframework.util.Assert; +import org.springframework.util.ClassUtils; +import org.springframework.util.ReflectionUtils; +import org.springframework.util.StringUtils; + +/** + * An endpoint implementation for classes annotated with {@link MessageEndpoint @MessageEndpoint}. + * + * @author Mark Fisher + */ +public class AnnotationAwareMessageEndpoint extends GenericMessageEndpoint implements InitializingBean { + + private Object targetObject; + + private MessageBus messageBus; + + + public AnnotationAwareMessageEndpoint(Object targetObject, MessageBus messageBus) { + Assert.notNull(targetObject, "targetObject must not be null"); + Assert.notNull(messageBus, "messageBus must not be null"); + this.targetObject = targetObject; + this.messageBus = messageBus; + } + + + public void afterPropertiesSet() { + this.configureInputChannel(); + this.configureDefaultOutputChannel(); + this.createHandlerChain(); + this.messageBus.registerEndpoint(ClassUtils.getShortNameAsProperty(targetObject.getClass()), this); + } + + private void configureInputChannel() { + MessageEndpoint endpointAnnotation = this.targetObject.getClass().getAnnotation(MessageEndpoint.class); + String channelName = endpointAnnotation.input(); + if (StringUtils.hasText(channelName)) { + this.setInputChannelName(channelName); + ConsumerPolicy consumerPolicy = new ConsumerPolicy(); + consumerPolicy.setPeriod(endpointAnnotation.pollPeriod()); + this.setConsumerPolicy(consumerPolicy); + } + else { + ReflectionUtils.doWithMethods(targetObject.getClass(), new ReflectionUtils.MethodCallback() { + public void doWith(Method method) throws IllegalArgumentException, IllegalAccessException { + Annotation annotation = AnnotationUtils.getAnnotation(method, Polled.class); + if (annotation != null) { + InboundMethodInvokingChannelAdapter adapter = + new InboundMethodInvokingChannelAdapter(); + adapter.setObject(targetObject); + adapter.setMethod(method.getName()); + adapter.afterPropertiesSet(); + String channelName = ClassUtils.getShortNameAsProperty( + targetObject.getClass()) + "-inputChannel"; + messageBus.registerChannel(channelName, adapter); + setInputChannelName(channelName); + return; + } + } + }); + } + } + + private void configureDefaultOutputChannel() { + MessageEndpoint endpointAnnotation = this.targetObject.getClass().getAnnotation(MessageEndpoint.class); + String channelName = endpointAnnotation.defaultOutput(); + if (StringUtils.hasText(channelName)) { + this.setDefaultOutputChannelName(channelName); + } + else { + ReflectionUtils.doWithMethods(targetObject.getClass(), new ReflectionUtils.MethodCallback() { + public void doWith(Method method) throws IllegalArgumentException, IllegalAccessException { + Annotation annotation = AnnotationUtils.getAnnotation(method, DefaultOutput.class); + if (annotation != null) { + OutboundMethodInvokingChannelAdapter adapter = + new OutboundMethodInvokingChannelAdapter(); + adapter.setObject(targetObject); + adapter.setMethod(method.getName()); + adapter.afterPropertiesSet(); + String channelName = ClassUtils.getShortNameAsProperty( + targetObject.getClass()) + "-defaultOutputChannel"; + messageBus.registerChannel(channelName, adapter); + setDefaultOutputChannelName(channelName); + return; + } + } + }); + } + } + + @SuppressWarnings("unchecked") + private void createHandlerChain() { + final List> handlers = new ArrayList>(); + ReflectionUtils.doWithMethods(targetObject.getClass(), new ReflectionUtils.MethodCallback() { + public void doWith(Method method) throws IllegalArgumentException, IllegalAccessException { + Annotation annotation = AnnotationUtils.getAnnotation(method, Handler.class); + if (annotation != null) { + MessageHandlerAdapter adapter = new MessageHandlerAdapter(); + adapter.setObject(targetObject); + adapter.setMethod(method.getName()); + adapter.afterPropertiesSet(); + Order orderAnnotation = (Order) AnnotationUtils.getAnnotation(method, Order.class); + if (orderAnnotation != null) { + adapter.setOrder(orderAnnotation.value()); + } + handlers.add(adapter); + } + } + }); + if (handlers.size() > 0) { + MessageHandlerChain handlerChain = new MessageHandlerChain(); + Collections.sort(handlers, new OrderComparator()); + for (MessageHandler handler : handlers) { + handlerChain.add(handler); + } + this.setHandler(handlerChain); + } + } + +} diff --git a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Target.java b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/DefaultOutput.java similarity index 93% rename from spring-eai-core/src/main/java/org/springframework/integration/annotation/Target.java rename to spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/DefaultOutput.java index 7d75d91899..0277425f2a 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Target.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/DefaultOutput.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.integration.annotation; +package org.springframework.integration.endpoint.annotation; import java.lang.annotation.Documented; import java.lang.annotation.ElementType; @@ -36,6 +36,6 @@ import org.springframework.integration.message.Message; @Retention(RetentionPolicy.RUNTIME) @Inherited @Documented -public @interface Target { +public @interface DefaultOutput { } diff --git a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Handler.java b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/Handler.java similarity index 95% rename from spring-eai-core/src/main/java/org/springframework/integration/annotation/Handler.java rename to spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/Handler.java index 63d0e1f3f5..ad42588372 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Handler.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/Handler.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.integration.annotation; +package org.springframework.integration.endpoint.annotation; import java.lang.annotation.Documented; import java.lang.annotation.ElementType; diff --git a/spring-eai-core/src/main/java/org/springframework/integration/annotation/MessageEndpoint.java b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/MessageEndpoint.java similarity index 88% rename from spring-eai-core/src/main/java/org/springframework/integration/annotation/MessageEndpoint.java rename to spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/MessageEndpoint.java index 4ddbefb3b2..b0b199e460 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/annotation/MessageEndpoint.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/MessageEndpoint.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.integration.annotation; +package org.springframework.integration.endpoint.annotation; import java.lang.annotation.Documented; import java.lang.annotation.ElementType; @@ -37,8 +37,10 @@ import org.springframework.stereotype.Component; @Component public @interface MessageEndpoint { - String input(); + String input() default ""; - String output(); + String defaultOutput() default ""; + + int pollPeriod() default 0; } diff --git a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Source.java b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/Polled.java similarity index 91% rename from spring-eai-core/src/main/java/org/springframework/integration/annotation/Source.java rename to spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/Polled.java index cd8cbe45ab..008e297777 100644 --- a/spring-eai-core/src/main/java/org/springframework/integration/annotation/Source.java +++ b/spring-eai-core/src/main/java/org/springframework/integration/endpoint/annotation/Polled.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.integration.annotation; +package org.springframework.integration.endpoint.annotation; import java.lang.annotation.Documented; import java.lang.annotation.ElementType; @@ -35,6 +35,8 @@ import java.lang.annotation.Target; @Retention(RetentionPolicy.RUNTIME) @Inherited @Documented -public @interface Source { +public @interface Polled { + + int period() default 100; } diff --git a/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationAdvisorTests.java b/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationAdvisorTests.java index 7739e4ee7b..e58f984a70 100644 --- a/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationAdvisorTests.java +++ b/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationAdvisorTests.java @@ -23,7 +23,6 @@ import static org.junit.Assert.assertNull; import org.junit.Test; import org.springframework.aop.framework.ProxyFactory; -import org.springframework.integration.annotation.Publisher; import org.springframework.integration.channel.ChannelMapping; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.PointToPointChannel; diff --git a/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationTestBean.java b/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationTestBean.java index b662e8c271..f27d499a52 100644 --- a/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationTestBean.java +++ b/spring-eai-core/src/test/java/org/springframework/integration/aop/PublisherAnnotationTestBean.java @@ -16,7 +16,6 @@ package org.springframework.integration.aop; -import org.springframework.integration.annotation.Publisher; /** * @author Mark Fisher diff --git a/spring-eai-core/src/test/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpointTests.java b/spring-eai-core/src/test/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpointTests.java new file mode 100644 index 0000000000..2bc5927b3e --- /dev/null +++ b/spring-eai-core/src/test/java/org/springframework/integration/endpoint/annotation/AnnotationAwareMessageEndpointTests.java @@ -0,0 +1,51 @@ +/* + * 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.endpoint.annotation; + +import org.junit.Test; + +import org.springframework.integration.bus.MessageBus; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.message.DocumentMessage; + +/** + * @author Mark Fisher + */ +public class AnnotationAwareMessageEndpointTests { + + @Test + public void testSimpleHandler() throws InterruptedException { + MessageBus messageBus = new MessageBus(); + messageBus.setAutoCreateChannels(true); + AnnotationAwareMessageEndpoint endpoint = + new AnnotationAwareMessageEndpoint(new SimpleEndpoint(), messageBus); + endpoint.afterPropertiesSet(); + MessageChannel channel = messageBus.getChannel("testChannel"); + messageBus.start(); + endpoint.afterPropertiesSet(); + channel.send(new DocumentMessage(1, "world"), 10); + } + + @MessageEndpoint(input="testChannel", pollPeriod=10) + public class SimpleEndpoint { + + @Handler + public void sayHello(String name) { + System.out.println("hello " + name); + } + } +}