Renamed Target to MessageTarget.

This commit is contained in:
Mark Fisher
2008-06-30 23:23:04 +00:00
parent 504b101213
commit 908966d975
26 changed files with 78 additions and 79 deletions

View File

@@ -58,7 +58,7 @@ import org.springframework.integration.endpoint.TargetEndpoint;
import org.springframework.integration.handler.MessageHandler;
import org.springframework.integration.message.CommandMessage;
import org.springframework.integration.message.PollCommand;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.scheduling.MessagePublishingErrorHandler;
import org.springframework.integration.scheduling.MessagingTask;
import org.springframework.integration.scheduling.MessagingTaskScheduler;
@@ -273,11 +273,11 @@ public class MessageBus implements ChannelRegistry, EndpointRegistry, Applicatio
this.doRegisterEndpoint(name, endpoint, subscription, concurrencyPolicy);
}
public void registerTarget(String name, Target target, Subscription subscription) {
public void registerTarget(String name, MessageTarget target, Subscription subscription) {
this.registerTarget(name, target, subscription, this.defaultConcurrencyPolicy);
}
public void registerTarget(String name, Target target, Subscription subscription,
public void registerTarget(String name, MessageTarget target, Subscription subscription,
ConcurrencyPolicy concurrencyPolicy) {
Assert.notNull(target, "'target' must not be null");
TargetEndpoint endpoint = new TargetEndpoint(target);

View File

@@ -32,7 +32,7 @@ import org.springframework.integration.dispatcher.DirectChannel;
import org.springframework.integration.dispatcher.PollingDispatcherTask;
import org.springframework.integration.endpoint.TargetEndpoint;
import org.springframework.integration.message.MessagingException;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.scheduling.MessagingTaskScheduler;
import org.springframework.integration.scheduling.PollingSchedule;
import org.springframework.integration.scheduling.Schedule;
@@ -40,7 +40,7 @@ import org.springframework.integration.util.ErrorHandler;
import org.springframework.util.Assert;
/**
* Manages subscriptions for {@link Target Targets} to a {@link MessageChannel}
* Manages subscriptions for {@link MessageTarget Targets} to a {@link MessageChannel}
* including the creation, scheduling, and lifecycle management of dispatchers.
*
* @author Mark Fisher
@@ -77,11 +77,11 @@ public class SubscriptionManager {
this.defaultSchedule = defaultSchedule;
}
public void addTarget(Target target) {
public void addTarget(MessageTarget target) {
this.addTarget(target, null);
}
public void addTarget(Target target, Schedule schedule) {
public void addTarget(MessageTarget target, Schedule schedule) {
Assert.notNull(target, "'target' must not be null");
if (schedule == null) {
schedule = this.defaultSchedule;
@@ -129,7 +129,7 @@ public class SubscriptionManager {
}
}
public boolean removeTarget(Target target) {
public boolean removeTarget(MessageTarget target) {
boolean removed = false;
Collection<PollingDispatcherTask> dispatcherTaskValues = this.dispatcherTasks.values();
for (PollingDispatcherTask dispatcherTask : dispatcherTaskValues) {

View File

@@ -31,7 +31,7 @@ import org.springframework.integration.bus.MessageBus;
import org.springframework.integration.channel.ChannelRegistryAware;
import org.springframework.integration.handler.MessageHandler;
import org.springframework.integration.message.MessageSource;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.util.Assert;
/**
@@ -66,7 +66,7 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Init
public void afterPropertiesSet() {
this.postProcessors.put(MessageHandler.class, new HandlerAnnotationPostProcessor(this.messageBus, this.beanClassLoader));
this.postProcessors.put(MessageSource.class, new SourceAnnotationPostProcessor(this.messageBus, this.beanClassLoader));
this.postProcessors.put(Target.class, new TargetAnnotationPostProcessor(this.messageBus, this.beanClassLoader));
this.postProcessors.put(MessageTarget.class, new TargetAnnotationPostProcessor(this.messageBus, this.beanClassLoader));
}
public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException {

View File

@@ -23,14 +23,13 @@ import java.util.List;
import org.springframework.core.annotation.AnnotationUtils;
import org.springframework.integration.ConfigurationException;
import org.springframework.integration.annotation.Concurrency;
import org.springframework.integration.annotation.MessageTarget;
import org.springframework.integration.annotation.Polled;
import org.springframework.integration.bus.MessageBus;
import org.springframework.integration.endpoint.ConcurrencyPolicy;
import org.springframework.integration.endpoint.MessageEndpoint;
import org.springframework.integration.endpoint.TargetEndpoint;
import org.springframework.integration.handler.MethodInvokingTarget;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.scheduling.Subscription;
/**
@@ -38,21 +37,21 @@ import org.springframework.integration.scheduling.Subscription;
*
* @author Mark Fisher
*/
public class TargetAnnotationPostProcessor extends AbstractAnnotationMethodPostProcessor<Target> {
public class TargetAnnotationPostProcessor extends AbstractAnnotationMethodPostProcessor<MessageTarget> {
public TargetAnnotationPostProcessor(MessageBus messageBus, ClassLoader beanClassLoader) {
super(MessageTarget.class, messageBus, beanClassLoader);
super(org.springframework.integration.annotation.MessageTarget.class, messageBus, beanClassLoader);
}
public Target processMethod(Object bean, Method method, Annotation annotation) {
public MessageTarget processMethod(Object bean, Method method, Annotation annotation) {
MethodInvokingTarget target = new MethodInvokingTarget();
target.setObject(bean);
target.setMethod(method);
return target;
}
protected Target processResults(List<Target> results) {
protected MessageTarget processResults(List<MessageTarget> results) {
if (results.size() > 1) {
throw new ConfigurationException("At most one @MessageTarget annotation is allowed per class.");
}
@@ -61,7 +60,7 @@ public class TargetAnnotationPostProcessor extends AbstractAnnotationMethodPostP
public MessageEndpoint createEndpoint(Object bean, String beanName, Class<?> originalBeanClass,
org.springframework.integration.annotation.MessageEndpoint endpointAnnotation) {
TargetEndpoint endpoint = new TargetEndpoint((Target) bean);
TargetEndpoint endpoint = new TargetEndpoint((MessageTarget) bean);
Polled polledAnnotation = AnnotationUtils.findAnnotation(originalBeanClass, Polled.class);
Subscription subscription = this.createSubscription(bean, beanName, endpointAnnotation, polledAnnotation);
endpoint.setSubscription(subscription);

View File

@@ -26,7 +26,7 @@ import org.springframework.integration.handler.MessageHandler;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageSource;
import org.springframework.integration.message.Subscribable;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.message.selector.MessageSelector;
/**
@@ -58,7 +58,7 @@ public class DirectChannel extends AbstractMessageChannel implements Subscribabl
}
public boolean subscribe(Target target) {
public boolean subscribe(MessageTarget target) {
boolean added = this.dispatcher.subscribe(target);
if (added) {
this.handlerCount.incrementAndGet();
@@ -66,7 +66,7 @@ public class DirectChannel extends AbstractMessageChannel implements Subscribabl
return added;
}
public boolean unsubscribe(Target target) {
public boolean unsubscribe(MessageTarget target) {
boolean removed = this.dispatcher.unsubscribe(target);
if (removed) {
this.handlerCount.decrementAndGet();

View File

@@ -19,7 +19,7 @@ package org.springframework.integration.dispatcher;
import org.springframework.integration.channel.MessageChannel;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.Subscribable;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.scheduling.MessagingTask;
import org.springframework.integration.scheduling.Schedule;
import org.springframework.util.Assert;
@@ -46,11 +46,11 @@ public class PollingDispatcherTask implements MessagingTask, Subscribable {
}
public boolean subscribe(Target target) {
public boolean subscribe(MessageTarget target) {
return this.dispatcher.subscribe(target);
}
public boolean unsubscribe(Target target) {
public boolean unsubscribe(MessageTarget target) {
return this.dispatcher.unsubscribe(target);
}

View File

@@ -31,7 +31,7 @@ import org.springframework.integration.message.BlockingTarget;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.Subscribable;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
/**
* Basic implementation of {@link MessageDispatcher}.
@@ -42,7 +42,7 @@ public class SimpleDispatcher implements MessageDispatcher, Subscribable {
protected final Log logger = LogFactory.getLog(this.getClass());
private final List<Target> targets = new CopyOnWriteArrayList<Target>();
private final List<MessageTarget> targets = new CopyOnWriteArrayList<MessageTarget>();
protected final DispatcherPolicy dispatcherPolicy;
@@ -58,17 +58,17 @@ public class SimpleDispatcher implements MessageDispatcher, Subscribable {
this.sendTimeout = sendTimeout;
}
public boolean subscribe(Target target) {
public boolean subscribe(MessageTarget target) {
return this.targets.add(target);
}
public boolean unsubscribe(Target target) {
public boolean unsubscribe(MessageTarget target) {
return this.targets.remove(target);
}
public boolean dispatch(Message<?> message) {
int attempts = 0;
List<Target> targetList = new ArrayList<Target>(this.targets);
List<MessageTarget> targetList = new ArrayList<MessageTarget>(this.targets);
while (attempts < this.dispatcherPolicy.getRejectionLimit()) {
if (attempts > 0) {
if (logger.isDebugEnabled()) {
@@ -84,7 +84,7 @@ public class SimpleDispatcher implements MessageDispatcher, Subscribable {
return false;
}
}
Iterator<Target> iter = targetList.iterator();
Iterator<MessageTarget> iter = targetList.iterator();
if (!iter.hasNext()) {
if (logger.isWarnEnabled()) {
logger.warn("no active targets");
@@ -93,7 +93,7 @@ public class SimpleDispatcher implements MessageDispatcher, Subscribable {
}
boolean rejected = false;
while (iter.hasNext()) {
Target target = iter.next();
MessageTarget target = iter.next();
try {
boolean sent = (target instanceof BlockingTarget && this.sendTimeout >= 0) ?
((BlockingTarget) target).send(message, this.sendTimeout) : target.send(message);

View File

@@ -27,29 +27,29 @@ import org.springframework.integration.handler.MessageHandlerNotRunningException
import org.springframework.integration.handler.MessageHandlerRejectedExecutionException;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.util.ErrorHandler;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.util.Assert;
/**
* A {@link Target} implementation that encapsulates an Executor and delegates
* A {@link MessageTarget} implementation that encapsulates an Executor and delegates
* to a wrapped target for concurrent, asynchronous message handling.
*
* @author Mark Fisher
*/
public class ConcurrentTarget implements Target, DisposableBean {
public class ConcurrentTarget implements MessageTarget, DisposableBean {
private final Log logger = LogFactory.getLog(this.getClass());
private final Target target;
private final MessageTarget target;
private final ExecutorService executor;
private volatile ErrorHandler errorHandler;
public ConcurrentTarget(Target target, ExecutorService executor) {
public ConcurrentTarget(MessageTarget target, ExecutorService executor) {
Assert.notNull(target, "'target' must not be null");
Assert.notNull(executor, "'executor' must not be null");
this.target = target;

View File

@@ -25,7 +25,7 @@ import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.MessageHandlingException;
import org.springframework.integration.message.MessageHeader;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
@@ -141,7 +141,7 @@ public class HandlerEndpoint extends TargetEndpoint {
}
private static class HandlerInvokingTarget implements Target {
private static class HandlerInvokingTarget implements MessageTarget {
private final MessageHandler handler;

View File

@@ -32,7 +32,7 @@ import org.springframework.integration.handler.MessageHandlerNotRunningException
import org.springframework.integration.handler.MessageHandlerRejectedExecutionException;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageHandlingException;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.message.selector.MessageSelector;
import org.springframework.integration.scheduling.Subscription;
import org.springframework.integration.util.ErrorHandler;
@@ -43,9 +43,9 @@ import org.springframework.util.Assert;
*
* @author Mark Fisher
*/
public class TargetEndpoint extends AbstractEndpoint implements Target, ChannelRegistryAware, InitializingBean, Lifecycle {
public class TargetEndpoint extends AbstractEndpoint implements MessageTarget, ChannelRegistryAware, InitializingBean, Lifecycle {
private volatile Target target;
private volatile MessageTarget target;
private volatile Subscription subscription;
@@ -65,17 +65,17 @@ public class TargetEndpoint extends AbstractEndpoint implements Target, ChannelR
public TargetEndpoint() {
}
public TargetEndpoint(Target target) {
public TargetEndpoint(MessageTarget target) {
Assert.notNull(target, "target must not be null");
this.target = target;
}
public Target getTarget() {
public MessageTarget getTarget() {
return this.target;
}
public void setTarget(Target target) {
public void setTarget(MessageTarget target) {
Assert.notNull(target, "target must not be null");
this.target = target;
}

View File

@@ -18,14 +18,14 @@ package org.springframework.integration.handler;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessagingException;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
/**
* A messaging target that invokes the specified method on the provided object.
*
* @author Mark Fisher
*/
public class MethodInvokingTarget extends AbstractMessageHandlerAdapter implements Target {
public class MethodInvokingTarget extends AbstractMessageHandlerAdapter implements MessageTarget {
public boolean send(Message<?> message) {
this.handle(message);

View File

@@ -17,11 +17,11 @@
package org.springframework.integration.message;
/**
* Extends {@link Target} and provides a timeout-aware send method.
* Extends {@link MessageTarget} and provides a timeout-aware send method.
*
* @author Mark Fisher
*/
public interface BlockingTarget extends Target {
public interface BlockingTarget extends MessageTarget {
/**
* Send a message, blocking indefinitely if necessary.

View File

@@ -21,7 +21,7 @@ package org.springframework.integration.message;
*
* @author Mark Fisher
*/
public interface Target {
public interface MessageTarget {
boolean send(Message<?> message);

View File

@@ -24,13 +24,13 @@ package org.springframework.integration.message;
public interface Subscribable {
/**
* Register a {@link Target} as a subscriber to this source.
* Register a {@link MessageTarget} as a subscriber to this source.
*/
boolean subscribe(Target target);
boolean subscribe(MessageTarget target);
/**
* Remove a {@link Target} from the subscribers of this source.
* Remove a {@link MessageTarget} from the subscribers of this source.
*/
boolean unsubscribe(Target target);
boolean unsubscribe(MessageTarget target);
}

View File

@@ -39,7 +39,7 @@ import org.springframework.integration.message.ErrorMessage;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.message.selector.PayloadTypeSelector;
import org.springframework.integration.scheduling.MessagePublishingErrorHandler;
import org.springframework.integration.scheduling.SimpleMessagingTaskScheduler;
@@ -157,7 +157,7 @@ public class SubscriptionManagerTests {
channel.send(new StringMessage(1, "test"));
SubscriptionManager manager = new SubscriptionManager(channel, scheduler);
manager.addTarget(createEndpoint(handler1, true));
manager.addTarget(new Target() {
manager.addTarget(new MessageTarget() {
public boolean send(Message<?> message) {
throw new MessageHandlerRejectedExecutionException(message);
}

View File

@@ -39,7 +39,7 @@ import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.MessagePriority;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
import org.springframework.integration.scheduling.SimpleMessagingTaskScheduler;
/**
@@ -272,7 +272,7 @@ public class ChannelParserTests {
}
private static class TestTarget implements Target {
private static class TestTarget implements MessageTarget {
private AtomicInteger counter;

View File

@@ -36,7 +36,7 @@ import org.springframework.integration.message.GenericMessage;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageHandlingException;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
/**
* @author Mark Fisher
@@ -109,7 +109,7 @@ public class EndpointParserTests {
public void testEndpointWithSelectorAccepts() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"endpointWithSelector.xml", this.getClass());
Target endpoint = (Target) context.getBean("endpoint");
MessageTarget endpoint = (MessageTarget) context.getBean("endpoint");
((Lifecycle) endpoint).start();
Message<?> message = new StringMessage("test");
MessageChannel replyChannel = new QueueChannel();
@@ -124,7 +124,7 @@ public class EndpointParserTests {
public void testEndpointWithSelectorRejects() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"endpointWithSelector.xml", this.getClass());
Target endpoint = (Target) context.getBean("endpoint");
MessageTarget endpoint = (MessageTarget) context.getBean("endpoint");
((Lifecycle) endpoint).start();
Message<?> message = new GenericMessage<Integer>(123);
MessageChannel replyChannel = new QueueChannel();
@@ -136,7 +136,7 @@ public class EndpointParserTests {
public void testCustomErrorHandler() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"endpointWithErrorHandler.xml", this.getClass());
Target endpoint = (Target) context.getBean("endpoint");
MessageTarget endpoint = (MessageTarget) context.getBean("endpoint");
TestErrorHandler errorHandler = (TestErrorHandler) context.getBean("errorHandler");
assertNull(errorHandler.getLastError());
Message<?> message = new StringMessage("test");
@@ -151,7 +151,7 @@ public class EndpointParserTests {
public void testCustomReplyHandler() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
"endpointWithReplyHandler.xml", this.getClass());
Target endpoint = (Target) context.getBean("endpoint");
MessageTarget endpoint = (MessageTarget) context.getBean("endpoint");
TestReplyHandler replyHandler = (TestReplyHandler) context.getBean("replyHandler");
assertNull(replyHandler.getLastMessage());
Message<?> message = new StringMessage("test");

View File

@@ -17,12 +17,12 @@
package org.springframework.integration.config;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
/**
* @author Mark Fisher
*/
public class TestTarget implements Target {
public class TestTarget implements MessageTarget {
public boolean send(Message<?> message) {
return true;

View File

@@ -30,7 +30,7 @@ import org.junit.Test;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageSource;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
/**
* @author Mark Fisher
@@ -124,7 +124,7 @@ public class DirectChannelTests {
}
private static class ThreadNameSettingTestTarget implements Target {
private static class ThreadNameSettingTestTarget implements MessageTarget {
private final CountDownLatch latch;

View File

@@ -29,7 +29,7 @@ import org.springframework.integration.endpoint.HandlerEndpoint;
import org.springframework.integration.handler.MessageHandler;
import org.springframework.integration.handler.TestHandlers;
import org.springframework.integration.message.StringMessage;
import org.springframework.integration.message.Target;
import org.springframework.integration.message.MessageTarget;
/**
* @author Mark Fisher
@@ -76,7 +76,7 @@ public class SimpleDispatcherTests {
}
private static Target createEndpoint(MessageHandler handler) {
private static MessageTarget createEndpoint(MessageHandler handler) {
HandlerEndpoint endpoint = new HandlerEndpoint(handler);
endpoint.start();
return endpoint;