From cfa00c5707b24cfbca4f8029e8087c8c068ada2f Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 1 Jul 2016 11:11:09 -0400 Subject: [PATCH] INT-1463 Fix PriorityChannel#getRemainingCapacity JIRA: https://jira.spring.io/browse/INT-1463 An existing `QueueChannel#getRemainingCapacity()` is based on the `BlockingQueue#getRemainingCapacity()` which is always `Integer.MAX_VALUE` and doesn't reflect reality for provided `capacity` * Fix `PriorityChannel` to return `upperBound.availablePermits()` for `getRemainingCapacity()` * Add test for `PriorityChannel` on the matter * Also increase `reply-timeout` for `` in `RedisQueueGatewayIntegrationTests-context.xml` from 1 sec to 10 secs --- .../integration/channel/PriorityChannel.java | 5 +++++ ...ChannelCapacityPlaceholderTests-context.xml | 7 ++++--- .../ChannelCapacityPlaceholderTests.java | 18 ++++++++++++------ ...disQueueGatewayIntegrationTests-context.xml | 2 +- 4 files changed, 22 insertions(+), 10 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/PriorityChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/PriorityChannel.java index f664e01440..60dd1f00fb 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/PriorityChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/PriorityChannel.java @@ -84,6 +84,11 @@ public class PriorityChannel extends QueueChannel { this(0, null); } + @Override + public int getRemainingCapacity() { + return this.upperBound.availablePermits(); + } + @Override protected boolean doSend(Message message, long timeout) { if (!this.upperBound.tryAcquire(timeout)) { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests-context.xml index 182d2d1dd2..d379d9c854 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests-context.xml @@ -13,11 +13,12 @@ - - + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests.java index 8cbc89e50a..4e0ed1f4cc 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/config/ChannelCapacityPlaceholderTests.java @@ -24,6 +24,7 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; +import org.springframework.integration.channel.PriorityChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.support.MessageBuilder; import org.springframework.test.context.ContextConfiguration; @@ -31,6 +32,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Mark Fisher + * @author Artem Bilan */ @ContextConfiguration @RunWith(SpringJUnit4ClassRunner.class) @@ -51,12 +53,16 @@ public class ChannelCapacityPlaceholderTests { assertEquals(98, channel.getRemainingCapacity()); } - - public interface TestService { - - @org.springframework.integration.annotation.Gateway(requestChannel = "channel") - void test(); - + @Test + public void testCapacityOnPriorityChannel() { + PriorityChannel channel = context.getBean("priorityChannel", PriorityChannel.class); + assertNotNull(channel); + assertEquals(99, channel.getRemainingCapacity()); + channel.send(MessageBuilder.withPayload("test1").build()); + channel.send(MessageBuilder.withPayload("test2").build()); + assertEquals(97, channel.getRemainingCapacity()); + assertNotNull(channel.receive(0)); + assertEquals(98, channel.getRemainingCapacity()); } } diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueGatewayIntegrationTests-context.xml b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueGatewayIntegrationTests-context.xml index 21ac79d59f..836fd30791 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueGatewayIntegrationTests-context.xml +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueGatewayIntegrationTests-context.xml @@ -26,7 +26,7 @@