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 4db824c40226 CAMEL-25071: camel-management - Route and CamelContext
MBeans: fix the remaining follow-ups from the deep review (#27072)
4db824c40226 is described below
commit 4db824c40226c07ba4560c766be8ebb22c7ce4b8
Author: Claus Ibsen <[email protected]>
AuthorDate: Wed Sep 30 12:42:56 2026 +0200
CAMEL-25071: camel-management - Route and CamelContext MBeans: fix the
remaining follow-ups from the deep review (#27072)
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
---
.../engine/DefaultSupervisingRouteController.java | 3 +
...ervisingRouteControllerExhaustedRemoveTest.java | 101 +++++++++++++++++++++
.../management/CompositePerformanceCounter.java | 36 +++-----
.../management/mbean/ManagedCamelContext.java | 24 +++--
.../mbean/ManagedPerformanceCounter.java | 37 +++++++-
.../camel/management/mbean/ManagedProcessor.java | 20 ++++
.../camel/management/mbean/ManagedRoute.java | 9 ++
.../camel/management/ManagedRedeliverTest.java | 44 +++++++++
.../ManagedStatisticsEnabledInflightTest.java | 87 ++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 8 ++
10 files changed, 335 insertions(+), 34 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 d85a2e391bcb..7960bb32c55d 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..e0bf112a7c0f 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
@@ -36,6 +36,14 @@ public abstract class ManagedPerformanceCounter extends
ManagedCounter
public static final String TIMESTAMP_FORMAT = "yyyy-MM-dd'T'HH:mm:ss.SSSZ";
+ /**
+ * The route and processor id of the processor whose last attempt failed,
so a redelivery is only counted by the
+ * processor (and route) that is redelivered, as the redelivered header
stays on the exchange after a redelivery
+ * succeeds.
+ */
+ static final String FAILED_ROUTE_ID = "CamelManagementFailedRouteId";
+ static final String FAILED_PROCESSOR_ID =
"CamelManagementFailedProcessorId";
+
private static final int PERCENTILE_WINDOW_SIZE = 1024;
private Statistic exchangesCompleted;
@@ -304,9 +312,22 @@ public abstract class ManagedPerformanceCounter extends
ManagedCounter
thp.update(getExchangesTotal());
}
+ /**
+ * Whether the exchange was redelivered by this counter, which is when the
exchange is redelivered after a processor
+ * failed. A processor and a route override this to only count the
redelivery of their own processor.
+ */
+ protected boolean isRedeliveredHere(Exchange exchange) {
+ return ExchangeHelper.isRedelivered(exchange) &&
exchange.getProperty(FAILED_PROCESSOR_ID) != null;
+ }
+
@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 +336,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 (isRedeliveredHere(exchange)) {
+ redeliveries.increment();
+ }
if (exchange.isExternalRedelivered()) {
externalRedeliveries.increment();
}
@@ -363,11 +391,14 @@ 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)) {
+ if (isRedeliveredHere(exchange)) {
redeliveries.increment();
}
if (exchange.isExternalRedelivered()) {
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedProcessor.java
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedProcessor.java
index 04935145ff82..79a6d2b1338f 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedProcessor.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedProcessor.java
@@ -16,8 +16,11 @@
*/
package org.apache.camel.management.mbean;
+import java.util.Objects;
+
import org.apache.camel.CamelContext;
import org.apache.camel.DisabledAware;
+import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.Route;
import org.apache.camel.ServiceStatus;
@@ -33,6 +36,7 @@ import org.apache.camel.model.StepDefinition;
import org.apache.camel.spi.ManagementStrategy;
import org.apache.camel.spi.NodeIdFactory;
import org.apache.camel.spi.RouteIdAware;
+import org.apache.camel.support.ExchangeHelper;
import org.apache.camel.support.LoggerHelper;
import org.apache.camel.support.PluginHelper;
import org.apache.camel.support.service.ServiceHelper;
@@ -260,4 +264,20 @@ public class ManagedProcessor extends
ManagedPerformanceCounter
public String dumpProcessorAsXml() throws Exception {
return
PluginHelper.getModelToXMLDumper(context).dumpModelAsXml(context, definition);
}
+
+ @Override
+ protected boolean isRedeliveredHere(Exchange exchange) {
+ // only the processor that failed is redelivered (the later processors
see the redelivered header as well)
+ return ExchangeHelper.isRedelivered(exchange)
+ && id.equals(exchange.getProperty(FAILED_PROCESSOR_ID))
+ && Objects.equals(getRouteId(),
exchange.getProperty(FAILED_ROUTE_ID));
+ }
+
+ @Override
+ public void failedExchange(Exchange exchange) {
+ super.failedExchange(exchange);
+ // remember the processor that failed, as the error handler may
redeliver it
+ exchange.setProperty(FAILED_ROUTE_ID, getRouteId());
+ exchange.setProperty(FAILED_PROCESSOR_ID, id);
+ }
}
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
index 780ea82cea6a..85efb01eb4cc 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
@@ -43,6 +43,7 @@ import javax.management.openmbean.TabularData;
import javax.management.openmbean.TabularDataSupport;
import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
import org.apache.camel.ExtendedCamelContext;
import org.apache.camel.ManagementStatisticsLevel;
import org.apache.camel.Route;
@@ -61,6 +62,7 @@ import org.apache.camel.model.RoutesDefinition;
import org.apache.camel.spi.InflightRepository;
import org.apache.camel.spi.ManagementStrategy;
import org.apache.camel.spi.RoutePolicy;
+import org.apache.camel.support.ExchangeHelper;
import org.apache.camel.support.PluginHelper;
import org.apache.camel.util.ObjectHelper;
import org.apache.camel.util.StringHelper;
@@ -1043,4 +1045,11 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
return StringHelper.xmlEncode(text);
}
+ @Override
+ protected boolean isRedeliveredHere(Exchange exchange) {
+ // only the route of the processor that failed is redelivered (the
later routes see the redelivered header too)
+ return ExchangeHelper.isRedelivered(exchange)
+ && exchange.getProperty(FAILED_PROCESSOR_ID) != null
+ && route.getId().equals(exchange.getProperty(FAILED_ROUTE_ID));
+ }
}
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..72e046d02a1a 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,35 @@ 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"));
+
+ // the redelivered header stays on the exchange, but the later
processors and routes were not redelivered
+ for (String id : new String[] { "after-flaky", "call-other",
"in-other" }) {
+ on = getCamelObjectName(TYPE_PROCESSOR, id);
+ assertEquals(1L, mbeanServer.getAttribute(on,
"ExchangesCompleted"), id);
+ assertEquals(0L, mbeanServer.getAttribute(on, "Redeliveries"), id);
+ }
+ on = getCamelObjectName(TYPE_ROUTE, "other");
+ assertEquals(1L, mbeanServer.getAttribute(on, "ExchangesCompleted"));
+ assertEquals(0L, mbeanServer.getAttribute(on, "Redeliveries"));
+ }
+
@Override
protected RouteBuilder createRouteBuilder() {
return new RouteBuilder() {
@@ -88,6 +119,19 @@ 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")
+ .log("after the flaky processor").id("after-flaky")
+ .to("direct:other").id("call-other");
+
+ from("direct:other").routeId("other")
+ .log("in the other route").id("in-other");
+
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 3b55ad81d6ef..077eeda00e73 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
@@ -526,6 +526,14 @@ 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, and it is only counted where the
+redelivery happened. For a processor it is the number of redelivery attempts
of that processor. For a route it is the
+number of exchanges that were redelivered by a processor of that route, and
for the CamelContext the number of
+exchanges that were redelivered. 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), and a processor or route that the exchange went through
after a redelivery could count that
+redelivery as well.
+
=== camel-groovy
A `GroovyShellFactory` is now looked up in the registry once per
`CamelContext`, when the first groovy expression is