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

davsclaus pushed a commit to branch fix/CAMEL-25071
in repository https://gitbox.apache.org/repos/asf/camel.git

commit e4c2dedab2cbafcb1122dedb8eb4a9f33d15808d
Author: Claus Ibsen <[email protected]>
AuthorDate: Tue Sep 29 11:53:49 2026 +0200

    CAMEL-25071: camel-management - Route and CamelContext MBeans: fix the 
remaining follow-ups from the deep review
    
    - Enabling or disabling StatisticsEnabled while exchanges are inflight no 
longer leaves the inflight count
      of the route, CamelContext and route group wrong. The inflight count is 
always kept, and the other
      statistics only when enabled.
    - Removing an exhausted route of the supervising route controller removes 
it from the restarting and
      exhausted routes, so it is no longer unhealthy and routeStatus no longer 
fails.
    - Redeliveries also counts a redelivery attempt that succeeds.
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../engine/DefaultSupervisingRouteController.java  |   3 +
 ...ervisingRouteControllerExhaustedRemoveTest.java | 101 +++++++++++++++++++++
 .../management/CompositePerformanceCounter.java    |  36 +++-----
 .../management/mbean/ManagedCamelContext.java      |  24 +++--
 .../mbean/ManagedPerformanceCounter.java           |  19 +++-
 .../camel/management/ManagedRedeliverTest.java     |  29 ++++++
 .../ManagedStatisticsEnabledInflightTest.java      |  87 ++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |   4 +
 8 files changed, 270 insertions(+), 33 deletions(-)

diff --git 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
index 6fcdc7a3a567..08f116d7b328 100644
--- 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
+++ 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
@@ -964,6 +964,9 @@ public class DefaultSupervisingRouteController extends 
DefaultRouteController im
             try {
                 routes.removeIf(
                         r -> ObjectHelper.equal(r.get(), route) || 
ObjectHelper.equal(r.getId(), route.getId()));
+                nonSupervisedRoutes.remove(route.getId());
+                // a removed route is no longer restarting or exhausted (and 
unhealthy)
+                routeManager.release(new RouteHolder(route, 0));
             } finally {
                 lock.unlock();
             }
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/impl/engine/DefaultSupervisingRouteControllerExhaustedRemoveTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/impl/engine/DefaultSupervisingRouteControllerExhaustedRemoveTest.java
new file mode 100644
index 000000000000..a09b9e931c97
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/impl/engine/DefaultSupervisingRouteControllerExhaustedRemoveTest.java
@@ -0,0 +1,101 @@
+/*
+ * 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.impl.engine;
+
+import java.time.Duration;
+import java.util.Map;
+
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.SupervisingRouteController;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultConsumer;
+import org.apache.camel.support.DefaultEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A route whose restarts were exhausted, and which is then removed, is no 
longer regarded as exhausted (or unhealthy).
+ */
+public class DefaultSupervisingRouteControllerExhaustedRemoveTest extends 
ContextTestSupport {
+
+    @Override
+    public boolean isUseRouteBuilder() {
+        return false;
+    }
+
+    @Test
+    public void testRemoveAfterExhausted() throws Exception {
+        context.addComponent("flaky", new DefaultComponent() {
+            @Override
+            protected Endpoint createEndpoint(String uri, String remaining, 
Map<String, Object> parameters) {
+                return new DefaultEndpoint(uri, this) {
+                    @Override
+                    public Producer createProducer() {
+                        throw new UnsupportedOperationException();
+                    }
+
+                    @Override
+                    public Consumer createConsumer(Processor processor) {
+                        return new DefaultConsumer(this, processor) {
+                            @Override
+                            protected void doStart() throws Exception {
+                                throw new IllegalStateException("Cannot 
connect");
+                            }
+                        };
+                    }
+                };
+            }
+        });
+        context.addRoutes(new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("flaky:start").routeId("flaky").to("mock:result");
+            }
+        });
+
+        SupervisingRouteController src = 
context.getRouteController().supervising();
+        src.setBackOffDelay(10);
+        src.setBackOffMaxAttempts(2);
+        src.setInitialDelay(10);
+        src.setUnhealthyOnExhausted(true);
+        context.start();
+
+        await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> 
assertEquals(1, src.getExhaustedRoutes().size()));
+        assertTrue(((DefaultSupervisingRouteController) 
src).hasUnhealthyRoutes());
+
+        // the route is given up and removed
+        
assertTrue(context.getRouteController().getRouteStatus("flaky").isStopped());
+        assertTrue(context.removeRoute("flaky"));
+
+        assertNull(context.getRoute("flaky"));
+        assertEquals(0, src.getExhaustedRoutes().size(), "the removed route 
should no longer be exhausted");
+        assertEquals(0, src.getControlledRoutes().size());
+        assertNull(src.getRestartException("flaky"));
+        assertFalse(((DefaultSupervisingRouteController) 
src).hasUnhealthyRoutes(),
+                "the removed route should no longer be unhealthy");
+    }
+}
diff --git 
a/core/camel-management/src/main/java/org/apache/camel/management/CompositePerformanceCounter.java
 
