diff --git a/core/camel-core/src/test/java/org/apache/camel/component/direct/DirectProducerBlockingTest.java b/core/camel-core/src/test/java/org/apache/camel/component/direct/DirectProducerBlockingTest.java index 4ad06aa15819c..a692a4685b416 100644 --- a/core/camel-core/src/test/java/org/apache/camel/component/direct/DirectProducerBlockingTest.java +++ b/core/camel-core/src/test/java/org/apache/camel/component/direct/DirectProducerBlockingTest.java @@ -16,8 +16,8 @@ */ package org.apache.camel.component.direct; -import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import org.apache.camel.CamelExchangeException; @@ -28,7 +28,6 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; -import static org.awaitility.Awaitility.await; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -74,32 +73,26 @@ public void testProducerBlocksWithNoConsumers() throws Exception { @Test public void testProducerBlocksResumeTest() throws Exception { - getMockEndpoint("mock:result").expectedMessageCount(1); - context.getRouteController().suspendRoute("foo"); - Thread mainThread = Thread.currentThread(); - ExecutorService executor = Executors.newSingleThreadExecutor(); - executor.submit(new Runnable() { - @Override - public void run() { - try { - // Wait for the main thread to enter TIMED_WAITING state - // (blocked on condition in DirectComponent.getConsumer) - await().atMost(2, TimeUnit.SECONDS) - .pollInterval(10, TimeUnit.MILLISECONDS) - .until(() -> mainThread.getState() == Thread.State.TIMED_WAITING); - - log.info("Resuming consumer"); - context.getRouteController().resumeRoute("foo"); - } catch (Exception e) { - log.error("Error in background thread", e); - } + // Schedule route resume after 200ms. This is race-tolerant by design: + // - If resume fires before sendBody starts blocking: route is already + // resumed, sendBody finds the consumer immediately and succeeds + // - If resume fires while sendBody is blocking: sendBody gets unblocked + // Either outcome produces a passing test. + ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); + executor.schedule(() -> { + try { + log.info("Resuming consumer"); + context.getRouteController().resumeRoute("foo"); + } catch (Exception e) { + log.error("Error resuming route", e); } - }); + }, 200, TimeUnit.MILLISECONDS); + + getMockEndpoint("mock:result").expectedMessageCount(1); - // This call will block until the route is resumed by the background thread - template.sendBody("direct:suspended?block=true&timeout=2000", "hello world"); + template.sendBody("direct:suspended?block=true&timeout=1000", "hello world"); assertMockEndpointsSatisfied();