From 7227544364fb13c393778b17fb48fd48724264fe Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Thu, 10 Jan 2008 21:55:10 +0000 Subject: [PATCH] Removed MessageReceiver and refactored to a more unified model with MessageHandler as the central interface. --- .../adapter/AbstractTargetAdapter.java | 5 +- .../integration/adapter/TargetAdapter.java | 4 +- .../bus/AbstractMessageDispatcher.java | 36 +++-- .../bus/DefaultMessageDispatcher.java | 48 +++--- .../integration/bus/MessageBus.java | 34 ++-- .../MessageHandlerNotRunningException.java} | 17 +- ...sageHandlerRejectedExecutionException.java | 37 +++++ .../bus/MessageSelectorRejectedException.java | 33 ++++ ...xecutor.java => PooledMessageHandler.java} | 51 +++--- .../endpoint/GenericMessageEndpoint.java | 7 +- .../integration/endpoint/MessageEndpoint.java | 6 +- .../bus/DefaultMessageDispatcherTests.java | 150 ++++++++++-------- .../bus/FixedDelayConsumerTests.java | 8 +- .../bus/FixedRateConsumerTests.java | 12 +- .../integration/bus/MessageBusTests.java | 6 +- 15 files changed, 286 insertions(+), 168 deletions(-) rename spring-integration-core/src/main/java/org/springframework/integration/{message/MessageReceiver.java => bus/MessageHandlerNotRunningException.java} (65%) create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/bus/MessageHandlerRejectedExecutionException.java create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/bus/MessageSelectorRejectedException.java rename spring-integration-core/src/main/java/org/springframework/integration/bus/{MessageReceivingExecutor.java => PooledMessageHandler.java} (77%) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java index 105684fdca..825f8861d4 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/adapter/AbstractTargetAdapter.java @@ -28,7 +28,7 @@ import org.springframework.util.Assert; * * @author Mark Fisher */ -public abstract class AbstractTargetAdapter implements TargetAdapter { +public abstract class AbstractTargetAdapter implements TargetAdapter { private String name; @@ -74,8 +74,9 @@ public abstract class AbstractTargetAdapter implements TargetAdapter { return this.policy; } - public final void messageReceived(Message message) { + public final Message handle(Message message) { this.sendToTarget(this.mapper.fromMessage(message)); + return null; } protected abstract boolean sendToTarget(T object); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/adapter/TargetAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/adapter/TargetAdapter.java index f9e3ffdab9..3e8ffdf627 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/adapter/TargetAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/adapter/TargetAdapter.java @@ -18,14 +18,14 @@ package org.springframework.integration.adapter; import org.springframework.integration.bus.ConsumerPolicy; import org.springframework.integration.channel.MessageChannel; -import org.springframework.integration.message.MessageReceiver; +import org.springframework.integration.handler.MessageHandler; /** * Base interface for target adapters. * * @author Mark Fisher */ -public interface TargetAdapter extends MessageReceiver { +public interface TargetAdapter extends MessageHandler { String getName(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/AbstractMessageDispatcher.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/AbstractMessageDispatcher.java index c7a89944d0..a3e5ca33af 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/bus/AbstractMessageDispatcher.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/AbstractMessageDispatcher.java @@ -22,7 +22,8 @@ import java.util.concurrent.CopyOnWriteArrayList; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; - +import org.springframework.context.Lifecycle; +import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.Message; import org.springframework.util.Assert; @@ -38,7 +39,7 @@ public abstract class AbstractMessageDispatcher implements MessageDispatcher { private MessageRetriever retriever; - private List executors = new CopyOnWriteArrayList(); + private List handlers = new CopyOnWriteArrayList(); private volatile boolean running; @@ -50,14 +51,16 @@ public abstract class AbstractMessageDispatcher implements MessageDispatcher { } - public void addExecutor(MessageReceivingExecutor executor) { - Assert.notNull(executor, "'executor' must not be null"); - executor.start(); - this.executors.add(executor); + public void addHandler(MessageHandler handler) { + Assert.notNull(handler, "'handler' must not be null"); + if (this.isRunning() && handler instanceof Lifecycle) { + ((Lifecycle) handler).start(); + } + this.handlers.add(handler); } - protected List getExecutors() { - return this.executors; + protected List getHandlers() { + return this.handlers; } public boolean isRunning() { @@ -67,9 +70,12 @@ public abstract class AbstractMessageDispatcher implements MessageDispatcher { public void start() { synchronized (this.lifecycleMonitor) { if (!this.isRunning()) { - for (MessageReceivingExecutor executor : this.executors) { - executor.start(); + for (MessageHandler handler : this.handlers) { + if (handler instanceof Lifecycle) { + ((Lifecycle) handler).start(); + } } + this.running = true; } } } @@ -77,9 +83,12 @@ public abstract class AbstractMessageDispatcher implements MessageDispatcher { public void stop() { synchronized (this.lifecycleMonitor) { if (this.isRunning()) { - for (MessageReceivingExecutor executor : this.executors) { - executor.stop(); + for (MessageHandler handler : this.handlers) { + if (handler instanceof Lifecycle) { + ((Lifecycle) handler).stop(); + } } + this.running = false; } } } @@ -90,6 +99,9 @@ public abstract class AbstractMessageDispatcher implements MessageDispatcher { * @return the number of messages processed */ public int dispatch() { + if (!this.isRunning()) { + return 0; + } int messagesProcessed = 0; Collection> messages = this.retriever.retrieveMessages(); if (messages == null) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/DefaultMessageDispatcher.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/DefaultMessageDispatcher.java index ba64012fcb..f3cf3dc5a7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/bus/DefaultMessageDispatcher.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/DefaultMessageDispatcher.java @@ -19,17 +19,18 @@ package org.springframework.integration.bus; import java.util.ArrayList; import java.util.Iterator; import java.util.List; -import java.util.concurrent.RejectedExecutionException; import org.springframework.integration.MessageDeliveryException; +import org.springframework.integration.MessageHandlingException; +import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.Message; import org.springframework.util.Assert; /** * The default implementation of {@link MessageDispatcher}. If * {@link #broadcast} is set to false (the default), each message - * will be sent to a single {@link MessageReceivingExecutor}. Otherwise, each - * retrieved {@link Message} will be sent to all executors. + * will be sent to a single {@link MessageHandler}. Otherwise, each + * retrieved {@link Message} will be sent to all handlers. * * @author Mark Fisher */ @@ -74,11 +75,11 @@ public class DefaultMessageDispatcher extends AbstractMessageDispatcher { @Override protected boolean dispatchMessage(Message message) { int attempts = 0; - List targets = new ArrayList(this.getExecutors()); + List targets = new ArrayList(this.getHandlers()); while (attempts < this.rejectionLimit) { if (attempts > 0) { if (logger.isDebugEnabled()) { - logger.debug("executor(s) rejected message after " + attempts + logger.debug("handler(s) rejected message after " + attempts + " attempt(s), will try again after 'retryInterval' of " + this.retryInterval + " milliseconds"); } @@ -90,38 +91,41 @@ public class DefaultMessageDispatcher extends AbstractMessageDispatcher { return false; } } - Iterator iter = targets.iterator(); + Iterator iter = targets.iterator(); if (!iter.hasNext()) { if (logger.isWarnEnabled()) { - logger.warn("dispatcher has no active executors"); + logger.warn("dispatcher has no active handlers"); } return false; } - boolean encounteredRejection = false; + boolean encounteredHandlerException = false; while (iter.hasNext()) { - MessageReceivingExecutor executor = iter.next(); - if (executor == null || !executor.isRunning()) { - if (logger.isInfoEnabled()) { - logger.info("skipping inactive executor"); - } - iter.remove(); - continue; - } + MessageHandler handler = iter.next(); try { - boolean accepted = executor.acceptMessage(message); - if (accepted && !this.broadcast) { + handler.handle(message); + if (!this.broadcast) { return true; } iter.remove(); } - catch (RejectedExecutionException rex) { - encounteredRejection = true; + catch (MessageSelectorRejectedException e) { if (logger.isDebugEnabled()) { - logger.debug("executor rejected task, continuing with other executors if available", rex); + logger.debug("selector rejected task, continuing with other handlers if available", e); + } + } + catch (MessageHandlerNotRunningException e) { + if (logger.isDebugEnabled()) { + logger.debug("handler not running, continuing with other handlers if available", e); + } + } + catch (MessageHandlingException e) { + encounteredHandlerException = true; + if (logger.isDebugEnabled()) { + logger.debug("handler threw exception, continuing with other handlers if available", e); } } } - if (!encounteredRejection) { + if (!encounteredHandlerException) { return true; } attempts++; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java index 0901520605..cc5c5343b0 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java @@ -38,8 +38,8 @@ import org.springframework.integration.channel.DefaultChannelRegistry; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.SimpleChannel; import org.springframework.integration.endpoint.MessageEndpoint; +import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.ErrorMessage; -import org.springframework.integration.message.MessageReceiver; import org.springframework.scheduling.concurrent.CustomizableThreadFactory; import org.springframework.util.Assert; import org.springframework.util.StringUtils; @@ -56,7 +56,7 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif private ChannelRegistry channelRegistry = new DefaultChannelRegistry(); - private Map> receivers = new ConcurrentHashMap>(); + private Map handlers = new ConcurrentHashMap(); private List dispatchers = new CopyOnWriteArrayList(); @@ -180,11 +180,11 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif this.channelRegistry.registerChannel(name, channel); } - public void registerEndpoint(String name, MessageEndpoint endpoint) { + public void registerEndpoint(String name, MessageEndpoint endpoint) { Assert.notNull(name, "'name' must not be null"); Assert.notNull(endpoint, "'endpoint' must not be null"); endpoint.setName(name); - this.receivers.put(name, endpoint); + this.handlers.put(name, endpoint); if (logger.isInfoEnabled()) { logger.info("registered endpoint '" + name + "'"); } @@ -217,11 +217,11 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif } } - public void registerTargetAdapter(String name, TargetAdapter targetAdapter) { + public void registerTargetAdapter(String name, TargetAdapter targetAdapter) { if (targetAdapter instanceof AbstractTargetAdapter) { AbstractTargetAdapter adapter = (AbstractTargetAdapter) targetAdapter; adapter.setName(name); - this.receivers.put(name, targetAdapter); + this.handlers.put(name, targetAdapter); MessageChannel channel = adapter.getChannel(); ConsumerPolicy policy = adapter.getConsumerPolicy(); this.doActivate(channel, adapter, policy); @@ -233,11 +233,11 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif public void activateSubscription(Subscription subscription) { String channelName = subscription.getChannel(); - String receiverName = subscription.getReceiver(); + String handlerName = subscription.getReceiver(); ConsumerPolicy policy = subscription.getPolicy(); - MessageReceiver receiver = this.receivers.get(receiverName); - if (receiver == null) { - throw new MessagingException("Cannot activate subscription, unknown receiver '" + receiverName + "'"); + MessageHandler handler = this.handlers.get(handlerName); + if (handler == null) { + throw new MessagingException("Cannot activate subscription, unknown handler '" + handlerName + "'"); } MessageChannel channel = this.lookupChannel(channelName); if (channel == null) { @@ -251,29 +251,29 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif channel = new SimpleChannel(); this.registerChannel(channelName, channel); } - this.doActivate(channel, receiver, policy); + this.doActivate(channel, handler, policy); if (logger.isInfoEnabled()) { logger.info("activated subscription to channel '" + channelName + - "' for receiver '" + receiverName + "'"); + "' for handler '" + handlerName + "'"); } } - private void doActivate(MessageChannel channel, MessageReceiver receiver, ConsumerPolicy policy) { - MessageReceivingExecutor executor = new MessageReceivingExecutor(receiver, policy.getConcurrency(), policy.getMaxConcurrency()); + private void doActivate(MessageChannel channel, MessageHandler handler, ConsumerPolicy policy) { + PooledMessageHandler pooledHandler = new PooledMessageHandler(handler, policy.getConcurrency(), policy.getMaxConcurrency()); MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); dispatcher.setRejectionLimit(policy.getRejectionLimit()); dispatcher.setRetryInterval(policy.getRetryInterval()); - dispatcher.addExecutor(executor); + dispatcher.addHandler(pooledHandler); DispatcherTask dispatcherTask = new DispatcherTask(dispatcher, policy); if (this.isRunning()) { - executor.start(); + dispatcher.start(); } this.dispatchers.add(dispatcher); this.addDispatcherTask(dispatcherTask); if (this.logger.isInfoEnabled()) { logger.info("registered dispatcher task: channel='" + - channel.getName() + "' receiver='" + receiver.getName() + "'"); + channel.getName() + "' handler='" + handler + "'"); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/message/MessageReceiver.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageHandlerNotRunningException.java similarity index 65% rename from spring-integration-core/src/main/java/org/springframework/integration/message/MessageReceiver.java rename to spring-integration-core/src/main/java/org/springframework/integration/bus/MessageHandlerNotRunningException.java index edbb8891f0..2406f4b4ce 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/message/MessageReceiver.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageHandlerNotRunningException.java @@ -14,20 +14,19 @@ * limitations under the License. */ -package org.springframework.integration.message; +package org.springframework.integration.bus; + +import org.springframework.integration.MessageHandlingException; /** - * The primary callback interface for any component capable of receiving - * messages. This includes message endpoints as well as target adapters. + * An exception indicating that a handler is not currently running. * * @author Mark Fisher */ -public interface MessageReceiver { +public class MessageHandlerNotRunningException extends MessageHandlingException { - String getName(); - - void setName(String name); - - void messageReceived(Message message); + public MessageHandlerNotRunningException() { + super("handler is not running"); + } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageHandlerRejectedExecutionException.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageHandlerRejectedExecutionException.java new file mode 100644 index 0000000000..b2e98ca635 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageHandlerRejectedExecutionException.java @@ -0,0 +1,37 @@ +/* + * 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.bus; + +import org.springframework.integration.MessageHandlingException; + +/** + * An exception indicating that a message was rejected by a handler; typically + * this would be the result of a thread pool executor rejecting a handler task. + * + * @author Mark Fisher + */ +public class MessageHandlerRejectedExecutionException extends MessageHandlingException { + + public MessageHandlerRejectedExecutionException(Throwable cause) { + super("handler rejected execution", cause); + } + + public MessageHandlerRejectedExecutionException(String message, Throwable cause) { + super(message, cause); + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageSelectorRejectedException.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageSelectorRejectedException.java new file mode 100644 index 0000000000..77e4dccbb2 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageSelectorRejectedException.java @@ -0,0 +1,33 @@ +/* + * 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.bus; + +import org.springframework.integration.MessageHandlingException; + +/** + * An exception indicating that a message was rejected by an implementation of + * {@link org.springframework.integration.message.selector.MessageSelector}. + * + * @author Mark Fisher + */ +public class MessageSelectorRejectedException extends MessageHandlingException { + + public MessageSelectorRejectedException(String message) { + super(message); + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageReceivingExecutor.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/PooledMessageHandler.java similarity index 77% rename from spring-integration-core/src/main/java/org/springframework/integration/bus/MessageReceivingExecutor.java rename to spring-integration-core/src/main/java/org/springframework/integration/bus/PooledMessageHandler.java index 8232bca37e..d15d22fe7a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageReceivingExecutor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/PooledMessageHandler.java @@ -18,6 +18,7 @@ package org.springframework.integration.bus; import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.SynchronousQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; @@ -27,8 +28,8 @@ import org.apache.commons.logging.LogFactory; import org.springframework.context.Lifecycle; import org.springframework.integration.MessageHandlingException; +import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.Message; -import org.springframework.integration.message.MessageReceiver; import org.springframework.integration.message.selector.MessageSelector; import org.springframework.scheduling.concurrent.CustomizableThreadFactory; import org.springframework.util.Assert; @@ -38,11 +39,11 @@ import org.springframework.util.Assert; * * @author Mark Fisher */ -public class MessageReceivingExecutor implements Lifecycle { +public class PooledMessageHandler implements MessageHandler, Lifecycle { private Log logger = LogFactory.getLog(this.getClass()); - private MessageReceiver receiver; + private MessageHandler handler; private List selectors = new CopyOnWriteArrayList(); @@ -65,12 +66,12 @@ public class MessageReceivingExecutor implements Lifecycle { private int totalErrorThreshold = -1; - public MessageReceivingExecutor(MessageReceiver receiver, int corePoolSize, int maxPoolSize) { - Assert.notNull(receiver, "'receiver' must not be null"); + public PooledMessageHandler(MessageHandler handler, int corePoolSize, int maxPoolSize) { + Assert.notNull(handler, "'handler' must not be null"); Assert.isTrue(corePoolSize > 0, "'corePoolSize' must be at least 1"); Assert.isTrue(maxPoolSize > 0, "'maxPoolSize' must be at least 1"); Assert.isTrue(maxPoolSize >= corePoolSize, "'corePoolSize' cannot exceed 'maxPoolSize'"); - this.receiver = receiver; + this.handler = handler; this.corePoolSize = corePoolSize; this.maxPoolSize = maxPoolSize; } @@ -95,7 +96,7 @@ public class MessageReceivingExecutor implements Lifecycle { public void start() { synchronized (this.lifecycleMonitor) { if (!this.running) { - this.threadPoolExecutor = new MessageReceivingThreadPoolExecutor(this.corePoolSize, this.maxPoolSize); + this.threadPoolExecutor = new MessageHandlerThreadPoolExecutor(this.corePoolSize, this.maxPoolSize); } this.running = true; } @@ -111,17 +112,23 @@ public class MessageReceivingExecutor implements Lifecycle { } } - public boolean acceptMessage(Message message) { - if (threadPoolExecutor == null) { - throw new MessageHandlingException("executor is not running"); + public Message handle(Message message) { + if (!this.isRunning()) { + throw new MessageHandlerNotRunningException(); } for (MessageSelector selector : this.selectors) { if (!selector.accept(message)) { - return false; + throw new MessageSelectorRejectedException("selector rejected message"); } } - this.threadPoolExecutor.execute(new MessageReceivingTask(this.receiver, message)); - return true; + try { + this.threadPoolExecutor.execute(new HandlerTask(this.handler, message)); + return null; + } + catch (RejectedExecutionException e) { + throw new MessageHandlerRejectedExecutionException( + "handler executor rejected message", e); + } } /** @@ -153,17 +160,17 @@ public class MessageReceivingExecutor implements Lifecycle { } - private static class MessageReceivingTask implements Runnable { + private static class HandlerTask implements Runnable { - private MessageReceiver receiver; + private MessageHandler handler; private Message message; private Throwable error; - MessageReceivingTask(MessageReceiver receiver, Message message) { - this.receiver = receiver; + HandlerTask(MessageHandler handler, Message message) { + this.handler = handler; this.message = message; } @@ -173,7 +180,7 @@ public class MessageReceivingExecutor implements Lifecycle { public void run() { try { - this.receiver.messageReceived(this.message); + this.handler.handle(this.message); } catch (Throwable t) { this.error = t; @@ -182,18 +189,18 @@ public class MessageReceivingExecutor implements Lifecycle { } - private class MessageReceivingThreadPoolExecutor extends ThreadPoolExecutor { + private class MessageHandlerThreadPoolExecutor extends ThreadPoolExecutor { - public MessageReceivingThreadPoolExecutor(int corePoolSize, int maximumPoolSize) { + public MessageHandlerThreadPoolExecutor(int corePoolSize, int maximumPoolSize) { super(corePoolSize, maximumPoolSize, 0, TimeUnit.MILLISECONDS, new SynchronousQueue()); CustomizableThreadFactory threadFactory = new CustomizableThreadFactory(); - threadFactory.setThreadNamePrefix("endpoint-executor-"); + threadFactory.setThreadNamePrefix("handler-pool-"); this.setThreadFactory(threadFactory); } @Override protected void afterExecute(Runnable r, Throwable t) { - MessageReceivingTask task = (MessageReceivingTask) r; + HandlerTask task = (HandlerTask) r; if (task.getError() != null) { if (logger.isWarnEnabled()) { logger.warn("Exception occurred during task execution", task.getError()); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/GenericMessageEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/GenericMessageEndpoint.java index 62d815d8ff..aec445f8bb 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/GenericMessageEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/GenericMessageEndpoint.java @@ -38,7 +38,7 @@ import org.springframework.integration.message.Message; * * @author Mark Fisher */ -public class GenericMessageEndpoint implements MessageEndpoint, BeanNameAware { +public class GenericMessageEndpoint implements MessageEndpoint, BeanNameAware { private String name; @@ -113,7 +113,7 @@ public class GenericMessageEndpoint implements MessageEndpoint, BeanNameAw } - public void messageReceived(Message message) { + public Message handle(Message message) { if (this.handler == null) { if (this.defaultOutputChannelName == null) { throw new MessagingConfigurationException( @@ -121,7 +121,7 @@ public class GenericMessageEndpoint implements MessageEndpoint, BeanNameAw } MessageChannel replyChannel = this.channelRegistry.lookupChannel(this.defaultOutputChannelName); replyChannel.send(message); - return; + return null; } Message replyMessage = handler.handle(message); if (replyMessage != null) { @@ -133,6 +133,7 @@ public class GenericMessageEndpoint implements MessageEndpoint, BeanNameAw } replyChannel.send(replyMessage); } + return null; } private MessageChannel resolveReplyChannel(Message message) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java index a9d0dc161d..09872edd5e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java @@ -18,14 +18,16 @@ package org.springframework.integration.endpoint; import org.springframework.integration.bus.ConsumerPolicy; import org.springframework.integration.channel.ChannelRegistry; -import org.springframework.integration.message.MessageReceiver; +import org.springframework.integration.handler.MessageHandler; /** * Base interface for message endpoints. * * @author Mark Fisher */ -public interface MessageEndpoint extends MessageReceiver { +public interface MessageEndpoint extends MessageHandler { + + void setName(String name); void setInputChannelName(String inputChannelName); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/DefaultMessageDispatcherTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/DefaultMessageDispatcherTests.java index ae43290491..515956cd53 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/DefaultMessageDispatcherTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/DefaultMessageDispatcherTests.java @@ -49,8 +49,9 @@ public class DefaultMessageDispatcherTests { channel.send(new StringMessage(1, "test")); MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1)); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1)); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("exactly one endpoint should have received message", 1, counter1.get() + counter2.get()); @@ -69,8 +70,9 @@ public class DefaultMessageDispatcherTests { MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); dispatcher.setBroadcast(true); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1)); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1)); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("both endpoints should have received message", 2, counter1.get() + counter2.get()); @@ -90,13 +92,14 @@ public class DefaultMessageDispatcherTests { channel.send(new StringMessage(1, "test")); MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1) { @Override public void start() { } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1)); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint3, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint3, 1, 1)); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("inactive endpoint should not have received message", 0, counter1.get()); @@ -118,13 +121,14 @@ public class DefaultMessageDispatcherTests { MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); dispatcher.setBroadcast(true); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1)); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) { @Override public void start() { } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint3, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint3, 1, 1)); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("inactive endpoint should not have received message", 0, counter2.get()); @@ -158,14 +162,15 @@ public class DefaultMessageDispatcherTests { dispatcher.setBroadcast(true); dispatcher.setRejectionLimit(2); dispatcher.setRetryInterval(3); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1)); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) { @Override - public boolean acceptMessage(Message message) { - throw new RejectedExecutionException(); + public Message handle(Message message) { + throw new MessageHandlerRejectedExecutionException(null); } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint3, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint3, 1, 1)); + dispatcher.start(); dispatcher.dispatch(); } @@ -187,14 +192,15 @@ public class DefaultMessageDispatcherTests { dispatcher.setRejectionLimit(2); dispatcher.setRetryInterval(3); dispatcher.setShouldFailOnRejectionLimit(false); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1)); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) { @Override - public boolean acceptMessage(Message message) { - throw new RejectedExecutionException(); + public Message handle(Message message) { + throw new MessageHandlerRejectedExecutionException(null); } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint3, 1, 1)); + dispatcher.addHandler(new PooledMessageHandler(endpoint3, 1, 1)); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("rejecting endpoint should not have received message", 0, counter2.get()); @@ -215,18 +221,19 @@ public class DefaultMessageDispatcherTests { DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); dispatcher.setRejectionLimit(2); dispatcher.setRetryInterval(3); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1) { @Override - public boolean acceptMessage(Message message) { - throw new RejectedExecutionException(); + public Message handle(Message message) { + throw new MessageHandlerRejectedExecutionException(null); } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) { @Override - public boolean acceptMessage(Message message) { - throw new RejectedExecutionException(); + public Message handle(Message message) { + throw new MessageHandlerRejectedExecutionException(null); } }); + dispatcher.start(); dispatcher.dispatch(); } @@ -247,20 +254,21 @@ public class DefaultMessageDispatcherTests { dispatcher.setRejectionLimit(2); dispatcher.setRetryInterval(3); dispatcher.setShouldFailOnRejectionLimit(false); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { rejectedCounter1.incrementAndGet(); - throw new RejectedExecutionException(); + throw new MessageHandlerRejectedExecutionException(null); } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { rejectedCounter2.incrementAndGet(); - throw new RejectedExecutionException(); + throw new MessageHandlerRejectedExecutionException(null); } }); + dispatcher.start(); dispatcher.dispatch(); assertEquals("rejecting endpoints should not have received message", 0, counter1.get() + counter2.get()); assertEquals("endpoint1 should have rejected two times", 2, rejectedCounter1.get()); @@ -287,30 +295,31 @@ public class DefaultMessageDispatcherTests { dispatcher.setRejectionLimit(2); dispatcher.setRetryInterval(3); dispatcher.setShouldFailOnRejectionLimit(false); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { rejectedCounter1.incrementAndGet(); - throw new RejectedExecutionException(); + throw new MessageHandlerRejectedExecutionException(null); } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { if (rejectedCounter2.get() == 1) { - return super.acceptMessage(message); + return super.handle(message); } rejectedCounter2.incrementAndGet(); - throw new RejectedExecutionException(); + throw new MessageHandlerRejectedExecutionException(null); } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint3, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint3, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { rejectedCounter3.incrementAndGet(); - throw new RejectedExecutionException(); + throw new MessageHandlerRejectedExecutionException(null); } }); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("endpoint1 should not have received message", 0, counter1.get()); @@ -339,26 +348,27 @@ public class DefaultMessageDispatcherTests { dispatcher.setRejectionLimit(5); dispatcher.setRetryInterval(3); dispatcher.setShouldFailOnRejectionLimit(false); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint1, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint1, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { if (rejectedCounter1.get() == 2) { - return super.acceptMessage(message); + return super.handle(message); } rejectedCounter1.incrementAndGet(); - throw new RejectedExecutionException(); + throw new MessageHandlerRejectedExecutionException(null); } }); - dispatcher.addExecutor(new MessageReceivingExecutor(endpoint2, 1, 1) { + dispatcher.addHandler(new PooledMessageHandler(endpoint2, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { if (rejectedCounter2.get() == 4) { - return super.acceptMessage(message); + return super.handle(message); } rejectedCounter2.incrementAndGet(); - throw new RejectedExecutionException(); + throw new MessageHandlerRejectedExecutionException(null); } }); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("endpoint1 should have received one message", 1, counter1.get()); @@ -379,12 +389,13 @@ public class DefaultMessageDispatcherTests { channel.send(new StringMessage(1, "test")); MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); - MessageReceivingExecutor executor1 = new MessageReceivingExecutor(endpoint1, 1, 1); - MessageReceivingExecutor executor2 = new MessageReceivingExecutor(endpoint2, 1, 1); + PooledMessageHandler executor1 = new PooledMessageHandler(endpoint1, 1, 1); + PooledMessageHandler executor2 = new PooledMessageHandler(endpoint2, 1, 1); executor1.addMessageSelector(new PayloadTypeSelector(Integer.class)); executor2.addMessageSelector(new PayloadTypeSelector(String.class)); - dispatcher.addExecutor(executor1); - dispatcher.addExecutor(executor2); + dispatcher.addHandler(executor1); + dispatcher.addHandler(executor2); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("endpoint1 should not have accepted the message", 0, counter1.get()); @@ -406,26 +417,27 @@ public class DefaultMessageDispatcherTests { channel.send(new StringMessage(1, "test")); MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); - MessageReceivingExecutor executor1 = new MessageReceivingExecutor(endpoint1, 1, 1) { + PooledMessageHandler executor1 = new PooledMessageHandler(endpoint1, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { attemptedCounter1.incrementAndGet(); attemptedLatch.countDown(); - return super.acceptMessage(message); + return super.handle(message); } }; - MessageReceivingExecutor executor2 = new MessageReceivingExecutor(endpoint2, 1, 1) { + PooledMessageHandler executor2 = new PooledMessageHandler(endpoint2, 1, 1) { @Override - public boolean acceptMessage(Message message) { + public Message handle(Message message) { attemptedCounter2.incrementAndGet(); attemptedLatch.countDown(); - return super.acceptMessage(message); + return super.handle(message); } }; executor1.addMessageSelector(new PayloadTypeSelector(Integer.class)); executor2.addMessageSelector(new PayloadTypeSelector(Integer.class)); - dispatcher.addExecutor(executor1); - dispatcher.addExecutor(executor2); + dispatcher.addHandler(executor1); + dispatcher.addHandler(executor2); + dispatcher.start(); dispatcher.dispatch(); attemptedLatch.await(100, TimeUnit.MILLISECONDS); assertEquals("endpoint1 should not have accepted the message", 0, counter1.get()); @@ -449,12 +461,13 @@ public class DefaultMessageDispatcherTests { MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(retriever); dispatcher.setBroadcast(true); - MessageReceivingExecutor executor1 = new MessageReceivingExecutor(endpoint1, 1, 1); - MessageReceivingExecutor executor2 = new MessageReceivingExecutor(endpoint2, 1, 1); + PooledMessageHandler executor1 = new PooledMessageHandler(endpoint1, 1, 1); + PooledMessageHandler executor2 = new PooledMessageHandler(endpoint2, 1, 1); executor1.addMessageSelector(new PayloadTypeSelector(Integer.class)); executor2.addMessageSelector(new PayloadTypeSelector(String.class)); - dispatcher.addExecutor(executor1); - dispatcher.addExecutor(executor2); + dispatcher.addHandler(executor1); + dispatcher.addHandler(executor2); + dispatcher.start(); dispatcher.dispatch(); latch.await(100, TimeUnit.MILLISECONDS); assertEquals("endpoint1 should not have accepted the message", 0, counter1.get()); @@ -462,7 +475,7 @@ public class DefaultMessageDispatcherTests { } - private static class TestEndpoint extends GenericMessageEndpoint { + private static class TestEndpoint extends GenericMessageEndpoint { private AtomicInteger counter; @@ -475,9 +488,10 @@ public class DefaultMessageDispatcherTests { } @Override - public void messageReceived(Message message) { + public Message handle(Message message) { counter.incrementAndGet(); latch.countDown(); + return null; } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedDelayConsumerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedDelayConsumerTests.java index b520f3c839..4f3a1152ff 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedDelayConsumerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedDelayConsumerTests.java @@ -42,9 +42,11 @@ public class FixedDelayConsumerTests { final CountDownLatch latch = new CountDownLatch(messagesToSend); SimpleChannel channel = new SimpleChannel(); MessageEndpoint endpoint = new GenericMessageEndpoint() { - public void messageReceived(Message message) { + @Override + public Message handle(Message message) { counter.incrementAndGet(); latch.countDown(); + return null; } }; MessageBus bus = new MessageBus(); @@ -76,9 +78,11 @@ public class FixedDelayConsumerTests { final CountDownLatch latch = new CountDownLatch(messagesToSend); SimpleChannel channel = new SimpleChannel(); MessageEndpoint endpoint = new GenericMessageEndpoint() { - public void messageReceived(Message message) { + @Override + public Message handle(Message message) { counter.incrementAndGet(); latch.countDown(); + return null; } }; MessageBus bus = new MessageBus(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedRateConsumerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedRateConsumerTests.java index 24f88a43f1..8d2f65fb07 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedRateConsumerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/FixedRateConsumerTests.java @@ -42,10 +42,12 @@ public class FixedRateConsumerTests { final AtomicInteger counter = new AtomicInteger(0); final CountDownLatch latch = new CountDownLatch(messagesToSend); SimpleChannel channel = new SimpleChannel(); - MessageEndpoint endpoint = new GenericMessageEndpoint() { - public void messageReceived(Message message) { + MessageEndpoint endpoint = new GenericMessageEndpoint() { + @Override + public Message handle(Message message) { counter.incrementAndGet(); latch.countDown(); + return null; } }; MessageBus bus = new MessageBus(); @@ -73,10 +75,12 @@ public class FixedRateConsumerTests { final AtomicInteger counter = new AtomicInteger(0); final CountDownLatch latch = new CountDownLatch(messagesToSend); SimpleChannel channel = new SimpleChannel(); - MessageEndpoint endpoint = new GenericMessageEndpoint() { - public void messageReceived(Message message) { + MessageEndpoint endpoint = new GenericMessageEndpoint() { + @Override + public Message handle(Message message) { counter.incrementAndGet(); latch.countDown(); + return null; } }; MessageBus bus = new MessageBus(); 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 284e579948..3c1e5086d7 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 @@ -46,7 +46,7 @@ public class MessageBusTests { bus.registerChannel("sourceChannel", sourceChannel); sourceChannel.send(new StringMessage("123", "test")); bus.registerChannel("targetChannel", targetChannel); - GenericMessageEndpoint endpoint = new GenericMessageEndpoint(); + GenericMessageEndpoint endpoint = new GenericMessageEndpoint(); endpoint.setInputChannelName("sourceChannel"); endpoint.setDefaultOutputChannelName("targetChannel"); bus.registerEndpoint("endpoint", endpoint); @@ -93,10 +93,10 @@ public class MessageBusTests { SimpleChannel inputChannel = new SimpleChannel(); SimpleChannel outputChannel1 = new SimpleChannel(); SimpleChannel outputChannel2 = new SimpleChannel(); - GenericMessageEndpoint endpoint1 = new GenericMessageEndpoint(); + GenericMessageEndpoint endpoint1 = new GenericMessageEndpoint(); endpoint1.setDefaultOutputChannelName("output1"); endpoint1.setInputChannelName("input"); - GenericMessageEndpoint endpoint2 = new GenericMessageEndpoint(); + GenericMessageEndpoint endpoint2 = new GenericMessageEndpoint(); endpoint2.setDefaultOutputChannelName("output2"); endpoint2.setInputChannelName("input"); MessageBus bus = new MessageBus();