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) {