INT-1275 PriorityChannel now works properly with timeout values
This commit is contained in:
@@ -28,25 +28,29 @@ import java.util.concurrent.TimeUnit;
|
||||
* @since 2.0
|
||||
*/
|
||||
public final class UpperBound {
|
||||
|
||||
public final Semaphore semaphore;
|
||||
|
||||
/**
|
||||
* Create an UpperBound with the given capacity. If the given capacity is less than 1
|
||||
* an infinite UpperBound is created.
|
||||
*
|
||||
* @param capacity
|
||||
*/
|
||||
public UpperBound(int capacity) {
|
||||
this.semaphore = (capacity > 0) ? new Semaphore(capacity, true) : null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Acquires a lock on the underlying semaphore if this UpperBound is bounded and returns true
|
||||
* if it succeeds within the given timeout
|
||||
* Acquires a permit from the underlying semaphore if this UpperBound is bounded and returns
|
||||
* true if it succeeds within the given timeout. If the timeout is less than 0, it will block
|
||||
* indefinitely.
|
||||
*/
|
||||
public boolean tryAcquire(long timeoutInMilliseconds) {
|
||||
if (this.semaphore != null) {
|
||||
try {
|
||||
if (timeoutInMilliseconds < 0) {
|
||||
this.semaphore.acquire();
|
||||
return true;
|
||||
}
|
||||
return this.semaphore.tryAcquire(timeoutInMilliseconds, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
@@ -66,4 +70,5 @@ public final class UpperBound {
|
||||
this.semaphore.release();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2008 the original author or authors.
|
||||
* Copyright 2002-2010 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.
|
||||
@@ -18,9 +18,16 @@ 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.Comparator;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Executor;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
@@ -113,6 +120,78 @@ public class PriorityChannelTests {
|
||||
assertEquals("test-LOW", channel.receive(0).getPayload());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testTimeoutElapses() throws InterruptedException {
|
||||
final PriorityChannel channel = new PriorityChannel(1);
|
||||
final AtomicBoolean sentSecondMessage = new AtomicBoolean(false);
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
Executor executor = Executors.newSingleThreadScheduledExecutor();
|
||||
channel.send(new StringMessage("test-1"));
|
||||
executor.execute(new Runnable() {
|
||||
public void run() {
|
||||
sentSecondMessage.set(channel.send(new StringMessage("test-2"), 10));
|
||||
latch.countDown();
|
||||
}
|
||||
});
|
||||
assertFalse(sentSecondMessage.get());
|
||||
Thread.sleep(500);
|
||||
Message<?> message1 = channel.receive();
|
||||
assertNotNull(message1);
|
||||
assertEquals("test-1", message1.getPayload());
|
||||
latch.await(1000, TimeUnit.MILLISECONDS);
|
||||
assertFalse(sentSecondMessage.get());
|
||||
assertNull(channel.receive(0));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testTimeoutDoesNotElapse() throws InterruptedException {
|
||||
final PriorityChannel channel = new PriorityChannel(1);
|
||||
final AtomicBoolean sentSecondMessage = new AtomicBoolean(false);
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
Executor executor = Executors.newSingleThreadScheduledExecutor();
|
||||
channel.send(new StringMessage("test-1"));
|
||||
executor.execute(new Runnable() {
|
||||
public void run() {
|
||||
sentSecondMessage.set(channel.send(new StringMessage("test-2"), 3000));
|
||||
latch.countDown();
|
||||
}
|
||||
});
|
||||
assertFalse(sentSecondMessage.get());
|
||||
Thread.sleep(500);
|
||||
Message<?> message1 = channel.receive();
|
||||
assertNotNull(message1);
|
||||
assertEquals("test-1", message1.getPayload());
|
||||
latch.await(1000, TimeUnit.MILLISECONDS);
|
||||
assertTrue(sentSecondMessage.get());
|
||||
Message<?> message2 = channel.receive();
|
||||
assertNotNull(message2);
|
||||
assertEquals("test-2", message2.getPayload());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testIndefiniteTimeout() throws InterruptedException {
|
||||
final PriorityChannel channel = new PriorityChannel(1);
|
||||
final AtomicBoolean sentSecondMessage = new AtomicBoolean(false);
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
Executor executor = Executors.newSingleThreadScheduledExecutor();
|
||||
channel.send(new StringMessage("test-1"));
|
||||
executor.execute(new Runnable() {
|
||||
public void run() {
|
||||
sentSecondMessage.set(channel.send(new StringMessage("test-2"), -1));
|
||||
latch.countDown();
|
||||
}
|
||||
});
|
||||
assertFalse(sentSecondMessage.get());
|
||||
Message<?> message1 = channel.receive();
|
||||
assertNotNull(message1);
|
||||
assertEquals("test-1", message1.getPayload());
|
||||
latch.await(1000, TimeUnit.MILLISECONDS);
|
||||
assertTrue(sentSecondMessage.get());
|
||||
Message<?> message2 = channel.receive();
|
||||
assertNotNull(message2);
|
||||
assertEquals("test-2", message2.getPayload());
|
||||
}
|
||||
|
||||
|
||||
private static Message<String> createPriorityMessage(MessagePriority priority) {
|
||||
return MessageBuilder.withPayload("test-" + priority).setPriority(priority).build();
|
||||
|
||||
Reference in New Issue
Block a user