This is an automated email from the ASF dual-hosted git repository.

apupier pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git

commit 12bd1f030d41cf99243b897fc011d30f910dda85
Author: Guillaume Nodet <[email protected]>
AuthorDate: Tue Jul 21 14:18:39 2026 +0000

    CAMEL-24063: Revert Awaitility-based DirectProducerBlockingTest to 
race-tolerant design
    
    The Awaitility-based TIMED_WAITING approach (introduced in CAMEL-19549,
    Jun 25/29) imposed a strict thread-ordering dependency that introduced
    flakiness where the original code had none (stable for 6 years).
    
    Replace with ScheduledExecutorService.schedule() — same 200ms delay
    semantics as the original Thread.sleep(200), without using Thread.sleep:
    - If resume fires before sendBody starts blocking: route is already
      resumed, sendBody finds the consumer immediately (pass)
    - If resume fires while sendBody is blocking: sendBody gets unblocked
      (pass)
    Both race outcomes produce a passing test. Restore timeout=1000 (800ms
    of headroom is plenty).
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
 .../direct/DirectProducerBlockingTest.java         | 45 ++++++++--------------
 1 file changed, 17 insertions(+), 28 deletions(-)

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 160777ee4ecb..a692a4685b41 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.apache.camel.util.StopWatch;
 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,36 +73,26 @@ public class DirectProducerBlockingTest extends 
ContextTestSupport {
 
     @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).
-                    // Use a generous timeout — on slow CI the thread state
-                    // detection can take longer than 2 s.
-                    await().atMost(10, 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.
-        // Use a generous timeout so the background thread has enough headroom 
to
-        // detect the TIMED_WAITING state and resume the route even under CI 
load.
-        template.sendBody("direct:suspended?block=true&timeout=10000", "hello 
world");
+        template.sendBody("direct:suspended?block=true&timeout=1000", "hello 
world");
 
         assertMockEndpointsSatisfied();
 

Reply via email to