b/core/camel-management/src/main/java/org/apache/camel/management/CompositePerformanceCounter.java
index 78162c85d5d9..d6280e596afe 100644
--- 
a/core/camel-management/src/main/java/org/apache/camel/management/CompositePerformanceCounter.java
+++ 
b/core/camel-management/src/main/java/org/apache/camel/management/CompositePerformanceCounter.java
@@ -39,47 +39,37 @@ public class CompositePerformanceCounter implements 
PerformanceCounter {
 
     @Override
     public void processExchange(Exchange exchange, String type) {
-        if (counter1.isStatisticsEnabled()) {
-            counter1.processExchange(exchange, type);
-        }
-        if (counter2.isStatisticsEnabled()) {
-            counter2.processExchange(exchange, type);
-        }
-        if (counter3 != null && counter3.isStatisticsEnabled()) {
+        counter1.processExchange(exchange, type);
+        counter2.processExchange(exchange, type);
+        if (counter3 != null) {
             counter3.processExchange(exchange, type);
         }
     }
 
     @Override
     public void completedExchange(Exchange exchange, long time) {
-        if (counter1.isStatisticsEnabled()) {
-            counter1.completedExchange(exchange, time);
-        }
-        if (counter2.isStatisticsEnabled()) {
-            counter2.completedExchange(exchange, time);
-        }
-        if (counter3 != null && counter3.isStatisticsEnabled()) {
+        counter1.completedExchange(exchange, time);
+        counter2.completedExchange(exchange, time);
+        if (counter3 != null) {
             counter3.completedExchange(exchange, time);
         }
     }
 
     @Override
     public void failedExchange(Exchange exchange) {
-        if (counter1.isStatisticsEnabled()) {
-            counter1.failedExchange(exchange);
-        }
-        if (counter2.isStatisticsEnabled()) {
-            counter2.failedExchange(exchange);
-        }
-        if (counter3 != null && counter3.isStatisticsEnabled()) {
+        counter1.failedExchange(exchange);
+        counter2.failedExchange(exchange);
+        if (counter3 != null) {
             counter3.failedExchange(exchange);
         }
     }
 
     @Override
     public boolean isStatisticsEnabled() {
-        // this method is not used
-        return true;
+        // an exchange is counted when any of the counters has statistics 
enabled; each counter keeps its inflight
+        // count, and only gathers the other statistics when it is enabled 
itself
+        return counter1.isStatisticsEnabled() || counter2.isStatisticsEnabled()
+                || counter3 != null && counter3.isStatisticsEnabled();
     }
 
     @Override
diff --git 
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
 
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
index c6ca85c46fc1..7d561b61fcc2 100644
--- 
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
+++ 
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
@@ -115,17 +115,21 @@ public class ManagedCamelContext extends 
ManagedPerformanceCounter implements Ma
             if (level <= 1) {
                 super.completedExchange(exchange, time);
                 if (exchange.getFromEndpoint() != null && 
exchange.getFromEndpoint().isRemote()) {
-                    remoteExchangesTotal.increment();
-                    remoteExchangesCompleted.increment();
                     remoteExchangesInflight.decrement();
+                    if (isStatisticsEnabled()) {
+                        remoteExchangesTotal.increment();
+                        remoteExchangesCompleted.increment();
+                    }
                 }
             }
         } else {
             super.completedExchange(exchange, time);
             if (exchange.getFromEndpoint() != null && 
exchange.getFromEndpoint().isRemote()) {
-                remoteExchangesTotal.increment();
-                remoteExchangesCompleted.increment();
                 remoteExchangesInflight.decrement();
+                if (isStatisticsEnabled()) {
+                    remoteExchangesTotal.increment();
+                    remoteExchangesCompleted.increment();
+                }
             }
         }
     }
@@ -142,17 +146,21 @@ public class ManagedCamelContext extends 
ManagedPerformanceCounter implements Ma
             if (level <= 1) {
                 super.failedExchange(exchange);
                 if (exchange.getFromEndpoint() != null && 
exchange.getFromEndpoint().isRemote()) {
-                    remoteExchangesTotal.increment();
-                    remoteExchangesFailed.increment();
                     remoteExchangesInflight.decrement();
+                    if (isStatisticsEnabled()) {
+                        remoteExchangesTotal.increment();
+                        remoteExchangesFailed.increment();
+                    }
                 }
             }
         } else {
             super.failedExchange(exchange);
             if (exchange.getFromEndpoint() != null && 
exchange.getFromEndpoint().isRemote()) {
-                remoteExchangesTotal.increment();
-                remoteExchangesFailed.increment();
                 remoteExchangesInflight.decrement();
+                if (isStatisticsEnabled()) {
+                    remoteExchangesTotal.increment();
+                    remoteExchangesFailed.increment();
+                }
             }
         }
     }
