diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/registry/ChannelRegistry.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/registry/ChannelRegistry.java new file mode 100644 index 0000000000..22a14da11c --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/registry/ChannelRegistry.java @@ -0,0 +1,50 @@ +/* + * Copyright 2002-2013 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.channel.registry; + +import org.springframework.integration.MessageChannel; + +/** + * A strategy interface used to bind a {@link MessageChannel} to a logical name. The name + * is intended to identify a logical consumer or producer of messages. This may be a + * queue, a channel adapter, another message channel, a Spring bean, etc. + * + * @author Mark Fisher + * @author David Turanski + * @since 3.0 + */ +public interface ChannelRegistry { + + /** + * Register a message consumer + * @param name the logical identity of the message source + * @param channel the channel bound as a consumer + */ + void inbound(String name, MessageChannel channel); + + /** + * Register a message producer + * @param name the logical identity of the message target + * @param channel the channel bound as a producer + */ + void outbound(String name, MessageChannel channel); + + /** + * Create a tap on an already registered inbound channel + * @param the registered name + * @param channel the channel that will receive messages from the tap + */ + void tap(String name, MessageChannel channel); + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/registry/LocalChannelRegistry.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/registry/LocalChannelRegistry.java new file mode 100644 index 0000000000..928ab1e661 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/registry/LocalChannelRegistry.java @@ -0,0 +1,168 @@ +/* + * Copyright 2002-2013 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.channel.registry; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.support.AbstractApplicationContext; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.channel.AbstractMessageChannel; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.channel.PublishSubscribeChannel; +import org.springframework.integration.channel.interceptor.WireTap; +import org.springframework.integration.core.SubscribableChannel; +import org.springframework.integration.handler.BridgeHandler; +import org.springframework.util.Assert; + +/** + * A simple implementation of {@link ChannelRegistry} for in-process use. For inbound and + * outbound, creates a {@link DirectChannel} and bridges the passed + * {@link MessageChannel} to the channel which is registered in the given application + * context. If that channel does not yet exist, it will be created. For tap, it adds a + * {@link WireTap} for an inbound channel whose name matches the one provided. If no such + * inbound channel exists at the time of the method invocation, it will throw an + * Exception. Otherwise the provided channel instance will receive messages from the wire + * tap on that inbound channel. + * + * @author David Turanski + * @author Mark Fisher + * @since 3.0 + */ +public class LocalChannelRegistry implements ChannelRegistry, ApplicationContextAware, InitializingBean { + + private volatile AbstractApplicationContext applicationContext; + + @Override + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + Assert.isInstanceOf(AbstractApplicationContext.class, applicationContext); + this.applicationContext = (AbstractApplicationContext) applicationContext; + } + + @Override + public void afterPropertiesSet() throws Exception { + Assert.notNull(applicationContext, "The 'applicationContext' property cannot be null"); + } + + /** + * Looks up or creates a DirectChannel with the given name and creates a bridge from + * that channel to the provided channel instance. Also registers a wire tap if the + * channel for the given name had been created. The target of the wire tap is a + * publish-subscribe channel. + */ + @Override + public void inbound(String name, MessageChannel channel) { + Assert.hasText(name, "a valid name is required to register an inbound channel"); + Assert.notNull(channel, "channel must not be null"); + DirectChannel registeredChannel = lookupOrCreateSharedChannel(name, DirectChannel.class); + bridge(registeredChannel, channel); + createSharedTapChannelIfNecessary(registeredChannel); + } + + /** + * Looks up or creates a DirectChannel with the given name and creates a bridge to + * that channel from the provided channel instance. + */ + @Override + public void outbound(String name, MessageChannel channel) { + Assert.hasText(name, "a valid name is required to register an outbound channel"); + Assert.notNull(channel, "channel must not be null"); + Assert.isTrue(channel instanceof SubscribableChannel, + "channel must be of type " + SubscribableChannel.class.getName()); + DirectChannel registeredChannel = lookupOrCreateSharedChannel(name, DirectChannel.class); + bridge((SubscribableChannel) channel, registeredChannel); + } + + /** + * Looks up a wiretap for the inbound channel with the given name and creates a + * bridge from that wiretap's output channel to the provided channel instance. + * Will throw an Exception if no such wiretap exists. + */ + @Override + public void tap(String name, MessageChannel channel) { + Assert.hasText(name, "a valid name is required to register a tap channel"); + Assert.notNull(channel, "channel must not be null"); + SubscribableChannel tapChannel = null; + String tapName = name + ".tap"; + try { + tapChannel = applicationContext.getBean(tapName, SubscribableChannel.class); + } + catch (Exception e) { + throw new IllegalArgumentException("No tap channel exists for '" + name + + "'. A tap is only valid for a registered inbound channel."); + } + bridge(tapChannel, channel); + } + + private synchronized T lookupOrCreateSharedChannel(String name, Class requiredType) { + T channel = null; + if (applicationContext.containsBean(name)) { + try { + channel = applicationContext.getBean(name, requiredType); + } + catch (Exception e) { + throw new IllegalArgumentException("bean '" + name + + "' is already registered but does not match the required type"); + } + } + else { + channel = createSharedChannel(name, requiredType); + } + return channel; + } + + private T createSharedChannel(String name, Class requiredType) { + try { + T channel = requiredType.newInstance(); + channel.setComponentName(name); + channel.setBeanFactory(applicationContext); + channel.setBeanName(name); + channel.afterPropertiesSet(); + applicationContext.getBeanFactory().registerSingleton(name, channel); + return channel; + } + catch (Exception e) { + throw new IllegalArgumentException("failed to create channel: " + name, e); + } + } + + private synchronized void createSharedTapChannelIfNecessary(AbstractMessageChannel channel) { + String tapName = channel.getComponentName() + ".tap"; + PublishSubscribeChannel tapChannel = null; + if (!applicationContext.containsBean(tapName)) { + tapChannel = createSharedChannel(tapName, PublishSubscribeChannel.class); + WireTap wireTap = new WireTap(tapChannel); + channel.addInterceptor(wireTap); + } + else { + try { + tapChannel = applicationContext.getBean(tapName, PublishSubscribeChannel.class); + } + catch (Exception e) { + throw new IllegalArgumentException("bean '" + tapName + + "' is already registered but does not match the required type"); + } + } + } + + private BridgeHandler bridge(SubscribableChannel from, MessageChannel to) { + BridgeHandler handler = new BridgeHandler(); + handler.setOutputChannel(to); + handler.afterPropertiesSet(); + from.subscribe(handler); + return handler; + } + +} diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/registry/LocalChannelRegistryTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/registry/LocalChannelRegistryTests.java new file mode 100644 index 0000000000..1cd601dd13 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/registry/LocalChannelRegistryTests.java @@ -0,0 +1,171 @@ +/* + * Copyright 2002-2013 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.channel.registry; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.Before; +import org.junit.Test; + +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.GenericApplicationContext; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.MessagingException; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.core.MessageHandler; +import org.springframework.integration.core.SubscribableChannel; +import org.springframework.integration.message.GenericMessage; + +/** + * @author David Turanski + * @author Mark Fisher + * @since 3.0 + */ +public class LocalChannelRegistryTests { + + private LocalChannelRegistry registry = new LocalChannelRegistry(); + + private ApplicationContext context = new GenericApplicationContext(); + + @Before + public void setUp() { + registry.setApplicationContext(context); + } + + @Test + public void testInbound() { + DirectChannel channel = new DirectChannel(); + registry.inbound("inbound", channel); + assertTrue(context.containsBean("inbound")); + + final AtomicBoolean messageReceived = new AtomicBoolean(); + channel.subscribe(new MessageHandler() { + @Override + public void handleMessage(Message message) throws MessagingException { + messageReceived.set(true); + assertEquals("hello", message.getPayload()); + } + }); + SubscribableChannel registeredChannel = context.getBean("inbound", SubscribableChannel.class); + registeredChannel.send(new GenericMessage("hello")); + assertTrue(messageReceived.get()); + } + + @Test + public void testOutbound() { + DirectChannel channel = new DirectChannel(); + registry.outbound("outbound", channel); + assertTrue(context.containsBean("outbound")); + + final AtomicBoolean messageReceived = new AtomicBoolean(); + SubscribableChannel registeredChannel = context.getBean("outbound", SubscribableChannel.class); + registeredChannel.subscribe(new MessageHandler() { + @Override + public void handleMessage(Message message) throws MessagingException { + messageReceived.set(true); + assertEquals("hello", message.getPayload()); + } + }); + channel.send(new GenericMessage("hello")); + assertTrue(messageReceived.get()); + } + + @Test(expected = IllegalArgumentException.class) + public void testOutboundTapShouldFail() { + DirectChannel channel = new DirectChannel(); + registry.outbound("outbound", channel); + DirectChannel tapChannel = new DirectChannel(); + registry.tap("outbound", tapChannel); + } + + @Test + public void testInboundTap() { + DirectChannel channel = new DirectChannel(); + registry.inbound("inbound", channel); + DirectChannel tapChannel = new DirectChannel(); + registry.tap("inbound", tapChannel); + final AtomicBoolean originalMessageReceived = new AtomicBoolean(); + final AtomicBoolean tapMessageReceived = new AtomicBoolean(); + tapChannel.subscribe(new MessageHandler() { + @Override + public void handleMessage(Message message) throws MessagingException { + tapMessageReceived.set(true); + assertEquals("hello", message.getPayload()); + } + }); + channel.subscribe(new MessageHandler() { + @Override + public void handleMessage(Message message) throws MessagingException { + originalMessageReceived.set(true); + assertEquals("hello", message.getPayload()); + } + }); + MessageChannel registeredChannel = context.getBean("inbound", MessageChannel.class); + registeredChannel.send(new GenericMessage("hello")); + assertTrue(originalMessageReceived.get()); + assertTrue(tapMessageReceived.get()); + } + + @Test + public void testFlowThroughRegisteredChannelFromOutboundToInbound() { + DirectChannel outbound = new DirectChannel(); + DirectChannel inbound = new DirectChannel(); + registry.outbound("foo", outbound); + registry.inbound("foo", inbound); + final AtomicBoolean messageReceived = new AtomicBoolean(); + inbound.subscribe(new MessageHandler() { + @Override + public void handleMessage(Message message) throws MessagingException { + messageReceived.set(true); + assertEquals("hello", message.getPayload()); + } + }); + outbound.send(new GenericMessage("hello")); + assertTrue(messageReceived.get()); + } + + @Test + public void testFlowThroughRegisteredChannelFromOutboundToInboundWithTap() { + DirectChannel outbound = new DirectChannel(); + DirectChannel inbound = new DirectChannel(); + DirectChannel tap = new DirectChannel(); + registry.outbound("foo", outbound); + registry.inbound("foo", inbound); + registry.tap("foo", tap); + final AtomicBoolean originalMessageReceived = new AtomicBoolean(); + final AtomicBoolean tapMessageReceived = new AtomicBoolean(); + inbound.subscribe(new MessageHandler() { + @Override + public void handleMessage(Message message) throws MessagingException { + originalMessageReceived.set(true); + assertEquals("hello", message.getPayload()); + } + }); + tap.subscribe(new MessageHandler() { + @Override + public void handleMessage(Message message) throws MessagingException { + tapMessageReceived.set(true); + assertEquals("hello", message.getPayload()); + } + }); + outbound.send(new GenericMessage("hello")); + assertTrue(originalMessageReceived.get()); + assertTrue(tapMessageReceived.get()); + } + +}