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 5f1974682bfe CAMEL-24038: Fix flaky AsyncWiretapTest - use thread-safe 
collections in mock telemetry infrastructure (#27226)
5f1974682bfe is described below

commit 5f1974682bfecca0744d3357aa29cc7744220be4
Author: Guillaume Nodet <[email protected]>
AuthorDate: Mon Oct 5 10:40:40 2026 +0200

    CAMEL-24038: Fix flaky AsyncWiretapTest - use thread-safe collections in 
mock telemetry infrastructure (#27226)
    
    The AsyncWiretapTest in camel-telemetry and camel-telemetry-dev was flaky
    due to thread-safety violations in the mock infrastructure used by tests.
    Wiretap creates async exchanges processed on different threads, leading to
    concurrent access on non-thread-safe collections.
    
    Root causes fixed:
    - MockTracer.MockSpanLifecycleManager.inMemoryStorageMap: plain HashMap
      accessed concurrently from exchange threads (activate/close) and the test
      thread (traces()). Replaced with ConcurrentHashMap to prevent silent entry
      loss, which directly caused the 'expected: <7> but was: <6>' failures.
    - MockSpanAdapter.tags: plain HashMap whose isDone tag is written from an
      exchange thread and read from the test thread without memory barriers.
      Replaced with ConcurrentHashMap.
    - MockSpanAdapter.logEntries: plain ArrayList mutated from exchange threads
      and read from the test thread. Replaced with a synchronized list.
    - DevSpanAdapter.tags: same pattern as MockSpanAdapter.tags. Replaced with
      ConcurrentHashMap.
    - DevSpanAdapter.logEntries: same pattern as MockSpanAdapter.logEntries.
      Replaced with a synchronized list.
    
    Also guard setComponent() and setTag() against null values since
    ConcurrentHashMap rejects null keys/values (unlike HashMap), and some
    decorator tests pass null values via Mockito mocks.
    
    InMemoryCollector is unaffected: its HashMap access is already protected by
    a ReentrantLock.
    
    Co-authored-by: Claude Sonnet 4.5 <[email protected]>
---
 .../apache/camel/telemetrydev/DevSpanAdapter.java  | 23 ++++++++++++++++------
 .../camel/telemetry/mock/MockSpanAdapter.java      | 22 +++++++++++++++------
 .../apache/camel/telemetry/mock/MockTracer.java    |  8 ++++++--
 3 files changed, 39 insertions(+), 14 deletions(-)

diff --git 
a/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/DevSpanAdapter.java
 
b/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/DevSpanAdapter.java
index 8f37cf8a035c..3e3566d3655d 100644
--- 
a/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/DevSpanAdapter.java
+++ 
b/components/camel-telemetry-dev/src/main/java/org/apache/camel/telemetrydev/DevSpanAdapter.java
@@ -17,10 +17,12 @@
 package org.apache.camel.telemetrydev;
 
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
+import java.util.concurrent.ConcurrentHashMap;
 
 import com.fasterxml.jackson.annotation.JsonAnyGetter;
 import com.fasterxml.jackson.annotation.JsonAnySetter;
@@ -29,8 +31,11 @@ import org.apache.camel.telemetry.TagConstants;
 
 public class DevSpanAdapter implements Span {
 
-    private List<LogEntry> logEntries = new ArrayList<>();
-    private final Map<String, String> tags = new HashMap<>();
+    // ConcurrentHashMap: tags (including isDone) are written from exchange 
threads and read
+    // from the collector thread without explicit synchronization.
+    private final Map<String, String> tags = new ConcurrentHashMap<>();
+    // synchronizedList: log() is called from exchange threads; 
getLogEntries() copies under the lock.
+    private List<LogEntry> logEntries = Collections.synchronizedList(new 
ArrayList<>());
 
     public static long nowMicros() {
         return System.currentTimeMillis() * 1000;
@@ -47,7 +52,9 @@ public class DevSpanAdapter implements Span {
 
     @Override
     public void setComponent(String component) {
-        this.tags.put(TagConstants.COMPONENT, component);
+        if (component != null) {
+            this.tags.put(TagConstants.COMPONENT, component);
+        }
     }
 
     @Override
@@ -58,7 +65,9 @@ public class DevSpanAdapter implements Span {
     @JsonAnySetter
     @Override
     public void setTag(String key, String value) {
-        this.tags.put(key, value);
+        if (key != null && value != null) {
+            this.tags.put(key, value);
+        }
     }
 
     public String getTag(String key) {
@@ -71,11 +80,13 @@ public class DevSpanAdapter implements Span {
     }
 
     public List<LogEntry> getLogEntries() {
-        return new ArrayList<>(this.logEntries);
+        synchronized (logEntries) {
+            return new ArrayList<>(this.logEntries);
+        }
     }
 
     public void setLogEntries(List<LogEntry> logEntries) {
-        this.logEntries = logEntries;
+        this.logEntries = Collections.synchronizedList(new 
ArrayList<>(logEntries));
     }
 
     public static final class LogEntry {
diff --git 
a/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockSpanAdapter.java
 
b/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockSpanAdapter.java
index 0b2d92cbcc9e..d7d1bfc7edb2 100644
--- 
a/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockSpanAdapter.java
+++ 
b/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockSpanAdapter.java
@@ -17,18 +17,22 @@
 package org.apache.camel.telemetry.mock;
 
 import java.util.ArrayList;
-import java.util.HashMap;
+import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
+import java.util.concurrent.ConcurrentHashMap;
 
 import org.apache.camel.telemetry.Span;
 import org.apache.camel.telemetry.TagConstants;
 
 public class MockSpanAdapter implements Span {
 
-    private final List<LogEntry> logEntries = new ArrayList<>();
-    private final Map<String, String> tags = new HashMap<>();
+    // ConcurrentHashMap: tags (including isDone) are written from exchange 
threads and read
+    // from the test thread without explicit synchronization.
+    private final Map<String, String> tags = new ConcurrentHashMap<>();
+    // synchronizedList: log() is called from exchange threads; logEntries() 
copies under the lock.
+    private final List<LogEntry> logEntries = Collections.synchronizedList(new 
ArrayList<>());
 
     public static long nowMicros() {
         return System.currentTimeMillis() * 1000;
@@ -44,7 +48,9 @@ public class MockSpanAdapter implements Span {
 
     @Override
     public void setComponent(String component) {
-        this.tags.put(TagConstants.COMPONENT, component);
+        if (component != null) {
+            this.tags.put(TagConstants.COMPONENT, component);
+        }
     }
 
     @Override
@@ -54,7 +60,9 @@ public class MockSpanAdapter implements Span {
 
     @Override
     public void setTag(String key, String value) {
-        this.tags.put(key, value);
+        if (key != null && value != null) {
+            this.tags.put(key, value);
+        }
     }
 
     public String getTag(String key) {
@@ -67,7 +75,9 @@ public class MockSpanAdapter implements Span {
     }
 
     public List<LogEntry> logEntries() {
-        return new ArrayList<>(this.logEntries);
+        synchronized (logEntries) {
+            return new ArrayList<>(this.logEntries);
+        }
     }
 
     public static final class LogEntry {
diff --git 
a/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockTracer.java
 
b/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockTracer.java
index 6479e8ebf213..303fedfed926 100644
--- 
a/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockTracer.java
+++ 
b/components/camel-telemetry/src/test/java/org/apache/camel/telemetry/mock/MockTracer.java
@@ -20,6 +20,7 @@ package org.apache.camel.telemetry.mock;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
 
 import org.apache.camel.api.management.ManagedResource;
 import org.apache.camel.spi.Configurer;
@@ -49,8 +50,11 @@ public class MockTracer extends Tracer {
 
     private class MockSpanLifecycleManager implements SpanLifecycleManager {
 
-        // Used to collect the traces for later evaluation as traces
-        Map<String, Span> inMemoryStorageMap = new HashMap<>();
+        // Used to collect the traces for later evaluation as traces.
+        // ConcurrentHashMap is required because wiretap creates async 
exchanges processed on
+        // different threads: activate() and close() are called from exchange 
threads while
+        // traces() is called from the test thread.
+        Map<String, Span> inMemoryStorageMap = new ConcurrentHashMap<>();
 
         @Override
         public Span create(String spanName, String spanKind, Span parentSpan, 
SpanContextPropagationExtractor extractor) {

Reply via email to