diff --git a/concurrency-limits-core/src/main/java/com/netflix/concurrency/limits/executors/BlockingAdaptiveExecutor.java b/concurrency-limits-core/src/main/java/com/netflix/concurrency/limits/executors/BlockingAdaptiveExecutor.java index 8c06b51e..a76bef88 100644 --- a/concurrency-limits-core/src/main/java/com/netflix/concurrency/limits/executors/BlockingAdaptiveExecutor.java +++ b/concurrency-limits-core/src/main/java/com/netflix/concurrency/limits/executors/BlockingAdaptiveExecutor.java @@ -86,10 +86,12 @@ public Thread newThread(Runnable r) { } if (limiter == null) { - limiter = SimpleLimiter.newBuilder() + limiter = BlockingLimiter.wrap(SimpleLimiter.newBuilder() .metricRegistry(metricRegistry) .limit(AIMDLimit.newBuilder().build()) - .build(); + .build()); + } else if (!(limiter instanceof BlockingLimiter)) { + limiter = BlockingLimiter.wrap(limiter); } return new BlockingAdaptiveExecutor(this); diff --git a/concurrency-limits-core/src/test/java/com/netflix/concurrency/limits/executor/BlockingAdaptiveExecutorTest.java b/concurrency-limits-core/src/test/java/com/netflix/concurrency/limits/executor/BlockingAdaptiveExecutorTest.java new file mode 100644 index 00000000..7880328d --- /dev/null +++ b/concurrency-limits-core/src/test/java/com/netflix/concurrency/limits/executor/BlockingAdaptiveExecutorTest.java @@ -0,0 +1,77 @@ +package com.netflix.concurrency.limits.executor; + +import com.netflix.concurrency.limits.Limiter; +import com.netflix.concurrency.limits.executors.BlockingAdaptiveExecutor; +import com.netflix.concurrency.limits.limit.SettableLimit; +import com.netflix.concurrency.limits.limiter.SimpleLimiter; +import org.junit.Assert; +import org.junit.Test; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +public class BlockingAdaptiveExecutorTest { + + @Test + public void testBuilderWithNoLimiter() { + BlockingAdaptiveExecutor executor = BlockingAdaptiveExecutor.newBuilder().build(); + AtomicBoolean executed = new AtomicBoolean(false); + executor.execute(() -> executed.set(true)); + + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + + Assert.assertTrue("Task should have been executed", executed.get()); + } + + @Test + public void testBuilderWithSimpleLimiter() throws InterruptedException { + SettableLimit limit = SettableLimit.startingAt(2); + Limiter simpleLimiter = SimpleLimiter.newBuilder() + .limit(limit) + .build(); + + BlockingAdaptiveExecutor executor = BlockingAdaptiveExecutor.newBuilder() + .limiter(simpleLimiter) + .build(); + + CountDownLatch taskStarted = new CountDownLatch(2); + CountDownLatch taskComplete = new CountDownLatch(1); + CountDownLatch thirdTaskComplete = new CountDownLatch(1); + + for (int i = 0; i < 2; i++) { + executor.execute(() -> { + taskStarted.countDown(); + try { + taskComplete.await(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); + } + + Assert.assertTrue("Tasks should start", taskStarted.await(1, TimeUnit.SECONDS)); + + AtomicBoolean thirdTaskStarted = new AtomicBoolean(false); + Thread blockedThread = new Thread(() -> + executor.execute(() -> { + thirdTaskStarted.set(true); + thirdTaskComplete.countDown(); + })); + + blockedThread.start(); + Thread.sleep(200); + + Assert.assertFalse("Third task should be blocked", thirdTaskStarted.get()); + + taskComplete.countDown(); + + Assert.assertTrue("Third task should eventually execute", + thirdTaskComplete.await(1, TimeUnit.SECONDS)); + } + +} diff --git a/concurrency-limits-servlet-jakarta/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java b/concurrency-limits-servlet-jakarta/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java index 98fcc3b7..29cb42ce 100644 --- a/concurrency-limits-servlet-jakarta/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java +++ b/concurrency-limits-servlet-jakarta/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java @@ -45,8 +45,8 @@ public class ConcurrencyLimitServletFilterSimulationTest { @Ignore public void simulation() throws Exception { Limit limit = VegasLimit.newDefault(); - BlockingAdaptiveExecutor executor = new BlockingAdaptiveExecutor( - SimpleLimiter.newBuilder().limit(limit).build()); + BlockingAdaptiveExecutor executor = BlockingAdaptiveExecutor.newBuilder().limiter( + SimpleLimiter.newBuilder().limit(limit).build()).build(); AtomicInteger errors = new AtomicInteger(); AtomicInteger success = new AtomicInteger(); diff --git a/concurrency-limits-servlet/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java b/concurrency-limits-servlet/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java index 0012a662..695b71a7 100644 --- a/concurrency-limits-servlet/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java +++ b/concurrency-limits-servlet/src/test/java/com/netflix/concurrency/limits/ConcurrencyLimitServletFilterSimulationTest.java @@ -46,8 +46,8 @@ public class ConcurrencyLimitServletFilterSimulationTest { @Ignore public void simulation() throws Exception { Limit limit = VegasLimit.newDefault(); - BlockingAdaptiveExecutor executor = new BlockingAdaptiveExecutor( - SimpleLimiter.newBuilder().limit(limit).build()); + BlockingAdaptiveExecutor executor = BlockingAdaptiveExecutor.newBuilder().limiter( + SimpleLimiter.newBuilder().limit(limit).build()).build(); AtomicInteger errors = new AtomicInteger(); AtomicInteger success = new AtomicInteger();