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 c86580a7a2d0 CAMEL-25055: camel-core - Error registry: fix bugs found
in a deep review (#26938)
c86580a7a2d0 is described below
commit c86580a7a2d096c440239cb90f9cfb6df23940d8
Author: Claus Ibsen <[email protected]>
AuthorDate: Tue Sep 29 08:11:06 2026 +0200
CAMEL-25055: camel-core - Error registry: fix bugs found in a deep review
(#26938)
- an error is only recorded as handled when the failure processor handled it
- only captures of the same failure (same exception, or one wrapping the
other) are merged,
so a second failure, the failures of split parts and a failure in
onCompletion are all recorded
- with noErrorHandler the route id is taken from the message history, as
the node is
- forRoute(id).clear() resets the repeat counts of the route
- the JMX browse table is indexed by the uid of the entry, as an exchange
can have more than one
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../camel/impl/engine/DefaultErrorRegistry.java | 72 +++++++--
.../camel/impl/ErrorRegistryEdgeCasesTest.java | 173 +++++++++++++++++++++
.../api/management/mbean/CamelOpenMBeanTypes.java | 10 +-
.../management/mbean/ManagedErrorRegistry.java | 3 +-
.../camel/management/ManagedErrorRegistryTest.java | 76 +++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 9 ++
6 files changed, 320 insertions(+), 23 deletions(-)
diff --git
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java
index 4ad90c829991..c47a53cdfcbe 100644
---
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java
+++
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java
@@ -110,6 +110,10 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
Throwable exception;
if (handled) {
exception =
exchange.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT, Throwable.class);
+ // the event means a failure processor (onException, dead letter
channel, doCatch) has run, which has
+ // only handled the failure when the exchange no longer has an
exception (not with handled(false) or a
+ // doCatch that throws again)
+ handled = exchange.getException() == null;
} else {
exception = exchange.getException();
}
@@ -125,19 +129,6 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
long timestamp = System.currentTimeMillis();
String exchangeId = correlationId != null ? correlationId :
exchange.getExchangeId();
String routeId =
exchange.getProperty(ExchangePropertyKey.FAILURE_ROUTE_ID, String.class);
- if (routeId == null) {
- routeId = exchange.getFromRouteId();
- }
- String fromRouteId = exchange.getFromRouteId();
- String routeGroup = null;
- if (routeId != null) {
- org.apache.camel.Route route =
exchange.getContext().getRoute(routeId);
- if (route != null) {
- routeGroup = route.getGroup();
- }
- }
- String endpointUri =
exchange.getProperty(ExchangePropertyKey.FAILURE_ENDPOINT, String.class);
-
// capture node id and location where the exchange actually failed
// (captured up-front by the error handler / doCatch, before any
failure processor such as
// onException or a dead letter channel ran its own processing steps -
otherwise those steps
@@ -155,8 +146,25 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
toNode = last.getNode().getId();
location =
LoggerHelper.getLineNumberLoggerName(last.getNode());
}
+ if (routeId == null) {
+ // the route of the node (which is not the route the
exchange came from when it failed in a
+ // route it was sent to)
+ routeId = last.getRouteId();
+ }
+ }
+ }
+ if (routeId == null) {
+ routeId = exchange.getFromRouteId();
+ }
+ String fromRouteId = exchange.getFromRouteId();
+ String routeGroup = null;
+ if (routeId != null) {
+ org.apache.camel.Route route =
exchange.getContext().getRoute(routeId);
+ if (route != null) {
+ routeGroup = route.getGroup();
}
}
+ String endpointUri =
exchange.getProperty(ExchangePropertyKey.FAILURE_ENDPOINT, String.class);
// capture step id (set by Step EIP)
String stepId = exchange.getProperty(ExchangePropertyKey.STEP_ID,
String.class);
@@ -196,16 +204,24 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
endpointUri, toNode, stepId, fromEndpointUri, routeUptime,
elapsed,
threadName, data, exception, handled, messageHistory);
- // deduplicate by exchange ID:
+ // deduplicate the same failure of an exchange (the same exception, or
one that wraps the other), as it is
+ // reported by both a correlated copy and the original exchange. Other
failures of the same exchange (the
+ // failures of the parts of a split, a failure after a doCatch, a
failure in onCompletion) are kept.
// - correlated copy (inner): has more specific node info (e.g.,
throwException inside circuit breaker),
- // so it replaces any existing entry for the same original exchange
+ // so it replaces an existing entry of the same failure
// - original exchange (outer): if already captured from a correlated
copy, skip it
// since the copy has more specific info about where the error
actually occurred
if (correlationId != null) {
- entries.removeIf(e -> exchangeId.equals(e.getExchangeId()));
+ entries.removeIf(e -> exchangeId.equals(e.getExchangeId()) &&
isSameFailure(e.getException(), exception));
+ entry.fromCopy = true;
} else {
for (BacklogErrorEventMessage e : entries) {
- if (exchangeId.equals(e.getExchangeId())) {
+ // the same exception again, or the failure of a correlated
copy wrapped (or unwrapped) by the
+ // original exchange. A doCatch of the original exchange that
throws a new exception wrapping the
+ // caught one is another failure, recorded besides the caught
one
+ if (exchangeId.equals(e.getExchangeId()) && (e.getException()
== exception
+ || e instanceof DefaultBacklogErrorEventMessage dbe &&
dbe.fromCopy
+ && isSameFailure(e.getException(),
exception))) {
// the copy's entry stays (it names the node), but the
original reporting the failure as
// handled (a circuit breaker's fallback, a doCatch around
a multicast) means the exchange
// recovered: the entry is an error that was handled, not
an error (CAMEL-24863)
@@ -227,6 +243,24 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
evict();
}
+ /**
+ * Whether the two exceptions are the same failure: the same exception, or
one is a cause of the other (such as the
+ * exception of a split part wrapped by the splitter).
+ */
+ private static boolean isSameFailure(Throwable a, Throwable b) {
+ return isCauseOf(a, b) || isCauseOf(b, a);
+ }
+
+ private static boolean isCauseOf(Throwable cause, Throwable exception) {
+ int depth = 0;
+ for (Throwable t = exception; t != null && depth < 20; t =
t.getCause(), depth++) {
+ if (t == cause) {
+ return true;
+ }
+ }
+ return false;
+ }
+
/**
* What makes two errors the same kind: the route, the node that failed,
and the type of the exception. The
* exception message is deliberately left out, because a real storm
usually carries the failing payload in its
@@ -529,6 +563,8 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
@Override
public void clear() {
entries.removeIf(entry -> routeId.equals(entry.getRouteId()));
+ // and the counts of the kinds of errors of this route (as clear
of the registry does)
+ repeats.keySet().removeIf(kind -> kind.startsWith(routeId + "|"));
}
}
@@ -554,6 +590,8 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
private final JsonObject data;
private final Throwable exception;
private volatile boolean handled;
+ // recorded from a correlated copy of the exchange (which the original
exchange reports again)
+ private volatile boolean fromCopy;
private volatile long repeatCount = 1;
private volatile long repeatFirstTimestamp;
private volatile long repeatLastTimestamp;
diff --git
a/core/camel-core/src/test/java/org/apache/camel/impl/ErrorRegistryEdgeCasesTest.java
b/core/camel-core/src/test/java/org/apache/camel/impl/ErrorRegistryEdgeCasesTest.java
new file mode 100644
index 000000000000..a34749f94330
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/impl/ErrorRegistryEdgeCasesTest.java
@@ -0,0 +1,173 @@
+/*
+ * 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;
+
+import java.io.IOException;
+import java.util.List;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.BacklogErrorEventMessage;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+public class ErrorRegistryEdgeCasesTest extends ContextTestSupport {
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext context = super.createCamelContext();
+ context.getErrorRegistry().setEnabled(true);
+ context.setMessageHistory(true);
+ return context;
+ }
+
+ private void send(String uri, Object body) {
+ try {
+ template.sendBody(uri, body);
+ } catch (Exception e) {
+ // the tests are about what the registry recorded
+ }
+ }
+
+ private List<BacklogErrorEventMessage> entries() {
+ return List.copyOf(context.getErrorRegistry().browse());
+ }
+
+ @Test
+ public void testOnExceptionNotHandledIsNotRecordedAsHandled() {
+ send("direct:notHandled", "a");
+ assertThat(entries()).singleElement().satisfies(e -> {
+
assertThat(e.getException()).isInstanceOf(IllegalArgumentException.class);
+ assertThat(e.isHandled()).isFalse();
+ });
+ }
+
+ @Test
+ public void testDoCatchThatThrowsAgainIsNotRecordedAsHandled() {
+ send("direct:rethrow", "a");
+
assertThat(entries()).extracting(BacklogErrorEventMessage::getExceptionType)
+
.containsExactlyInAnyOrder(IllegalStateException.class.getName(),
IllegalArgumentException.class.getName());
+ assertThat(entries()).allSatisfy(e ->
assertThat(e.isHandled()).isFalse());
+ }
+
+ @Test
+ public void testDoCatchThatThrowsWrappedExceptionIsRecorded() {
+ send("direct:wrap", "a");
+ assertThat(entries()).hasSize(2);
+ assertThat(entries()).anySatisfy(e -> {
+
assertThat(e.getException()).isInstanceOf(IllegalStateException.class).hasMessage("wrapped");
+ assertThat(e.isHandled()).isFalse();
+ });
+ assertThat(entries()).anySatisfy(e -> {
+
assertThat(e.getException()).isInstanceOf(IllegalArgumentException.class).hasMessage("inner");
+ assertThat(e.getToNode()).isEqualTo("w1");
+ assertThat(e.isHandled()).isFalse();
+ });
+ }
+
+ @Test
+ public void testFailureAfterHandledFailureIsRecorded() {
+ send("direct:second", "a");
+ assertThat(entries()).hasSize(2);
+ assertThat(entries()).anySatisfy(e -> {
+
assertThat(e.getException()).isInstanceOf(IllegalStateException.class);
+ assertThat(e.getToNode()).isEqualTo("t2");
+ assertThat(e.isHandled()).isFalse();
+ });
+ assertThat(entries()).anySatisfy(e -> {
+
assertThat(e.getException()).isInstanceOf(IllegalArgumentException.class);
+ assertThat(e.isHandled()).isTrue();
+ });
+ }
+
+ @Test
+ public void testFailuresOfSplitPartsAreAllRecorded() {
+ send("direct:split", "a,b,c");
+
assertThat(entries()).extracting(BacklogErrorEventMessage::getToNode).containsExactlyInAnyOrder("n1",
"n2");
+ }
+
+ @Test
+ public void testFailureInOnCompletionKeepsTheRouteFailure() {
+ send("direct:oc", "a");
+
assertThat(entries()).extracting(BacklogErrorEventMessage::getExceptionType)
+
.containsExactlyInAnyOrder(IllegalArgumentException.class.getName(),
IOException.class.getName());
+ }
+
+ @Test
+ public void testNoErrorHandlerRecordsTheRouteThatFailed() {
+ send("direct:a", "a");
+ assertThat(entries()).singleElement().satisfies(e -> {
+ assertThat(e.getRouteId()).isEqualTo("b");
+ assertThat(e.getToNode()).isEqualTo("boom");
+ });
+ }
+
+ @Test
+ public void testClearOfRouteResetsRepeatCount() {
+ for (int i = 0; i < 5; i++) {
+ send("direct:storm", "a");
+ }
+ context.getErrorRegistry().forRoute("storm").clear();
+ send("direct:storm", "a");
+ assertThat(entries()).singleElement().satisfies(e ->
assertThat(e.getRepeatCount()).isEqualTo(1));
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:notHandled").routeId("nh")
+
.onException(IllegalArgumentException.class).to("mock:log").end()
+ .throwException(new IllegalArgumentException("x"));
+ from("direct:rethrow").routeId("rt")
+ .doTry().throwException(new
IllegalArgumentException("x"))
+
.doCatch(IllegalArgumentException.class).throwException(new
IllegalStateException("wrapped"))
+ .end();
+ from("direct:wrap").routeId("wrap")
+ .doTry().throwException(new
IllegalArgumentException("inner")).id("w1")
+ .doCatch(IllegalArgumentException.class)
+ .process(e -> {
+ throw new IllegalStateException(
+ "wrapped",
e.getProperty(Exchange.EXCEPTION_CAUGHT, Exception.class));
+ })
+ .end();
+ from("direct:second").routeId("second")
+ .doTry().throwException(new
IllegalArgumentException("first")).id("t1")
+ .doCatch(IllegalArgumentException.class).end()
+ .throwException(new
IllegalStateException("second")).id("t2");
+ from("direct:split").routeId("split")
+ .errorHandler(deadLetterChannel("mock:dead"))
+ .split(body().tokenize(","))
+ .choice()
+ .when(body().isEqualTo("a")).throwException(new
IllegalArgumentException("a")).id("n1")
+ .when(body().isEqualTo("c")).throwException(new
NullPointerException("c")).id("n2")
+ .end();
+ from("direct:oc").routeId("oc")
+ .onCompletion().onFailureOnly().throwException(new
IOException("alert down")).end()
+ .throwException(new IllegalArgumentException("real"));
+
from("direct:a").routeId("a").errorHandler(noErrorHandler()).to("direct:b");
+ from("direct:b").routeId("b").errorHandler(noErrorHandler())
+ .throwException(new
IllegalArgumentException("boom")).id("boom");
+ from("direct:storm").routeId("storm").throwException(new
IllegalArgumentException("storm"));
+ }
+ };
+ }
+}
diff --git
a/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java
b/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java
index 9c1b3eb77d63..c1a686c01abf 100644
---
a/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java
+++
b/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java
@@ -384,25 +384,25 @@ public final class CamelOpenMBeanTypes {
public static TabularType listErrorRegistryTabularType() throws
OpenDataException {
CompositeType ct = listErrorRegistryCompositeType();
- return new TabularType("listErrors", "Lists captured routing errors",
ct, new String[] { "exchangeId" });
+ return new TabularType("listErrors", "Lists captured routing errors",
ct, new String[] { "uid" });
}
public static CompositeType listErrorRegistryCompositeType() throws
OpenDataException {
return new CompositeType(
"errors", "Errors",
new String[] {
- "exchangeId", "routeId", "routeGroup", "nodeId",
"stepId",
+ "uid", "exchangeId", "routeId", "routeGroup",
"nodeId", "stepId",
"endpointUri", "fromEndpointUri", "timestamp",
"routeUptime", "elapsed",
"handled", "exceptionType", "exceptionMessage" },
new String[] {
- "Exchange Id", "Route Id", "Route Group", "Node Id",
"Step Id",
+ "Uid", "Exchange Id", "Route Id", "Route Group", "Node
Id", "Step Id",
"Endpoint Uri", "From Endpoint Uri", "Timestamp",
"Route Uptime", "Elapsed",
"Handled", "Exception Type", "Exception Message" },
new OpenType[] {
- SimpleType.STRING, SimpleType.STRING,
SimpleType.STRING, SimpleType.STRING, SimpleType.STRING,
- SimpleType.STRING, SimpleType.STRING,
SimpleType.STRING,
+ SimpleType.LONG, SimpleType.STRING, SimpleType.STRING,
SimpleType.STRING, SimpleType.STRING,
+ SimpleType.STRING, SimpleType.STRING,
SimpleType.STRING, SimpleType.STRING,
SimpleType.LONG, SimpleType.LONG,
SimpleType.BOOLEAN, SimpleType.STRING,
SimpleType.STRING });
}
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java
index dc651f38470a..759d45c46f37 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java
@@ -170,11 +170,12 @@ public class ManagedErrorRegistry extends ManagedService
implements ManagedError
CompositeData data = new CompositeDataSupport(
ct,
new String[] {
- "exchangeId", "routeId", "routeGroup",
"nodeId", "stepId",
+ "uid", "exchangeId", "routeId", "routeGroup",
"nodeId", "stepId",
"endpointUri", "fromEndpointUri", "timestamp",
"routeUptime", "elapsed",
"handled", "exceptionType", "exceptionMessage"
},
new Object[] {
+ entry.getUid(),
entry.getExchangeId(),
entry.getRouteId(),
entry.getRouteGroup(),
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedErrorRegistryTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedErrorRegistryTest.java
new file mode 100644
index 000000000000..7d591e40d3e5
--- /dev/null
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedErrorRegistryTest.java
@@ -0,0 +1,76 @@
+/*
+ * 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.Set;
+
+import javax.management.MBeanServer;
+import javax.management.ObjectName;
+import javax.management.openmbean.TabularData;
+
+import org.apache.camel.CamelContext;
+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.junit.jupiter.api.Assertions.assertEquals;
+
+@DisabledOnOs(OS.AIX)
+public class ManagedErrorRegistryTest extends ManagementTestSupport {
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext context = super.createCamelContext();
+ context.getErrorRegistry().setEnabled(true);
+ return context;
+ }
+
+ @Test
+ public void testBrowseErrorsOfSameExchange() throws Exception {
+ getMockEndpoint("mock:dead").expectedMessageCount(2);
+ template.sendBody("direct:start", "a,b,c");
+ assertMockEndpointsSatisfied();
+
+ // the failures of the two parts of the split are recorded under the
same exchange id
+ assertEquals(2, context.getErrorRegistry().browse().size());
+
+ MBeanServer mbeanServer = getMBeanServer();
+ Set<ObjectName> names = mbeanServer.queryNames(
+ new ObjectName("org.apache.camel:context=" +
context.getManagementName() + ",*"), null);
+ ObjectName on = names.stream().filter(n -> n.getKeyProperty("name") !=
null
+ &&
n.getKeyProperty("name").contains("ErrorRegistry")).findFirst().orElseThrow();
+
+ TabularData data = (TabularData) mbeanServer.invoke(on, "browse",
null, null);
+ assertEquals(2, data.size());
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ errorHandler(deadLetterChannel("mock:dead"));
+
+ from("direct:start")
+ .split(body().tokenize(","))
+ .filter(body().isNotEqualTo("b"))
+ .throwException(IllegalArgumentException.class,
"Forced ${body}");
+ }
+ };
+ }
+}
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 5cd5607798de..9b1ea418f8a0 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
@@ -365,6 +365,15 @@ through the transformer. The defaults are unchanged when
the header is absent.
If a route sets `CamelAwsDdbReturnValues` before the transformer runs and
relies on it being discarded, remove the
header instead.
+=== Error registry
+
+- An error is only recorded as handled when the failure processor has handled
it. Prior to Camel 4.23 an error was
+ recorded as handled whenever an `onException`, dead letter channel or
`doCatch` ran for it, also with
+ `handled(false)` or when the `doCatch` threw again.
+- The error registry records every failure of an exchange: the failures of the
parts of a split or multicast, a failure
+ after an error that was handled by a `doCatch`, and a failure in
`onCompletion`. Prior to Camel 4.23 it kept only
+ one entry per exchange, so such failures replaced or hid each other.
+
=== camel-aws2-s3
The producer no longer evaluates the S3 object key or bucket name taken from a
message header as a Simple