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

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


The following commit(s) were added to refs/heads/main by this push:
     new d1ac9fc7218b CAMEL-24763: camel-reactor - Fix flaky 
ReactorStreamsServiceTest (#27012)
d1ac9fc7218b is described below

commit d1ac9fc7218b92216759297bf1a1ad789a3f0255
Author: Urmila Unni <[email protected]>
AuthorDate: Tue Sep 29 11:12:34 2026 +0530

    CAMEL-24763: camel-reactor - Fix flaky ReactorStreamsServiceTest (#27012)
    
    Wait for the stream to complete instead of a 2 second latch.
    
    Co-authored-by: Claude <[email protected]>
---
 .../reactor/engine/ReactorStreamsServiceTest.java  | 54 ++++++++--------------
 1 file changed, 18 insertions(+), 36 deletions(-)

diff --git 
a/components/camel-reactor/src/test/java/org/apache/camel/component/reactor/engine/ReactorStreamsServiceTest.java
 
b/components/camel-reactor/src/test/java/org/apache/camel/component/reactor/engine/ReactorStreamsServiceTest.java
index 5f73da236f6e..ec356e014f52 100644
--- 
a/components/camel-reactor/src/test/java/org/apache/camel/component/reactor/engine/ReactorStreamsServiceTest.java
+++ 
b/components/camel-reactor/src/test/java/org/apache/camel/component/reactor/engine/ReactorStreamsServiceTest.java
@@ -16,9 +16,9 @@
  */
 package org.apache.camel.component.reactor.engine;
 
+import java.time.Duration;
 import java.util.Arrays;
-import java.util.Collections;
-import java.util.Set;
+import java.util.List;
 import java.util.TreeSet;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.TimeUnit;
@@ -283,74 +283,56 @@ public class ReactorStreamsServiceTest extends 
ReactorStreamsServiceTestSupport
     public void testTo() throws Exception {
         context.start();
 
-        Set<String> values = Collections.synchronizedSet(new TreeSet<>());
-        CountDownLatch latch = new CountDownLatch(3);
-
-        Flux.just(1, 2, 3)
+        List<String> values = Flux.just(1, 2, 3)
                 .flatMap(e -> crs.to("bean:hello", e, String.class))
-                .doOnNext(values::add)
-                .doOnNext(res -> latch.countDown())
-                .subscribe();
+                .collectList()
+                .block(Duration.ofSeconds(30));
 
-        assertTrue(latch.await(2, TimeUnit.SECONDS));
-        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), values);
+        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), new TreeSet<>(values));
     }
 
     @Test
     public void testToWithExchange() throws Exception {
         context.start();
 
-        Set<String> values = Collections.synchronizedSet(new TreeSet<>());
-        CountDownLatch latch = new CountDownLatch(3);
-
-        Flux.just(1, 2, 3)
+        List<String> values = Flux.just(1, 2, 3)
                 .flatMap(e -> crs.to("bean:hello", e))
                 .map(Exchange::getMessage)
                 .map(e -> e.getBody(String.class))
-                .doOnNext(values::add)
-                .doOnNext(res -> latch.countDown())
-                .subscribe();
+                .collectList()
+                .block(Duration.ofSeconds(30));
 
-        assertTrue(latch.await(2, TimeUnit.SECONDS));
-        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), values);
+        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), new TreeSet<>(values));
     }
 
     @Test
     public void testToFunction() throws Exception {
         context.start();
 
-        Set<String> values = Collections.synchronizedSet(new TreeSet<>());
-        CountDownLatch latch = new CountDownLatch(3);
         Function<Object, Publisher<String>> fun = crs.to("bean:hello", 
String.class);
 
-        Flux.just(1, 2, 3)
+        List<String> values = Flux.just(1, 2, 3)
                 .flatMap(fun)
-                .doOnNext(values::add)
-                .doOnNext(res -> latch.countDown())
-                .subscribe();
+                .collectList()
+                .block(Duration.ofSeconds(30));
 
-        assertTrue(latch.await(2, TimeUnit.SECONDS));
-        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), values);
+        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), new TreeSet<>(values));
     }
 
     @Test
     public void testToFunctionWithExchange() throws Exception {
         context.start();
 
-        Set<String> values = Collections.synchronizedSet(new TreeSet<>());
-        CountDownLatch latch = new CountDownLatch(3);
         Function<Object, Publisher<Exchange>> fun = crs.to("bean:hello");
 
-        Flux.just(1, 2, 3)
+        List<String> values = Flux.just(1, 2, 3)
                 .flatMap(fun)
                 .map(Exchange::getMessage)
                 .map(e -> e.getBody(String.class))
-                .doOnNext(values::add)
-                .doOnNext(res -> latch.countDown())
-                .subscribe();
+                .collectList()
+                .block(Duration.ofSeconds(30));
 
-        assertTrue(latch.await(2, TimeUnit.SECONDS));
-        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), values);
+        assertEquals(new TreeSet<>(Arrays.asList("Hello 1", "Hello 2", "Hello 
3")), new TreeSet<>(values));
     }
 
     // ************************************************

Reply via email to