From d731f8fb70ad9ba681d1e7c42aebdcd3bcf3744e Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Sun, 6 Jul 2008 22:11:44 +0000 Subject: [PATCH] MessageBus is now an interface. The DefaultMessageBus class is the implementation. --- .../integration/bus/DefaultMessageBus.java | 544 ++++++++++++++++++ .../integration/bus/MessageBus.java | 514 +---------------- 2 files changed, 551 insertions(+), 507 deletions(-) create mode 100644 org.springframework.integration/src/main/java/org/springframework/integration/bus/DefaultMessageBus.java diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/bus/DefaultMessageBus.java b/org.springframework.integration/src/main/java/org/springframework/integration/bus/DefaultMessageBus.java new file mode 100644 index 0000000000..e2252d7da6 --- /dev/null +++ b/org.springframework.integration/src/main/java/org/springframework/integration/bus/DefaultMessageBus.java @@ -0,0 +1,544 @@ +/* + * Copyright 2002-2008 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.bus; + +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CopyOnWriteArraySet; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.DisposableBean; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.ApplicationEvent; +import org.springframework.context.ApplicationListener; +import org.springframework.context.Lifecycle; +import org.springframework.context.event.ApplicationEventMulticaster; +import org.springframework.context.event.ContextRefreshedEvent; +import org.springframework.context.event.SimpleApplicationEventMulticaster; +import org.springframework.context.support.AbstractApplicationContext; +import org.springframework.integration.ConfigurationException; +import org.springframework.integration.bus.interceptor.MessageBusInterceptor; +import org.springframework.integration.channel.ChannelRegistry; +import org.springframework.integration.channel.ChannelRegistryAware; +import org.springframework.integration.channel.DefaultChannelRegistry; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.channel.factory.ChannelFactory; +import org.springframework.integration.channel.factory.QueueChannelFactory; +import org.springframework.integration.endpoint.AbstractEndpoint; +import org.springframework.integration.endpoint.DefaultEndpointRegistry; +import org.springframework.integration.endpoint.EndpointRegistry; +import org.springframework.integration.endpoint.EndpointTrigger; +import org.springframework.integration.endpoint.HandlerEndpoint; +import org.springframework.integration.endpoint.MessageEndpoint; +import org.springframework.integration.endpoint.MessagingGateway; +import org.springframework.integration.endpoint.TargetEndpoint; +import org.springframework.integration.handler.MessageHandler; +import org.springframework.integration.message.MessageTarget; +import org.springframework.integration.message.Subscribable; +import org.springframework.integration.scheduling.MessagePublishingErrorHandler; +import org.springframework.integration.scheduling.PollingSchedule; +import org.springframework.integration.scheduling.Schedule; +import org.springframework.integration.scheduling.SimpleTaskScheduler; +import org.springframework.integration.scheduling.TaskScheduler; +import org.springframework.scheduling.concurrent.CustomizableThreadFactory; +import org.springframework.util.Assert; + +/** + * The messaging bus. Serves as a registry for channels and endpoints, manages their lifecycle, + * and activates subscriptions. + * + * @author Mark Fisher + * @author Marius Bogoevici + */ +public class DefaultMessageBus implements MessageBus, ApplicationContextAware, ApplicationListener { + + private static final int DEFAULT_DISPATCHER_POOL_SIZE = 10; + + private final Log logger = LogFactory.getLog(this.getClass()); + + private volatile ChannelFactory channelFactory = new QueueChannelFactory(); + + private final ChannelRegistry channelRegistry = new DefaultChannelRegistry(); + + private final EndpointRegistry endpointRegistry = new DefaultEndpointRegistry(); + + private final Set endpointTriggers = new CopyOnWriteArraySet(); + + private final List lifecycleEndpoints = new CopyOnWriteArrayList(); + + private final MessageBusInterceptorsList interceptors = new MessageBusInterceptorsList(); + + private volatile Schedule defaultPollerSchedule = new PollingSchedule(0); + + private volatile TaskScheduler taskScheduler; + + private volatile boolean configureAsyncEventMulticaster = false; + + private volatile boolean autoCreateChannels = false; + + private volatile boolean autoStartup = true; + + private volatile boolean initialized; + + private volatile boolean initializing; + + private volatile boolean starting; + + private volatile boolean running; + + private final Object lifecycleMonitor = new Object(); + + /** + * Set the {@link ChannelFactory} to use for auto-creating channels. + */ + public void setChannelFactory(ChannelFactory channelFactory) { + this.channelFactory = channelFactory; + } + + public ChannelFactory getChannelFactory() { + return channelFactory; + } + + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + Assert.notNull(applicationContext, "'applicationContext' must not be null"); + if (applicationContext.getBeanNamesForType(this.getClass()).length > 1) { + throw new ConfigurationException("Only one instance of '" + this.getClass().getSimpleName() + + "' is allowed per ApplicationContext."); + } + this.registerChannels(applicationContext); + } + + /** + * Set the {@link TaskScheduler} to use for scheduling message dispatchers. + */ + public void setTaskScheduler(TaskScheduler taskScheduler) { + this.taskScheduler = taskScheduler; + } + + /** + * Set whether to automatically start the bus after initialization. + *

Default is 'true'; set this to 'false' to allow for manual startup + * through the {@link #start()} method. + */ + public void setAutoStartup(boolean autoStartup) { + this.autoStartup = autoStartup; + } + + /** + * Set whether the bus should automatically create a channel when a + * subscription contains the name of a previously unregistered channel. + */ + public void setAutoCreateChannels(boolean autoCreateChannels) { + this.autoCreateChannels = autoCreateChannels; + } + + /** + * Set whether the bus should configure its asynchronous task executor + * to also be used by the ApplicationContext's 'applicationEventMulticaster'. + * This will only apply if the multicaster defined within the context + * is an instance of SimpleApplicationEventMulticaster (the default). + * This property is 'false' by default. + */ + public void setConfigureAsyncEventMulticaster(boolean configureAsyncEventMulticaster) { + this.configureAsyncEventMulticaster = configureAsyncEventMulticaster; + } + + @SuppressWarnings("unchecked") + private void registerChannels(ApplicationContext context) { + Map channelBeans = (Map) context + .getBeansOfType(MessageChannel.class); + for (Map.Entry entry : channelBeans.entrySet()) { + String channelName = entry.getKey(); + MessageChannel previousChannel = this.lookupChannel(channelName); + if (previousChannel == null) { + this.registerChannel(channelName, entry.getValue()); + } + else if (!previousChannel.equals(entry.getValue())) { + throw new ConfigurationException("A different channel instance has already " + + "been registered with the name '" + channelName + "'."); + } + } + } + + @SuppressWarnings("unchecked") + private void registerEndpoints(ApplicationContext context) { + Map endpointBeans = (Map) context + .getBeansOfType(MessageEndpoint.class); + for (Map.Entry entry : endpointBeans.entrySet()) { + this.registerEndpoint(entry.getValue()); + } + } + + @SuppressWarnings("unchecked") + private void registerGateways(ApplicationContext context) { + Map gatewayBeans = (Map) context + .getBeansOfType(MessagingGateway.class); + for (Map.Entry entry : gatewayBeans.entrySet()) { + this.registerGateway(entry.getKey(), entry.getValue()); + } + } + + public void initialize() { + synchronized (this.lifecycleMonitor) { + if (this.initialized || this.initializing) { + return; + } + this.initializing = true; + if (this.taskScheduler == null) { + ScheduledThreadPoolExecutor executor = new ScheduledThreadPoolExecutor(DEFAULT_DISPATCHER_POOL_SIZE); + executor.setThreadFactory(new CustomizableThreadFactory("message-bus-")); + executor.setRejectedExecutionHandler(new CallerRunsPolicy()); + this.taskScheduler = new SimpleTaskScheduler(executor); + } + if (this.getErrorChannel() == null) { + this.setErrorChannel(new DefaultErrorChannel()); + } + this.initialized = true; + this.initializing = false; + } + } + + public MessageChannel getErrorChannel() { + return this.lookupChannel(ERROR_CHANNEL_NAME); + } + + public void setErrorChannel(MessageChannel errorChannel) { + this.registerChannel(ERROR_CHANNEL_NAME, errorChannel); + } + + public MessageChannel lookupChannel(String channelName) { + return this.channelRegistry.lookupChannel(channelName); + } + + public void registerChannel(String name, MessageChannel channel) { + if (!this.initialized) { + this.initialize(); + } + channel.setName(name); + this.channelRegistry.registerChannel(name, channel); + if (logger.isInfoEnabled()) { + logger.info("registered channel '" + name + "'"); + } + } + + public MessageChannel unregisterChannel(String name) { + return this.channelRegistry.unregisterChannel(name); + } + + public void registerHandler(String name, MessageHandler handler, Object input, Schedule schedule) { + Assert.notNull(handler, "'handler' must not be null"); + HandlerEndpoint endpoint = new HandlerEndpoint(handler); + this.configureEndpoint(endpoint, name, input, schedule); + this.registerEndpoint(endpoint); + } + + public void registerTarget(String name, MessageTarget target, Object input, Schedule schedule) { + Assert.notNull(target, "'target' must not be null"); + TargetEndpoint endpoint = new TargetEndpoint(target); + this.configureEndpoint(endpoint, name, input, schedule); + this.registerEndpoint(endpoint); + } + + private void configureEndpoint(AbstractEndpoint endpoint, String name, Object input, Schedule schedule) { + endpoint.setName(name); + if (input instanceof MessageChannel) { + endpoint.setInputChannel((MessageChannel) input); + } + else if (input instanceof String) { + endpoint.setInputChannelName((String) input); + } + else { + throw new ConfigurationException("'input' must be a MessageChannel or String"); + } + endpoint.setSchedule(schedule); + } + + public void registerEndpoint(MessageEndpoint endpoint) { + if (!this.initialized) { + this.initialize(); + } + if (endpoint instanceof ChannelRegistryAware) { + ((ChannelRegistryAware) endpoint).setChannelRegistry(this.channelRegistry); + } + this.endpointRegistry.registerEndpoint(endpoint); + if (this.isRunning()) { + this.activateEndpoint(endpoint); + } + if (logger.isInfoEnabled()) { + logger.info("registered endpoint '" + endpoint + "'"); + } + } + + public MessageEndpoint unregisterEndpoint(String name) { + MessageEndpoint endpoint = this.endpointRegistry.unregisterEndpoint(name); + if (endpoint == null) { + return null; + } + this.deactivateEndpoint(endpoint); + return endpoint; + } + + public MessageEndpoint lookupEndpoint(String endpointName) { + return this.endpointRegistry.lookupEndpoint(endpointName); + } + + public Set getEndpointNames() { + return this.endpointRegistry.getEndpointNames(); + } + + private void activateEndpoints() { + Set endpointNames = this.endpointRegistry.getEndpointNames(); + for (String name : endpointNames) { + MessageEndpoint endpoint = this.endpointRegistry.lookupEndpoint(name); + if (endpoint != null) { + this.activateEndpoint(endpoint); + } + } + } + + private void activateEndpoint(MessageEndpoint endpoint) { + Assert.notNull(endpoint, "'endpoint' must not be null"); + if (endpoint.getOutputChannel() == null) { + this.lookupOrCreateChannel(endpoint.getOutputChannelName()); + } + try { + endpoint.afterPropertiesSet(); + } + catch (Exception e) { + throw new ConfigurationException("failed to initialize endpoint", e); + } + MessageChannel channel = endpoint.getInputChannel(); + if (channel == null) { + channel = this.lookupOrCreateChannel(endpoint.getInputChannelName()); + } + if (channel != null && channel instanceof Subscribable) { + ((Subscribable) channel).subscribe(endpoint); + if (logger.isInfoEnabled()) { + logger.info("activated subscription to channel '" + + channel.getName() + "' for endpoint '" + endpoint + "'"); + } + return; + } + Schedule schedule = endpoint.getSchedule(); + EndpointTrigger trigger = endpoint.getTrigger(); + if (trigger == null) { + trigger = new EndpointTrigger(schedule != null ? schedule : this.defaultPollerSchedule); + } + trigger.addTarget(endpoint); + if (this.endpointTriggers.add(trigger)) { + this.taskScheduler.schedule(trigger); + } + } + + private MessageChannel lookupOrCreateChannel(String channelName) { + if (channelName == null) { + return null; + } + MessageChannel channel = this.lookupChannel(channelName); + if (channel == null) { + if (!this.autoCreateChannels) { + throw new ConfigurationException("Cannot activate endpoint, unknown channel '" + channelName + + "'. Consider enabling the 'autoCreateChannels' option for the message bus."); + } + if (this.logger.isInfoEnabled()) { + logger.info("auto-creating channel '" + channelName + "'"); + } + channel = channelFactory.getChannel(channelName, null, null); + this.registerChannel(channelName, channel); + } + return channel; + } + + private void registerGateway(String name, MessagingGateway gateway) { + if (gateway instanceof Lifecycle) { + this.lifecycleEndpoints.add((Lifecycle) gateway); + if (this.isRunning()) { + ((Lifecycle) gateway).start(); + } + } + if (logger.isInfoEnabled()) { + logger.info("registered gateway '" + name + "'"); + } + } + + public void deactivateEndpoint(MessageEndpoint endpoint) { + Assert.notNull(endpoint, "'endpoint' must not be null"); + for (EndpointTrigger trigger : this.endpointTriggers) { + boolean removed = trigger.removeTarget(endpoint); + if (removed && this.logger.isInfoEnabled()) { + logger.info("removed endpoint '" + endpoint + "' from dispatcher"); + } + } + if (endpoint instanceof Lifecycle) { + ((Lifecycle) endpoint).stop(); + } + } + + public boolean isRunning() { + synchronized (this.lifecycleMonitor) { + return this.running; + } + } + + public void start() { + if (!this.initialized) { + this.initialize(); + } + if (this.isRunning() || this.starting) { + return; + } + this.interceptors.preStart(); + this.starting = true; + synchronized (this.lifecycleMonitor) { + this.activateEndpoints(); + this.taskScheduler.setErrorHandler(new MessagePublishingErrorHandler(this.getErrorChannel())); + this.taskScheduler.start(); + for (Lifecycle endpoint : this.lifecycleEndpoints) { + endpoint.start(); + if (logger.isInfoEnabled()) { + logger.info("started endpoint '" + endpoint + "'"); + } + } + } + this.running = true; + this.starting = false; + this.interceptors.postStart(); + if (logger.isInfoEnabled()) { + logger.info("message bus started"); + } + } + + public void stop() { + if (!this.isRunning()) { + return; + } + this.interceptors.preStop(); + synchronized (this.lifecycleMonitor) { + this.running = false; + this.taskScheduler.stop(); + for (Lifecycle endpoint : this.lifecycleEndpoints) { + endpoint.stop(); + if (logger.isInfoEnabled()) { + logger.info("stopped endpoint '" + endpoint + "'"); + } + } + } + this.interceptors.postStop(); + if (logger.isInfoEnabled()) { + logger.info("message bus stopped"); + } + } + + public void destroy() throws Exception { + if (this.taskScheduler instanceof DisposableBean) { + ((DisposableBean) this.taskScheduler).destroy(); + } + } + + public void onApplicationEvent(ApplicationEvent event) { + if (event instanceof ContextRefreshedEvent) { + ApplicationContext context = ((ContextRefreshedEvent) event).getApplicationContext(); + this.registerChannels(context); + this.registerEndpoints(context); + this.registerGateways(context); + if (this.configureAsyncEventMulticaster) { + this.initialize(); + this.doConfigureAsyncEventMulticaster(context); + } + if (this.autoStartup) { + this.start(); + } + } + } + + private void doConfigureAsyncEventMulticaster(ApplicationContext context) { + String multicasterBeanName = AbstractApplicationContext.APPLICATION_EVENT_MULTICASTER_BEAN_NAME; + if (context.containsBean(multicasterBeanName)) { + ApplicationEventMulticaster multicaster = (ApplicationEventMulticaster) context + .getBean(multicasterBeanName); + if (multicaster instanceof SimpleApplicationEventMulticaster) { + ((SimpleApplicationEventMulticaster) multicaster).setTaskExecutor(this.taskScheduler); + } + } + } + + public void addInterceptor(MessageBusInterceptor interceptor) { + this.interceptors.add(interceptor); + } + + public void removeInterceptor(MessageBusInterceptor interceptor) { + this.interceptors.remove(interceptor); + } + + public void setInterceptors(List interceptor) { + this.interceptors.set(interceptor); + } + + /* + * Wrapper class for the interceptor list + */ + private class MessageBusInterceptorsList { + + private CopyOnWriteArrayList messageBusInterceptors = new CopyOnWriteArrayList(); + + public void set(List interceptors) { + this.messageBusInterceptors.clear(); + this.messageBusInterceptors.addAll(interceptors); + } + + public void add(MessageBusInterceptor interceptor) { + this.messageBusInterceptors.add(interceptor); + } + + public void remove(MessageBusInterceptor interceptor) { + this.messageBusInterceptors.remove(interceptor); + } + + public void preStart() { + for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { + messageBusInterceptor.preStart(DefaultMessageBus.this); + } + } + + public void postStart() { + for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { + messageBusInterceptor.postStart(DefaultMessageBus.this); + } + } + + public void preStop() { + for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { + messageBusInterceptor.preStop(DefaultMessageBus.this); + } + } + + public void postStop() { + for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { + messageBusInterceptor.postStop(DefaultMessageBus.this); + } + } + } + +} diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/bus/MessageBus.java b/org.springframework.integration/src/main/java/org/springframework/integration/bus/MessageBus.java index b98141ad95..483c1501f4 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/bus/MessageBus.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/bus/MessageBus.java @@ -16,532 +16,32 @@ package org.springframework.integration.bus; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.CopyOnWriteArrayList; -import java.util.concurrent.CopyOnWriteArraySet; -import java.util.concurrent.ScheduledThreadPoolExecutor; -import java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy; - -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; - -import org.springframework.beans.BeansException; import org.springframework.beans.factory.DisposableBean; -import org.springframework.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; -import org.springframework.context.ApplicationEvent; -import org.springframework.context.ApplicationListener; import org.springframework.context.Lifecycle; -import org.springframework.context.event.ApplicationEventMulticaster; -import org.springframework.context.event.ContextRefreshedEvent; -import org.springframework.context.event.SimpleApplicationEventMulticaster; -import org.springframework.context.support.AbstractApplicationContext; -import org.springframework.integration.ConfigurationException; -import org.springframework.integration.bus.interceptor.MessageBusInterceptor; import org.springframework.integration.channel.ChannelRegistry; -import org.springframework.integration.channel.ChannelRegistryAware; -import org.springframework.integration.channel.DefaultChannelRegistry; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.factory.ChannelFactory; -import org.springframework.integration.channel.factory.QueueChannelFactory; -import org.springframework.integration.endpoint.AbstractEndpoint; -import org.springframework.integration.endpoint.DefaultEndpointRegistry; import org.springframework.integration.endpoint.EndpointRegistry; -import org.springframework.integration.endpoint.EndpointTrigger; -import org.springframework.integration.endpoint.HandlerEndpoint; -import org.springframework.integration.endpoint.MessageEndpoint; -import org.springframework.integration.endpoint.MessagingGateway; -import org.springframework.integration.endpoint.TargetEndpoint; import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.MessageTarget; -import org.springframework.integration.message.Subscribable; -import org.springframework.integration.scheduling.MessagePublishingErrorHandler; -import org.springframework.integration.scheduling.PollingSchedule; import org.springframework.integration.scheduling.Schedule; -import org.springframework.integration.scheduling.SimpleTaskScheduler; -import org.springframework.integration.scheduling.TaskScheduler; -import org.springframework.scheduling.concurrent.CustomizableThreadFactory; -import org.springframework.util.Assert; /** - * The messaging bus. Serves as a registry for channels and endpoints, manages their lifecycle, - * and activates subscriptions. + * The message bus interface. * * @author Mark Fisher - * @author Marius Bogoevici */ -public class MessageBus implements ChannelRegistry, EndpointRegistry, - ApplicationContextAware, ApplicationListener, Lifecycle, DisposableBean { +public interface MessageBus extends ChannelRegistry, EndpointRegistry, Lifecycle, DisposableBean { - public static final String ERROR_CHANNEL_NAME = "errorChannel"; + static final String ERROR_CHANNEL_NAME = "errorChannel"; - private static final int DEFAULT_DISPATCHER_POOL_SIZE = 10; - private final Log logger = LogFactory.getLog(this.getClass()); + MessageChannel getErrorChannel(); - private volatile ChannelFactory channelFactory = new QueueChannelFactory(); + ChannelFactory getChannelFactory(); - private final ChannelRegistry channelRegistry = new DefaultChannelRegistry(); + void registerHandler(String name, MessageHandler handler, Object input, Schedule schedule); - private final EndpointRegistry endpointRegistry = new DefaultEndpointRegistry(); - - private final Set endpointTriggers = new CopyOnWriteArraySet(); - - private final List lifecycleEndpoints = new CopyOnWriteArrayList(); - - private final MessageBusInterceptorsList interceptors = new MessageBusInterceptorsList(); - - private volatile Schedule defaultPollerSchedule = new PollingSchedule(0); - - private volatile TaskScheduler taskScheduler; - - private volatile boolean configureAsyncEventMulticaster = false; - - private volatile boolean autoCreateChannels = false; - - private volatile boolean autoStartup = true; - - private volatile boolean initialized; - - private volatile boolean initializing; - - private volatile boolean starting; - - private volatile boolean running; - - private final Object lifecycleMonitor = new Object(); - - /** - * Set the {@link ChannelFactory} to use for auto-creating channels. - */ - public void setChannelFactory(ChannelFactory channelFactory) { - this.channelFactory = channelFactory; - } - - public ChannelFactory getChannelFactory() { - return channelFactory; - } - - public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { - Assert.notNull(applicationContext, "'applicationContext' must not be null"); - if (applicationContext.getBeanNamesForType(this.getClass()).length > 1) { - throw new ConfigurationException("Only one instance of '" + this.getClass().getSimpleName() - + "' is allowed per ApplicationContext."); - } - this.registerChannels(applicationContext); - } - - /** - * Set the {@link TaskScheduler} to use for scheduling message dispatchers. - */ - public void setTaskScheduler(TaskScheduler taskScheduler) { - this.taskScheduler = taskScheduler; - } - - /** - * Set whether to automatically start the bus after initialization. - *

Default is 'true'; set this to 'false' to allow for manual startup - * through the {@link #start()} method. - */ - public void setAutoStartup(boolean autoStartup) { - this.autoStartup = autoStartup; - } - - /** - * Set whether the bus should automatically create a channel when a - * subscription contains the name of a previously unregistered channel. - */ - public void setAutoCreateChannels(boolean autoCreateChannels) { - this.autoCreateChannels = autoCreateChannels; - } - - /** - * Set whether the bus should configure its asynchronous task executor - * to also be used by the ApplicationContext's 'applicationEventMulticaster'. - * This will only apply if the multicaster defined within the context - * is an instance of SimpleApplicationEventMulticaster (the default). - * This property is 'false' by default. - */ - public void setConfigureAsyncEventMulticaster(boolean configureAsyncEventMulticaster) { - this.configureAsyncEventMulticaster = configureAsyncEventMulticaster; - } - - @SuppressWarnings("unchecked") - private void registerChannels(ApplicationContext context) { - Map channelBeans = (Map) context - .getBeansOfType(MessageChannel.class); - for (Map.Entry entry : channelBeans.entrySet()) { - String channelName = entry.getKey(); - MessageChannel previousChannel = this.lookupChannel(channelName); - if (previousChannel == null) { - this.registerChannel(channelName, entry.getValue()); - } - else if (!previousChannel.equals(entry.getValue())) { - throw new ConfigurationException("A different channel instance has already " - + "been registered with the name '" + channelName + "'."); - } - } - } - - @SuppressWarnings("unchecked") - private void registerEndpoints(ApplicationContext context) { - Map endpointBeans = (Map) context - .getBeansOfType(MessageEndpoint.class); - for (Map.Entry entry : endpointBeans.entrySet()) { - this.registerEndpoint(entry.getValue()); - } - } - - @SuppressWarnings("unchecked") - private void registerGateways(ApplicationContext context) { - Map gatewayBeans = (Map) context - .getBeansOfType(MessagingGateway.class); - for (Map.Entry entry : gatewayBeans.entrySet()) { - this.registerGateway(entry.getKey(), entry.getValue()); - } - } - - public void initialize() { - synchronized (this.lifecycleMonitor) { - if (this.initialized || this.initializing) { - return; - } - this.initializing = true; - if (this.taskScheduler == null) { - ScheduledThreadPoolExecutor executor = new ScheduledThreadPoolExecutor(DEFAULT_DISPATCHER_POOL_SIZE); - executor.setThreadFactory(new CustomizableThreadFactory("message-bus-")); - executor.setRejectedExecutionHandler(new CallerRunsPolicy()); - this.taskScheduler = new SimpleTaskScheduler(executor); - } - if (this.getErrorChannel() == null) { - this.setErrorChannel(new DefaultErrorChannel()); - } - this.initialized = true; - this.initializing = false; - } - } - - public MessageChannel getErrorChannel() { - return this.lookupChannel(ERROR_CHANNEL_NAME); - } - - public void setErrorChannel(MessageChannel errorChannel) { - this.registerChannel(ERROR_CHANNEL_NAME, errorChannel); - } - - public MessageChannel lookupChannel(String channelName) { - return this.channelRegistry.lookupChannel(channelName); - } - - public void registerChannel(String name, MessageChannel channel) { - if (!this.initialized) { - this.initialize(); - } - channel.setName(name); - this.channelRegistry.registerChannel(name, channel); - if (logger.isInfoEnabled()) { - logger.info("registered channel '" + name + "'"); - } - } - - public MessageChannel unregisterChannel(String name) { - return this.channelRegistry.unregisterChannel(name); - } - - public void registerHandler(String name, MessageHandler handler, Object input, Schedule schedule) { - Assert.notNull(handler, "'handler' must not be null"); - HandlerEndpoint endpoint = new HandlerEndpoint(handler); - this.configureEndpoint(endpoint, name, input, schedule); - this.registerEndpoint(endpoint); - } - - public void registerTarget(String name, MessageTarget target, Object input, Schedule schedule) { - Assert.notNull(target, "'target' must not be null"); - TargetEndpoint endpoint = new TargetEndpoint(target); - this.configureEndpoint(endpoint, name, input, schedule); - this.registerEndpoint(endpoint); - } - - private void configureEndpoint(AbstractEndpoint endpoint, String name, Object input, Schedule schedule) { - endpoint.setName(name); - if (input instanceof MessageChannel) { - endpoint.setInputChannel((MessageChannel) input); - } - else if (input instanceof String) { - endpoint.setInputChannelName((String) input); - } - else { - throw new ConfigurationException("'input' must be a MessageChannel or String"); - } - endpoint.setSchedule(schedule); - } - - public void registerEndpoint(MessageEndpoint endpoint) { - if (!this.initialized) { - this.initialize(); - } - if (endpoint instanceof ChannelRegistryAware) { - ((ChannelRegistryAware) endpoint).setChannelRegistry(this.channelRegistry); - } - this.endpointRegistry.registerEndpoint(endpoint); - if (this.isRunning()) { - this.activateEndpoint(endpoint); - } - if (logger.isInfoEnabled()) { - logger.info("registered endpoint '" + endpoint + "'"); - } - } - - public MessageEndpoint unregisterEndpoint(String name) { - MessageEndpoint endpoint = this.endpointRegistry.unregisterEndpoint(name); - if (endpoint == null) { - return null; - } - this.deactivateEndpoint(endpoint); - return endpoint; - } - - public MessageEndpoint lookupEndpoint(String endpointName) { - return this.endpointRegistry.lookupEndpoint(endpointName); - } - - public Set getEndpointNames() { - return this.endpointRegistry.getEndpointNames(); - } - - private void activateEndpoints() { - Set endpointNames = this.endpointRegistry.getEndpointNames(); - for (String name : endpointNames) { - MessageEndpoint endpoint = this.endpointRegistry.lookupEndpoint(name); - if (endpoint != null) { - this.activateEndpoint(endpoint); - } - } - } - - private void activateEndpoint(MessageEndpoint endpoint) { - Assert.notNull(endpoint, "'endpoint' must not be null"); - if (endpoint.getOutputChannel() == null) { - this.lookupOrCreateChannel(endpoint.getOutputChannelName()); - } - try { - endpoint.afterPropertiesSet(); - } - catch (Exception e) { - throw new ConfigurationException("failed to initialize endpoint", e); - } - MessageChannel channel = endpoint.getInputChannel(); - if (channel == null) { - channel = this.lookupOrCreateChannel(endpoint.getInputChannelName()); - } - if (channel != null && channel instanceof Subscribable) { - ((Subscribable) channel).subscribe(endpoint); - if (logger.isInfoEnabled()) { - logger.info("activated subscription to channel '" - + channel.getName() + "' for endpoint '" + endpoint + "'"); - } - return; - } - Schedule schedule = endpoint.getSchedule(); - EndpointTrigger trigger = endpoint.getTrigger(); - if (trigger == null) { - trigger = new EndpointTrigger(schedule != null ? schedule : this.defaultPollerSchedule); - } - trigger.addTarget(endpoint); - if (this.endpointTriggers.add(trigger)) { - this.taskScheduler.schedule(trigger); - } - } - - private MessageChannel lookupOrCreateChannel(String channelName) { - if (channelName == null) { - return null; - } - MessageChannel channel = this.lookupChannel(channelName); - if (channel == null) { - if (!this.autoCreateChannels) { - throw new ConfigurationException("Cannot activate endpoint, unknown channel '" + channelName - + "'. Consider enabling the 'autoCreateChannels' option for the message bus."); - } - if (this.logger.isInfoEnabled()) { - logger.info("auto-creating channel '" + channelName + "'"); - } - channel = channelFactory.getChannel(channelName, null, null); - this.registerChannel(channelName, channel); - } - return channel; - } - - private void registerGateway(String name, MessagingGateway gateway) { - if (gateway instanceof Lifecycle) { - this.lifecycleEndpoints.add((Lifecycle) gateway); - if (this.isRunning()) { - ((Lifecycle) gateway).start(); - } - } - if (logger.isInfoEnabled()) { - logger.info("registered gateway '" + name + "'"); - } - } - - public void deactivateEndpoint(MessageEndpoint endpoint) { - Assert.notNull(endpoint, "'endpoint' must not be null"); - for (EndpointTrigger trigger : this.endpointTriggers) { - boolean removed = trigger.removeTarget(endpoint); - if (removed && this.logger.isInfoEnabled()) { - logger.info("removed endpoint '" + endpoint + "' from dispatcher"); - } - } - if (endpoint instanceof Lifecycle) { - ((Lifecycle) endpoint).stop(); - } - } - - public boolean isRunning() { - synchronized (this.lifecycleMonitor) { - return this.running; - } - } - - public void start() { - if (!this.initialized) { - this.initialize(); - } - if (this.isRunning() || this.starting) { - return; - } - this.interceptors.preStart(); - this.starting = true; - synchronized (this.lifecycleMonitor) { - this.activateEndpoints(); - this.taskScheduler.setErrorHandler(new MessagePublishingErrorHandler(this.getErrorChannel())); - this.taskScheduler.start(); - for (Lifecycle endpoint : this.lifecycleEndpoints) { - endpoint.start(); - if (logger.isInfoEnabled()) { - logger.info("started endpoint '" + endpoint + "'"); - } - } - } - this.running = true; - this.starting = false; - this.interceptors.postStart(); - if (logger.isInfoEnabled()) { - logger.info("message bus started"); - } - } - - public void stop() { - if (!this.isRunning()) { - return; - } - this.interceptors.preStop(); - synchronized (this.lifecycleMonitor) { - this.running = false; - this.taskScheduler.stop(); - for (Lifecycle endpoint : this.lifecycleEndpoints) { - endpoint.stop(); - if (logger.isInfoEnabled()) { - logger.info("stopped endpoint '" + endpoint + "'"); - } - } - } - this.interceptors.postStop(); - if (logger.isInfoEnabled()) { - logger.info("message bus stopped"); - } - } - - public void destroy() throws Exception { - if (this.taskScheduler instanceof DisposableBean) { - ((DisposableBean) this.taskScheduler).destroy(); - } - } - - public void onApplicationEvent(ApplicationEvent event) { - if (event instanceof ContextRefreshedEvent) { - ApplicationContext context = ((ContextRefreshedEvent) event).getApplicationContext(); - this.registerChannels(context); - this.registerEndpoints(context); - this.registerGateways(context); - if (this.configureAsyncEventMulticaster) { - this.initialize(); - this.doConfigureAsyncEventMulticaster(context); - } - if (this.autoStartup) { - this.start(); - } - } - } - - private void doConfigureAsyncEventMulticaster(ApplicationContext context) { - String multicasterBeanName = AbstractApplicationContext.APPLICATION_EVENT_MULTICASTER_BEAN_NAME; - if (context.containsBean(multicasterBeanName)) { - ApplicationEventMulticaster multicaster = (ApplicationEventMulticaster) context - .getBean(multicasterBeanName); - if (multicaster instanceof SimpleApplicationEventMulticaster) { - ((SimpleApplicationEventMulticaster) multicaster).setTaskExecutor(this.taskScheduler); - } - } - } - - public void addInterceptor(MessageBusInterceptor interceptor) { - this.interceptors.add(interceptor); - } - - public void removeInterceptor(MessageBusInterceptor interceptor) { - this.interceptors.remove(interceptor); - } - - public void setInterceptors(List interceptor) { - this.interceptors.set(interceptor); - } - - /* - * Wrapper class for the interceptor list - */ - private class MessageBusInterceptorsList { - - private CopyOnWriteArrayList messageBusInterceptors = new CopyOnWriteArrayList(); - - public void set(List interceptors) { - this.messageBusInterceptors.clear(); - this.messageBusInterceptors.addAll(interceptors); - } - - public void add(MessageBusInterceptor interceptor) { - this.messageBusInterceptors.add(interceptor); - } - - public void remove(MessageBusInterceptor interceptor) { - this.messageBusInterceptors.remove(interceptor); - } - - public void preStart() { - for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { - messageBusInterceptor.preStart(MessageBus.this); - } - } - - public void postStart() { - for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { - messageBusInterceptor.postStart(MessageBus.this); - } - } - - public void preStop() { - for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { - messageBusInterceptor.preStop(MessageBus.this); - } - } - - public void postStop() { - for (MessageBusInterceptor messageBusInterceptor : messageBusInterceptors) { - messageBusInterceptor.postStop(MessageBus.this); - } - } - } + void registerTarget(String name, MessageTarget target, Object input, Schedule schedule); }