From a7cfbe7753a6a26d48b7b81428eed65e9965768d Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Fri, 28 Dec 2007 18:04:08 +0000 Subject: [PATCH] Moved test with channel and endpoint registration into MessageBusTests and added an isolated test for the UnicastMessageDispatcher. --- .../integration/bus/MessageBusTests.java | 26 +++++++++ .../bus/UnicastMessageDispatcherTests.java | 56 +++++++++++-------- 2 files changed, 60 insertions(+), 22 deletions(-) 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 37abf14c0d..aa212f73dd 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 @@ -18,6 +18,7 @@ package org.springframework.integration.bus; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; import org.junit.Test; @@ -84,4 +85,29 @@ public class MessageBusTests { assertEquals("test", result.getPayload()); } + @Test + public void testExactlyOneEndpointReceivesUnicastMessage() { + PointToPointChannel inputChannel = new PointToPointChannel(); + PointToPointChannel outputChannel1 = new PointToPointChannel(); + PointToPointChannel outputChannel2 = new PointToPointChannel(); + GenericMessageEndpoint endpoint1 = new GenericMessageEndpoint(); + endpoint1.setDefaultOutputChannelName("output1"); + endpoint1.setInputChannelName("input"); + GenericMessageEndpoint endpoint2 = new GenericMessageEndpoint(); + endpoint2.setDefaultOutputChannelName("output2"); + endpoint2.setInputChannelName("input"); + MessageBus bus = new MessageBus(); + bus.registerChannel("input", inputChannel); + bus.registerChannel("output1", outputChannel1); + bus.registerChannel("output2", outputChannel2); + bus.registerEndpoint("endpoint1", endpoint1); + bus.registerEndpoint("endpoint2", endpoint2); + bus.start(); + inputChannel.send(new StringMessage(1, "testing")); + Message message1 = outputChannel1.receive(100); + Message message2 = outputChannel2.receive(0); + bus.stop(); + assertTrue("exactly one message should be null", message1 == null ^ message2 == null); + } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/UnicastMessageDispatcherTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/UnicastMessageDispatcherTests.java index fd0889acef..45488c5865 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/UnicastMessageDispatcherTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/UnicastMessageDispatcherTests.java @@ -18,10 +18,15 @@ package org.springframework.integration.bus; import static org.junit.Assert.assertTrue; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + import org.junit.Test; import org.springframework.integration.channel.PointToPointChannel; import org.springframework.integration.endpoint.GenericMessageEndpoint; +import org.springframework.integration.endpoint.MessageEndpoint; import org.springframework.integration.message.Message; import org.springframework.integration.message.StringMessage; @@ -31,28 +36,35 @@ import org.springframework.integration.message.StringMessage; public class UnicastMessageDispatcherTests { @Test - public void testExactlyOneEndpointReceivesMessage() { - PointToPointChannel inputChannel = new PointToPointChannel(); - PointToPointChannel outputChannel1 = new PointToPointChannel(); - PointToPointChannel outputChannel2 = new PointToPointChannel(); - GenericMessageEndpoint endpoint1 = new GenericMessageEndpoint(); - endpoint1.setDefaultOutputChannelName("output1"); - endpoint1.setInputChannelName("input"); - GenericMessageEndpoint endpoint2 = new GenericMessageEndpoint(); - endpoint2.setDefaultOutputChannelName("output2"); - endpoint2.setInputChannelName("input"); - MessageBus bus = new MessageBus(); - bus.registerChannel("input", inputChannel); - bus.registerChannel("output1", outputChannel1); - bus.registerChannel("output2", outputChannel2); - bus.registerEndpoint("endpoint1", endpoint1); - bus.registerEndpoint("endpoint2", endpoint2); - bus.start(); - inputChannel.send(new StringMessage(1, "testing")); - Message message1 = outputChannel1.receive(100); - Message message2 = outputChannel2.receive(0); - bus.stop(); - assertTrue("exactly one message should be null", message1 == null ^ message2 == null); + public void testDispatcherSendsToExactlyOneEndpoint() throws InterruptedException { + final AtomicBoolean endpoint1Received = new AtomicBoolean(); + final AtomicBoolean endpoint2Received = new AtomicBoolean(); + final CountDownLatch latch = new CountDownLatch(1); + MessageEndpoint endpoint1 = new GenericMessageEndpoint() { + @Override + public void messageReceived(Message message) { + endpoint1Received.set(true); + latch.countDown(); + } + }; + MessageEndpoint endpoint2 = new GenericMessageEndpoint() { + @Override + public void messageReceived(Message message) { + endpoint2Received.set(true); + latch.countDown(); + } + }; + ConsumerPolicy policy = new ConsumerPolicy(); + PointToPointChannel channel = new PointToPointChannel(); + channel.send(new StringMessage(1, "test")); + MessageRetriever retriever = new ChannelPollingMessageRetriever(channel, policy); + UnicastMessageDispatcher dispatcher = new UnicastMessageDispatcher(retriever, policy); + dispatcher.addEndpointExecutor(new EndpointExecutor(endpoint1, 1, 1)); + dispatcher.addEndpointExecutor(new EndpointExecutor(endpoint2, 1, 1)); + dispatcher.receiveAndDispatch(); + latch.await(500, TimeUnit.MILLISECONDS); + assertTrue("exactly one endpoint should have received message", + endpoint1Received.get() ^ endpoint2Received.get()); } }