From 318a34593dfb167b0b799d4badd440307f988486 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Wed, 9 Jan 2008 13:17:23 +0000 Subject: [PATCH] Removed MessageSelector awareness from the PointToPointChannel implementation. --- .../channel/PointToPointChannel.java | 48 +------------ .../channel/PointToPointChannelTests.java | 72 +------------------ 2 files changed, 2 insertions(+), 118 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/PointToPointChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/PointToPointChannel.java index e89caf7100..f6f2fd74b6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/PointToPointChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/PointToPointChannel.java @@ -22,10 +22,9 @@ import java.util.concurrent.TimeUnit; import org.springframework.beans.factory.BeanNameAware; import org.springframework.integration.message.Message; -import org.springframework.integration.message.MessageSelector; /** - * Simple implementation of a point-to-point message channel. Each Messages is + * Simple implementation of a message channel. Each {@link Message} is * placed in a queue whose capacity may be specified upon construction. If no * capacity is specified, the {@link #DEFAULT_CAPACITY} will be used. * @@ -165,49 +164,4 @@ public class PointToPointChannel implements MessageChannel, BeanNameAware { } } - /** - * Receive the first message that is accepted by the specified selector - * starting from the head of the queue. If the queue is empty, this method - * will block until the allotted timeout elapses. If the specified timeout - * is less than 1, the method will return immediately. - * - * @param selector the selector to use - * @param timeout the timeout in milliseconds - * - * @return the first accepted message or null in case the - * selector does not accept any message within the allotted time or the - * receiving thread is interrupted. - */ - public Message receive(MessageSelector selector, long timeout) { - long start = System.currentTimeMillis(); - while (timeout <= 0 || System.currentTimeMillis() - start < timeout) { - Object[] elements = this.queue.toArray(); - for (int i = (elements.length - 1); i >= 0; i--) { - Message m = (Message) elements[i]; - if (selector.accept(m) && this.queue.remove(m)) { - return m; - } - } - if (timeout == 0) { - return null; - } - } - return null; - } - - /** - * Receive the first message that is accepted by the specified selector - * starting from the head of the queue. If the queue is empty, this method - * will block. - * - * @param selector the selector to use - * - * @return the first accepted message or null in case the - * selector does not accept any message or the receiving thread is - * interrupted. - */ - public Message receive(MessageSelector selector) { - return this.receive(selector, -1); - } - } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/PointToPointChannelTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/PointToPointChannelTests.java index 53a59db6e8..d4bc1c423a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/PointToPointChannelTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/PointToPointChannelTests.java @@ -16,10 +16,7 @@ package org.springframework.integration.channel; -import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; import java.util.concurrent.CountDownLatch; @@ -27,11 +24,10 @@ import java.util.concurrent.Executor; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicReference; import org.junit.Test; + import org.springframework.integration.message.Message; -import org.springframework.integration.message.MessageSelector; import org.springframework.integration.message.GenericMessage; /** @@ -193,70 +189,4 @@ public class PointToPointChannelTests { assertTrue(sendInterrupted.get()); } - @Test - public void testSelectorMatchesWithinTimeout() throws Exception { - final PointToPointChannel channel = new PointToPointChannel(); - final CountDownLatch latch = new CountDownLatch(1); - final AtomicReference> messageRef = new AtomicReference>(); - Thread receiver = new Thread(new Runnable() { - public void run() { - Message message = channel.receive(new MessageSelector() { - public boolean accept(Message message) { - return (((Integer)message.getId()).intValue() == 3); - } - }, 500); - messageRef.set(message); - latch.countDown(); - } - }); - receiver.start(); - Thread sender = new Thread(new Runnable() { - public void run() { - channel.send(new GenericMessage(1, "test-1")); - try { Thread.sleep(5); } catch (Exception e) {} - channel.send(new GenericMessage(2, "test-2")); - try { Thread.sleep(5); } catch (Exception e) {} - channel.send(new GenericMessage(3, "test-3")); - try { Thread.sleep(100); } catch (Exception e) {} - channel.send(new GenericMessage(4, "test-4")); - } - }); - sender.start(); - latch.await(); - assertNotNull("message reference should not be null", messageRef.get()); - String payload = (String) messageRef.get().getPayload(); - assertEquals("expected 'test-3', but message was '" + payload, "test-3", payload); - } - - @Test - public void testSelectorDoesNotMatchWithinTimeout() throws Exception { - final PointToPointChannel channel = new PointToPointChannel(); - final CountDownLatch latch = new CountDownLatch(1); - final AtomicReference> messageRef = new AtomicReference>(); - Thread receiver = new Thread(new Runnable() { - public void run() { - Message message = channel.receive(new MessageSelector() { - public boolean accept(Message message) { - return (((Integer)message.getId()).intValue() == 3); - } - }, 7); - messageRef.set(message); - latch.countDown(); - } - }); - receiver.start(); - Thread sender = new Thread(new Runnable() { - public void run() { - channel.send(new GenericMessage(1, "test-1")); - try { Thread.sleep(5); } catch (Exception e) {} - channel.send(new GenericMessage(2, "test-2")); - try { Thread.sleep(5); } catch (Exception e) {} - channel.send(new GenericMessage(3, "test-3")); - } - }); - sender.start(); - latch.await(); - assertNull(messageRef.get()); - } - }