diff --git 
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedPerformanceCounter.java
 
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedPerformanceCounter.java
index 40675eba41d5..87e8ee8ed5fc 100644
--- 
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedPerformanceCounter.java
+++ 
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedPerformanceCounter.java
@@ -306,7 +306,12 @@ public abstract class ManagedPerformanceCounter extends 
ManagedCounter
 
     @Override
     public void processExchange(Exchange exchange, String type) {
+        // the inflight count is kept also when statistics is disabled, so it 
stays correct when statistics
+        // is enabled or disabled while exchanges are inflight
         exchangesInflight.increment();
+        if (!statisticsEnabled) {
+            return;
+        }
         if ("route".equals(type)) {
             long now = System.currentTimeMillis();
             lastExchangeCreatedTimestamp.updateValue(now);
@@ -315,14 +320,21 @@ public abstract class ManagedPerformanceCounter extends 
ManagedCounter
 
     @Override
     public void completedExchange(Exchange exchange, long time) {
+        exchangesInflight.decrement();
+        if (!statisticsEnabled) {
+            return;
+        }
         increment();
         exchangesCompleted.increment();
-        exchangesInflight.decrement();
 
         if (ExchangeHelper.isFailureHandled(exchange)) {
             failuresHandled.increment();
             
lastExchangeFailureHandledTimestamp.updateValue(System.currentTimeMillis());
         }
+        // a redelivery attempt that succeeds is also a redelivery
+        if (ExchangeHelper.isRedelivered(exchange)) {
+            redeliveries.increment();
+        }
         if (exchange.isExternalRedelivered()) {
             externalRedeliveries.increment();
         }
@@ -363,9 +375,12 @@ public abstract class ManagedPerformanceCounter extends 
ManagedCounter
 
     @Override
     public void failedExchange(Exchange exchange) {
+        exchangesInflight.decrement();
+        if (!statisticsEnabled) {
+            return;
+        }
         increment();
         exchangesFailed.increment();
-        exchangesInflight.decrement();
 
         if (ExchangeHelper.isRedelivered(exchange)) {
             redeliveries.increment();
diff --git 
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedRedeliverTest.java
 
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRedeliverTest.java
index 4b134e6b5db4..45a01c85cf7c 100644
--- 
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedRedeliverTest.java
+++ 
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRedeliverTest.java
@@ -16,6 +16,8 @@
  */
 package org.apache.camel.management;
 
+import java.util.concurrent.atomic.AtomicInteger;
+
 import javax.management.MBeanServer;
 import javax.management.ObjectName;
 
@@ -78,6 +80,25 @@ public class ManagedRedeliverTest extends 
ManagementTestSupport {
         assertEquals(mock.getReceivedExchanges().get(0).getExchangeId(), last);
     }
 
+    @Test
+    public void testRedeliverSucceeds() throws Exception {
+        MBeanServer mbeanServer = getMBeanServer();
+
+        Object out = template.requestBody("direct:flaky", "Hello World");
+        assertEquals("Hello World", out);
+
+        // the processor failed twice, and the second redelivery succeeded
+        ObjectName on = getCamelObjectName(TYPE_PROCESSOR, "flaky-processor");
+        assertEquals(2L, mbeanServer.getAttribute(on, "ExchangesFailed"));
+        assertEquals(1L, mbeanServer.getAttribute(on, "ExchangesCompleted"));
+        assertEquals(2L, mbeanServer.getAttribute(on, "Redeliveries"));
+
+        // the exchange was redelivered and completed
+        on = getCamelObjectName(TYPE_ROUTE, "flaky");
+        assertEquals(1L, mbeanServer.getAttribute(on, "ExchangesCompleted"));
+        assertEquals(1L, mbeanServer.getAttribute(on, "Redeliveries"));
+    }
+
     @Override
     protected RouteBuilder createRouteBuilder() {
         return new RouteBuilder() {
@@ -88,6 +109,14 @@ public class ManagedRedeliverTest extends 
ManagementTestSupport {
                         .maximumRedeliveries(4).logStackTrace(false)
                         .setBody().constant("Error");
 
+                AtomicInteger attempts = new AtomicInteger();
+                from("direct:flaky").routeId("flaky")
+                        .process(exchange -> {
+                            if (attempts.incrementAndGet() < 3) {
+                                throw new IllegalArgumentException("Forced");
+                            }
+                        }).id("flaky-processor");
+
                 from("direct:start")
                         .to("mock:foo")
                         .process(exchange -> {
diff --git 
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedStatisticsEnabledInflightTest.java
 
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedStatisticsEnabledInflightTest.java
new file mode 100644
index 000000000000..909a76f47b9d
--- /dev/null
+++ 
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedStatisticsEnabledInflightTest.java
@@ -0,0 +1,87 @@
+/*
+ * 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.management;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+import javax.management.Attribute;
+import javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import org.apache.camel.builder.RouteBuilder;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.DisabledOnOs;
+import org.junit.jupiter.api.condition.OS;
+
+import static 
org.apache.camel.management.DefaultManagementObjectNameStrategy.TYPE_ROUTE;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * Enabling or disabling statistics while an exchange is inflight must not 
leave the inflight count wrong.
+ */
+@DisabledOnOs(OS.AIX)
+public class ManagedStatisticsEnabledInflightTest extends 
ManagementTestSupport {
+
+    private volatile CountDownLatch latch;
+
+    @Test
+    public void testDisableWhileInflight() throws Exception {
+        assertInflightAfterToggle(true, false);
+    }
+
+    @Test
+    public void testEnableWhileInflight() throws Exception {
+        assertInflightAfterToggle(false, true);
+    }
+
+    private void assertInflightAfterToggle(boolean before, boolean after) 
throws Exception {
+        MBeanServer mbeanServer = getMBeanServer();
+        ObjectName route = getCamelObjectName(TYPE_ROUTE, "foo");
+        ObjectName camelContext = getContextObjectName();
+
+        mbeanServer.setAttribute(route, new Attribute("StatisticsEnabled", 
before));
+        mbeanServer.setAttribute(camelContext, new 
Attribute("StatisticsEnabled", before));
+
+        latch = new CountDownLatch(1);
+        getMockEndpoint("mock:result").expectedMessageCount(1);
+        template.asyncSendBody("direct:start", "Hello World");
+        await().atMost(10, TimeUnit.SECONDS).until(() -> 
context.getInflightRepository().size("foo") == 1);
+
+        mbeanServer.setAttribute(route, new Attribute("StatisticsEnabled", 
after));
+        mbeanServer.setAttribute(camelContext, new 
Attribute("StatisticsEnabled", after));
+        latch.countDown();
+        assertMockEndpointsSatisfied();
+        await().atMost(10, TimeUnit.SECONDS).until(() -> 
context.getInflightRepository().size() == 0);
+
+        assertEquals(0L, mbeanServer.getAttribute(route, "ExchangesInflight"));
+        assertEquals(0L, mbeanServer.getAttribute(camelContext, 
"ExchangesInflight"));
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:start").routeId("foo")
+                        .process(e -> latch.await(20, TimeUnit.SECONDS))
+                        .to("mock:result");
+            }
+        };
+    }
+}
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 04d6027f2ded..980bb8b91c76 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -452,6 +452,10 @@ should be reviewed.
 `InflightRepository.InflightExchange` and 
`AsyncProcessorAwaitManager.AwaitThread` gained a `getNodeSource()` method
 for the same value. Both are `default` methods returning `null`, so existing 
implementations continue to compile.
 
+The `Redeliveries` statistic now also counts a redelivery attempt that 
succeeds. Before, only redelivery attempts
+that failed were counted, so a processor that succeeded on its second 
redelivery reported 1 instead of 2, and a route
+whose exchange was redelivered and then completed reported 0.
+
 === camel-groovy
 
 A `GroovyShellFactory` is now looked up in the registry once per 
`CamelContext`, when the first groovy expression is

Reply via email to