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

gnodet 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 31dc50538ad7 CAMEL-24265: Use AtomicLong for 
DefaultTracer.traceCounter (#25160)
31dc50538ad7 is described below

commit 31dc50538ad79c57a6ff75edcdf0e7f4392f77de
Author: Guillaume Nodet <[email protected]>
AuthorDate: Mon Jul 27 23:46:37 2026 +0200

    CAMEL-24265: Use AtomicLong for DefaultTracer.traceCounter (#25160)
    
    The traceCounter field was a plain long incremented with ++ in
    traceBeforeNode(), which is called from concurrent routing threads.
    This is a compound read-modify-write that loses updates under
    contention. Switch to AtomicLong, consistent with BacklogTracer
    which already uses AtomicLong for the same purpose.
    
    Co-authored-by: Claude Opus 4.6 <[email protected]>
---
 .../apache/camel/impl/engine/DefaultTracer.java    |   9 +-
 .../DefaultTracerTraceCounterConcurrencyTest.java  | 113 +++++++++++++++++++++
 2 files changed, 118 insertions(+), 4 deletions(-)

diff --git 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultTracer.java
 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultTracer.java
index f01d9af6b51a..a3a7b036f7da 100644
--- 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultTracer.java
+++ 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultTracer.java
@@ -19,6 +19,7 @@ package org.apache.camel.impl.engine;
 import java.util.Map;
 import java.util.Objects;
 import java.util.StringJoiner;
+import java.util.concurrent.atomic.AtomicLong;
 
 import org.apache.camel.CamelContext;
 import org.apache.camel.CamelContextAware;
@@ -59,7 +60,7 @@ public class DefaultTracer extends ServiceSupport implements 
CamelContextAware,
     private volatile boolean standby;
     private volatile boolean traceRests;
     private volatile boolean traceTemplates;
-    private long traceCounter;
+    private final AtomicLong traceCounter = new AtomicLong();
 
     // immutable holder for compound tracePattern+patterns that must be 
visible atomically
     private record TracePatternHolder(String tracePattern, String[] patterns) {
@@ -92,7 +93,7 @@ public class DefaultTracer extends ServiceSupport implements 
CamelContextAware,
     @Override
     public void traceBeforeNode(NamedNode node, Exchange exchange) {
         if (shouldTrace(node)) {
-            traceCounter++;
+            traceCounter.incrementAndGet();
             String routeId = 
ExpressionBuilder.routeIdExpression().evaluate(exchange, String.class);
 
             // we need to avoid leak the sensible information here
@@ -258,12 +259,12 @@ public class DefaultTracer extends ServiceSupport 
implements CamelContextAware,
 
     @Override
     public long getTraceCounter() {
-        return traceCounter;
+        return traceCounter.get();
     }
 
     @Override
     public void resetTraceCounter() {
-        traceCounter = 0;
+        traceCounter.set(0);
     }
 
     @Override
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/DefaultTracerTraceCounterConcurrencyTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/DefaultTracerTraceCounterConcurrencyTest.java
new file mode 100644
index 000000000000..dbe9a3b9ecf7
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/DefaultTracerTraceCounterConcurrencyTest.java
@@ -0,0 +1,113 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.processor;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.Tracer;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * Verifies that {@link 
org.apache.camel.impl.engine.DefaultTracer#traceCounter} is thread-safe under 
concurrent
+ * routing. Before the fix (CAMEL-24265) the counter was a plain {@code long} 
incremented with {@code ++}, which is a
+ * compound read-modify-write — a classic lost-update race under concurrency. 
The fix switches to
+ * {@link java.util.concurrent.atomic.AtomicLong}.
+ */
+class DefaultTracerTraceCounterConcurrencyTest extends ContextTestSupport {
+
+    private static final int THREADS = 8;
+    private static final int MESSAGES_PER_THREAD = 250;
+    private static final int TOTAL_MESSAGES = THREADS * MESSAGES_PER_THREAD;
+
+    @Override
+    protected CamelContext createCamelContext() throws Exception {
+        CamelContext ctx = super.createCamelContext();
+        ctx.setTracing(true);
+        return ctx;
+    }
+
+    @Test
+    void traceCounterShouldBeAccurateUnderConcurrentRouting() throws Exception 
{
+        // Each message traverses one traced node ("mock:result"), so we expect
+        // the trace counter to equal the total number of messages sent.
+        getMockEndpoint("mock:result").expectedMessageCount(TOTAL_MESSAGES);
+
+        CountDownLatch startGate = new CountDownLatch(1);
+        CountDownLatch doneLatch = new CountDownLatch(THREADS);
+        ExecutorService pool = Executors.newFixedThreadPool(THREADS);
+        try {
+            for (int t = 0; t < THREADS; t++) {
+                pool.submit(() -> {
+                    try {
+                        startGate.await();
+                        for (int i = 0; i < MESSAGES_PER_THREAD; i++) {
+                            template.sendBody("direct:start", "msg");
+                        }
+                    } catch (Exception e) {
+                        // let the assertion on the mock catch failures
+                    } finally {
+                        doneLatch.countDown();
+                    }
+                });
+            }
+            // release all threads at once to maximise contention
+            startGate.countDown();
+            doneLatch.await();
+        } finally {
+            pool.shutdown();
+        }
+
+        assertMockEndpointsSatisfied();
+
+        Tracer tracer = context.getTracer();
+        assertThat(tracer.getTraceCounter())
+                .as("Trace counter must equal total messages when using 
AtomicLong")
+                .isEqualTo(TOTAL_MESSAGES);
+    }
+
+    @Test
+    void resetTraceCounterShouldClearCount() throws Exception {
+        getMockEndpoint("mock:result").expectedMessageCount(1);
+
+        template.sendBody("direct:start", "Hello");
+
+        assertMockEndpointsSatisfied();
+
+        Tracer tracer = context.getTracer();
+        assertThat(tracer.getTraceCounter()).isGreaterThan(0);
+
+        tracer.resetTraceCounter();
+        assertThat(tracer.getTraceCounter()).isZero();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:start").to("mock:result");
+            }
+        };
+    }
+}

Reply via email to