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 e7d5f2c4d74f CAMEL-25115: camel-core - Error handler failure paths:
fix bugs found in a deep review (#27021)
e7d5f2c4d74f is described below
commit e7d5f2c4d74f036f2e64d08d74f7d1766bd34849
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 22:45:13 2026 +0200
CAMEL-25115: camel-core - Error handler failure paths: fix bugs found in a
deep review (#27021)
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../camel/catalog/docs/dead-letter-channel.adoc | 9 +
.../camel/impl/engine/DefaultErrorRegistry.java | 7 +-
.../modules/eips/pages/dead-letter-channel.adoc | 9 +
.../camel/processor/FatalFallbackErrorHandler.java | 49 +++-
.../camel/processor/OnCompletionProcessor.java | 25 +-
.../DefaultExceptionPolicyStrategy.java | 10 +-
.../errorhandler/RedeliveryErrorHandler.java | 128 +++++++++-
.../onexception/ErrorHandlerFailurePathsTest.java | 277 +++++++++++++++++++++
.../OnExceptionHandledThrowsExceptionTest.java | 9 +-
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 29 +++
10 files changed, 525 insertions(+), 27 deletions(-)
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/dead-letter-channel.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/dead-letter-channel.adoc
index 37d1a28e2fbc..93e755f663a7 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/dead-letter-channel.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/dead-letter-channel.adoc
@@ -414,6 +414,15 @@ YAML::
====
+If the `onPrepareFailure` processor throws an exception, then the Exchange is
*not* sent to the dead letter queue,
+and the exception is regarded as a new exception that occurred during the dead
letter channel, which is handled
+according to the `deadLetterHandleNewException` option:
+
+* `deadLetterHandleNewException=true` (default): the new exception is logged
at `WARN` level and handled, and the
+Exchange completes without an exception.
+* `deadLetterHandleNewException=false`: the Exchange fails with the new
exception, which has the original caught
+exception attached as a suppressed exception.
+
=== Calling a processor when an exception occurred
With the `onExceptionOccurred` you can call a custom processor right after an
exception was thrown,
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 d55dce232801..4ad90c829991 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
@@ -86,7 +86,12 @@ public class DefaultErrorRegistry extends
EventNotifierSupport implements ErrorR
if (event instanceof CamelEvent.ExchangeFailedEvent e) {
capture(e.getExchange(), false);
} else if (event instanceof CamelEvent.ExchangeFailureHandledEvent e) {
- capture(e.getExchange(), true);
+ // the failure processor (such as onException) may not have
handled the exception
+ // (a doCatch handles the exception without the error handler
marking it)
+ Exchange exchange = e.getExchange();
+ boolean handled =
!exchange.getExchangeExtension().isErrorHandlerHandledSet()
+ || exchange.getExchangeExtension().isErrorHandlerHandled();
+ capture(exchange, handled);
}
}
diff --git
a/core/camel-core-engine/src/main/docs/modules/eips/pages/dead-letter-channel.adoc
b/core/camel-core-engine/src/main/docs/modules/eips/pages/dead-letter-channel.adoc
index 37d1a28e2fbc..93e755f663a7 100644
---
a/core/camel-core-engine/src/main/docs/modules/eips/pages/dead-letter-channel.adoc
+++
b/core/camel-core-engine/src/main/docs/modules/eips/pages/dead-letter-channel.adoc
@@ -414,6 +414,15 @@ YAML::
====
+If the `onPrepareFailure` processor throws an exception, then the Exchange is
*not* sent to the dead letter queue,
+and the exception is regarded as a new exception that occurred during the dead
letter channel, which is handled
+according to the `deadLetterHandleNewException` option:
+
+* `deadLetterHandleNewException=true` (default): the new exception is logged
at `WARN` level and handled, and the
+Exchange completes without an exception.
+* `deadLetterHandleNewException=false`: the Exchange fails with the new
exception, which has the original caught
+exception attached as a suppressed exception.
+
=== Calling a processor when an exception occurred
With the `onExceptionOccurred` you can call a custom processor right after an
exception was thrown,
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/FatalFallbackErrorHandler.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/FatalFallbackErrorHandler.java
index e7c340109fc6..f235bd71ca86 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/FatalFallbackErrorHandler.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/FatalFallbackErrorHandler.java
@@ -16,8 +16,9 @@
*/
package org.apache.camel.processor;
-import java.util.ArrayDeque;
import java.util.Deque;
+import java.util.Objects;
+import java.util.concurrent.ConcurrentLinkedDeque;
import org.apache.camel.AsyncCallback;
import org.apache.camel.Exchange;
@@ -60,12 +61,14 @@ public class FatalFallbackErrorHandler extends
DelegateAsyncProcessor implements
final String id = routeIdExpression().evaluate(exchange, String.class);
// prevent endless looping if we end up coming back to ourself
- Deque<String> fatals =
exchange.getProperty(ExchangePropertyKey.FATAL_FALLBACK_ERROR_HANDLER,
Deque.class);
+ // (the fatals are shared by copies of the exchange, such as from
parallel processing in the splitter, so the
+ // entries are for this exchange instance only, and a copy is not
regarded as coming back to ourself)
+ Deque<FatalEntry> fatals =
exchange.getProperty(ExchangePropertyKey.FATAL_FALLBACK_ERROR_HANDLER,
Deque.class);
if (fatals == null) {
- fatals = new ArrayDeque<>();
+ fatals = new ConcurrentLinkedDeque<>();
exchange.setProperty(ExchangePropertyKey.FATAL_FALLBACK_ERROR_HANDLER, fatals);
}
- if (fatals.contains(id)) {
+ if (isCircular(fatals, id, exchange)) {
LOG.warn("Circular error-handler detected at route: {} - breaking
out processing Exchange: {}", id, exchange);
// mark this exchange as already been error handler handled (just
by having this property)
// the false value mean the caught exception will be kept on the
exchange, causing the
@@ -77,7 +80,8 @@ public class FatalFallbackErrorHandler extends
DelegateAsyncProcessor implements
}
// okay we run under this fatal error handler now
- fatals.push(id);
+ final FatalEntry entry = new FatalEntry(id, exchange);
+ fatals.push(entry);
// support the asynchronous routing engine
boolean sync = processor.process(exchange, new AsyncCallback() {
@@ -145,9 +149,13 @@ public class FatalFallbackErrorHandler extends
DelegateAsyncProcessor implements
}
} finally {
// no longer running under this fatal fallback error
handler
- Deque<String> fatals =
exchange.getProperty(ExchangePropertyKey.FATAL_FALLBACK_ERROR_HANDLER,
Deque.class);
+ Deque<FatalEntry> fatals
+ =
exchange.getProperty(ExchangePropertyKey.FATAL_FALLBACK_ERROR_HANDLER,
Deque.class);
if (fatals != null) {
- fatals.removeLastOccurrence(id);
+ fatals.remove(entry);
+ if (fatals.isEmpty()) {
+
exchange.removeProperty(ExchangePropertyKey.FATAL_FALLBACK_ERROR_HANDLER);
+ }
}
callback.done(doneSync);
}
@@ -157,6 +165,33 @@ public class FatalFallbackErrorHandler extends
DelegateAsyncProcessor implements
return sync;
}
+ private static boolean isCircular(Deque<FatalEntry> fatals, String id,
Exchange exchange) {
+ for (FatalEntry entry : fatals) {
+ if (entry.exchange == exchange && Objects.equals(entry.routeId,
id)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ /**
+ * The route (and exchange instance) running under a fatal fallback error
handler (compared by identity).
+ */
+ private static final class FatalEntry {
+ private final String routeId;
+ private final Exchange exchange;
+
+ private FatalEntry(String routeId, Exchange exchange) {
+ this.routeId = routeId;
+ this.exchange = exchange;
+ }
+
+ @Override
+ public String toString() {
+ return routeId;
+ }
+ }
+
private void log(String message) {
log(message, null);
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
index f6dd28ede396..ac2d5de89da2 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
@@ -229,8 +229,11 @@ public class OnCompletionProcessor extends
BaseProcessorSupport
boolean stop = exchange.isRouteStop();
exchange.setRouteStop(false);
boolean failureHandled =
exchange.getExchangeExtension().isFailureHandled();
+ // the onCompletion is not failure handled (so its own failures can be
handled by the error handler)
+ exchange.getExchangeExtension().setFailureHandled(false);
Boolean errorhandlerHandled =
exchange.getExchangeExtension().getErrorHandlerHandled();
exchange.getExchangeExtension().setErrorHandlerHandled(null);
+ Object caught =
exchange.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT);
boolean rollbackOnly = exchange.isRollbackOnly();
exchange.setRollbackOnly(false);
boolean rollbackOnlyLast = exchange.isRollbackOnlyLast();
@@ -251,11 +254,25 @@ public class OnCompletionProcessor extends
BaseProcessorSupport
} finally {
// restore the options
exchange.setRouteStop(stop);
- if (failureHandled) {
- exchange.getExchangeExtension().setFailureHandled(true);
- }
- if (errorhandlerHandled != null) {
+ boolean newFailure = cause == null && exchange.getException() !=
null;
+ if (newFailure) {
+ // the onCompletion failed (and was not handled) so keep its
error handler state
+ if (failureHandled) {
+ exchange.getExchangeExtension().setFailureHandled(true);
+ }
+ if (errorhandlerHandled != null) {
+
exchange.getExchangeExtension().setErrorHandlerHandled(errorhandlerHandled);
+ }
+ } else {
+ // restore the state as it was before the onCompletion (such
as when the onCompletion
+ // handled an exception by its error handler)
+
exchange.getExchangeExtension().setFailureHandled(failureHandled);
exchange.getExchangeExtension().setErrorHandlerHandled(errorhandlerHandled);
+ if (caught != null) {
+ exchange.setProperty(ExchangePropertyKey.EXCEPTION_CAUGHT,
caught);
+ } else {
+
exchange.removeProperty(ExchangePropertyKey.EXCEPTION_CAUGHT);
+ }
}
exchange.setRollbackOnly(rollbackOnly);
exchange.setRollbackOnlyLast(rollbackOnlyLast);
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/DefaultExceptionPolicyStrategy.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/DefaultExceptionPolicyStrategy.java
index e7c19c17f3da..4a1938f10ee0 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/DefaultExceptionPolicyStrategy.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/DefaultExceptionPolicyStrategy.java
@@ -227,7 +227,15 @@ public class DefaultExceptionPolicyStrategy implements
ExceptionPolicyStrategy {
// if no predicate then it's always a match
return true;
}
- return definition.getWhen().matches(exchange);
+ try {
+ return definition.getWhen().matches(exchange);
+ } catch (Exception e) {
+ // a failing predicate is not a match (so other exception policies
can be used)
+ LOG.warn("Error evaluating the onWhen predicate of onException: {}
on exchange: {}. The onException is regarded"
+ + " as not matching. Caused by: {}",
+ definition.getExceptionClass().getName(),
exchange.getExchangeId(), e.getMessage(), e);
+ return false;
+ }
}
/**
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
index 2e97808c98d8..651c559c8ad3 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
@@ -363,6 +363,61 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
return null;
}
+ /**
+ * The onPrepareFailure processor failed (threw an exception), so the
exchange is not delivered to the failure
+ * processor (such as the dead letter channel). A dead letter channel
handles the new exception if
+ * deadLetterHandleNewException is enabled (default), otherwise the
exchange fails with the new exception. The
+ * original caused exception is added as suppressed to the new exception.
+ */
+ protected void handleOnPrepareFailure(
+ Exchange exchange, Exception e, Processor processor, boolean
isDeadLetterChannel) {
+ Throwable original =
exchange.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT, Throwable.class);
+ if (original != null && original != e) {
+ e.addSuppressed(original);
+ }
+ ExchangeHelper.setFailureHandled(exchange);
+ boolean handled = isDeadLetterChannel && deadLetterHandleNewException;
+ String target = isDeadLetterChannel && deadLetterUri != null ?
URISupport.sanitizeUri(deadLetterUri) : "" + processor;
+ LOG.warn("Error during processing the onPrepareFailure processor on
exchange: {}. The exchange is not delivered to: {}"
+ + " and the new exception is {}. Caused by: {}",
+ exchange.getExchangeId(), target, handled ? "handled" : "not
handled", e.getMessage(), e);
+ if (handled) {
+ exchange.setException(null);
+ exchange.getExchangeExtension().setErrorHandlerHandled(true);
+ } else {
+ exchange.getExchangeExtension().setErrorHandlerHandled(false);
+ exchange.setException(e);
+ }
+ }
+
+ /**
+ * Evaluates a predicate of the error handler (such as handled or
continued). A predicate that fails (throws an
+ * exception) is regarded as <tt>false</tt>, as the error handler must
still handle the exchange (such as the dead
+ * letter channel which always handles the exchange). The original caused
exception is kept, and the exception from
+ * the predicate is logged and added as suppressed to the original
exception.
+ */
+ static boolean matchesPredicate(Predicate predicate, Exchange exchange,
String name) {
+ try {
+ return predicate.matches(exchange);
+ } catch (Exception e) {
+ onPredicateFailure(exchange, name, e);
+ return false;
+ }
+ }
+
+ static void onPredicateFailure(Exchange exchange, String name, Exception
e) {
+ Throwable original =
exchange.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT, Throwable.class);
+ if (original == null) {
+ original = exchange.getException();
+ }
+ if (original != null && original != e) {
+ original.addSuppressed(e);
+ }
+ LOG.warn("Error evaluating the {} predicate of the error handler on
exchange: {}. The predicate is regarded as false."
+ + " Caused by: {}",
+ name, exchange.getExchangeId(), e.getMessage(), e);
+ }
+
/**
* Simple task to perform calling the processor with no redelivery support
*/
@@ -594,12 +649,23 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
// invoke custom on prepare
if (onPrepareProcessor != null) {
+ Exception prepareException;
try {
LOG.trace("OnPrepare processor {} is processing
Exchange: {}", onPrepareProcessor, exchange);
onPrepareProcessor.process(exchange);
+ // the processor may be wrapped and set the exception
on the exchange instead of throwing
+ prepareException = exchange.getException();
} catch (Exception e) {
- // a new exception was thrown during prepare
- exchange.setException(e);
+ prepareException = e;
+ }
+ if (prepareException != null) {
+ // a new exception was thrown during prepare, then the
exchange is not delivered to the
+ // failure processor
+ handleOnPrepareFailure(exchange, prepareException,
processor, isDeadLetterChannel);
+ AsyncCallback cb = callback;
+ taskFactory.release(this);
+ reactiveExecutor.schedule(cb);
+ return;
}
}
@@ -693,14 +759,14 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
private boolean shouldHandle(Exchange exchange, Predicate
handledPredicate) {
if (handledPredicate != null) {
- return handledPredicate.matches(exchange);
+ return matchesPredicate(handledPredicate, exchange, "handled");
}
return false;
}
private boolean shouldContinue(Exchange exchange, Predicate
continuedPredicate) {
if (continuedPredicate != null) {
- return continuedPredicate.matches(exchange);
+ return matchesPredicate(continuedPredicate, exchange,
"continued");
}
return false;
}
@@ -804,6 +870,10 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
// e is never null
Throwable previous =
exchange.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT, Throwable.class);
+ if (previous != null && previous ==
exchange.getProperty(ExchangePropertyKey.EXCEPTION_HANDLED)) {
+ // the previous exception was handled (continued) so this is a
new failure
+ previous = null;
+ }
if (previous != null && previous != e && e != null) {
// a 2nd exception was thrown while handling a previous
exception
// so we need to add the previous as suppressed by the new
exception
@@ -1059,9 +1129,10 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
doRun();
} catch (Exception e) {
// unexpected exception during running so set exception and
trigger callback
- // (do not do taskFactory.release as that happens later)
exchange.setException(e);
- callback.done(false);
+ AsyncCallback cb = callback;
+ taskFactory.release(this);
+ cb.done(false);
}
}
@@ -1080,8 +1151,14 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
exhausted =
exchange.getExchangeExtension().isRedeliveryExhausted() ||
exchange.isRollbackOnly();
if (!exhausted && redeliveryCounter > 0) {
// its a potential redelivery so determine if we should
redeliver or not
- redeliverAllowed
- =
currentRedeliveryPolicy.shouldRedeliver(exchange, redeliveryCounter,
retryWhilePredicate);
+ try {
+ redeliverAllowed
+ =
currentRedeliveryPolicy.shouldRedeliver(exchange, redeliveryCounter,
retryWhilePredicate);
+ } catch (Exception e) {
+ // the retryWhile predicate failed, so do not
redeliver (the failure processor handles it)
+ onPredicateFailure(exchange, "retryWhile", e);
+ redeliverAllowed = false;
+ }
}
}
// if we are exhausted or redelivery is not allowed, then deliver
to failure processor (eg such as DLC)
@@ -1211,6 +1288,13 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
// letting onRedeliver be executed at first
deliverToOnRedeliveryProcessor();
+ if (exchange.getException() != null) {
+ // the on redelivery processor failed, which is a new failure
(the exchange must not be redelivered
+ // with the exception set), so loop back around to handle the
new exception
+ reactiveExecutor.schedule(this);
+ return;
+ }
+
if (exchange.isRouteStop()) {
// the on redelivery can mark that the exchange should stop
and therefore not perform a redelivery
// and if so then we are done so continue callback
@@ -1266,6 +1350,11 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
exchange.getIn().removeHeader(Exchange.REDELIVERY_MAX_COUNTER);
exchange.getExchangeExtension().setFailureHandled(false);
// keep the Exchange.EXCEPTION_CAUGHT as property so end user
knows the caused exception
+ // and mark it as handled (continued), so it is not regarded as a
previous exception of a later failure
+ Exception handled =
exchange.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT, Exception.class);
+ if (handled != null) {
+ exchange.setProperty(ExchangePropertyKey.EXCEPTION_HANDLED,
handled);
+ }
// create log message
String msg = "Failed delivery for " +
ExchangeHelper.logIds(exchange) + failureOrigin(exchange);
@@ -1317,6 +1406,10 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
// e is never null
Throwable previous =
exchange.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT, Throwable.class);
+ if (previous != null && previous ==
exchange.getProperty(ExchangePropertyKey.EXCEPTION_HANDLED)) {
+ // the previous exception was handled (continued) so this is a
new failure
+ previous = null;
+ }
if (previous != null && previous != e) {
// a 2nd exception was thrown while handling a previous
exception
// so we need to add the previous as suppressed by the new
exception
@@ -1526,12 +1619,23 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
// invoke custom on prepare
if (onPrepareProcessor != null) {
+ Exception prepareException;
try {
LOG.trace("OnPrepare processor {} is processing
Exchange: {}", onPrepareProcessor, exchange);
onPrepareProcessor.process(exchange);
+ // the processor may be wrapped and set the exception
on the exchange instead of throwing
+ prepareException = exchange.getException();
} catch (Exception e) {
- // a new exception was thrown during prepare
- exchange.setException(e);
+ prepareException = e;
+ }
+ if (prepareException != null) {
+ // a new exception was thrown during prepare, then the
exchange is not delivered to the
+ // failure processor
+ handleOnPrepareFailure(exchange, prepareException,
processor, isDeadLetterChannel);
+ AsyncCallback cb = callback;
+ taskFactory.release(this);
+ reactiveExecutor.schedule(cb);
+ return;
}
}
@@ -1845,7 +1949,7 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
*/
private boolean shouldContinue(Exchange exchange) {
if (continuedPredicate != null) {
- return continuedPredicate.matches(exchange);
+ return matchesPredicate(continuedPredicate, exchange,
"continued");
}
// do not continue by default
return false;
@@ -1859,7 +1963,7 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
*/
private boolean shouldHandle(Exchange exchange) {
if (handledPredicate != null) {
- return handledPredicate.matches(exchange);
+ return matchesPredicate(handledPredicate, exchange, "handled");
}
// do not handle by default
return false;
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/onexception/ErrorHandlerFailurePathsTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/onexception/ErrorHandlerFailurePathsTest.java
new file mode 100644
index 000000000000..8ea0d86e69d5
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/onexception/ErrorHandlerFailurePathsTest.java
@@ -0,0 +1,277 @@
+/*
+ * 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.onexception;
+
+import java.io.IOException;
+import java.util.Collection;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePropertyKey;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.BacklogErrorEventMessage;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+public class ErrorHandlerFailurePathsTest extends ContextTestSupport {
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Test
+ public void testOnWhenThrowsIsNoMatch() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ onException(IOException.class).onWhen(e -> {
+ throw new IllegalArgumentException("onWhen failed");
+ }).handled(true).to("mock:when");
+ onException(Exception.class).handled(true).to("mock:fallback");
+
+ from("direct:start").throwException(new IOException("Forced"));
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:when").expectedMessageCount(0);
+ getMockEndpoint("mock:fallback").expectedMessageCount(1);
+ template.sendBody("direct:start", "Hello");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testDeadLetterChannelRetryWhileThrows() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ errorHandler(deadLetterChannel("mock:dead"));
+
onException(IOException.class).maximumRedeliveries(2).retryWhile(e -> {
+ throw new IllegalArgumentException("retryWhile failed");
+ });
+
+ from("direct:start").throwException(new IOException("Forced"));
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:dead").expectedMessageCount(1);
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertMockEndpointsSatisfied();
+ // the dead letter channel handled the exchange
+ assertNull(out.getException());
+ Exception caught =
out.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT, Exception.class);
+ assertEquals(IOException.class, caught.getClass());
+ assertEquals(1, caught.getSuppressed().length);
+ }
+
+ @Test
+ public void testDeadLetterChannelHandledThrows() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ errorHandler(deadLetterChannel("mock:dead"));
+ onException(IOException.class).handled(e -> {
+ throw new IllegalArgumentException("handled failed");
+ });
+
+ from("direct:start").throwException(new IOException("Forced"));
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:dead").expectedMessageCount(1);
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertMockEndpointsSatisfied();
+ assertNull(out.getException());
+ }
+
+ @Test
+ public void testParallelSplitInsideOnException() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ onException(IllegalStateException.class).handled(true)
+
.split(body().tokenize(",")).parallelProcessing().to("direct:x");
+
onException(IllegalArgumentException.class).handled(true).to("direct:handler");
+
+ from("direct:start").throwException(new
IllegalStateException("Forced"));
+ from("direct:x").throwException(new
IllegalArgumentException("x failed"));
+ from("direct:handler").to("mock:x").delay(100);
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:x").expectedMessageCount(4);
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("a,b,c,d"));
+ assertMockEndpointsSatisfied();
+ assertNull(out.getException());
+ }
+
+ @Test
+ public void testOnPrepareFailureThrows() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+
.errorHandler(deadLetterChannel("mock:dead").onPrepareFailure(e -> {
+ throw new IllegalStateException("prepare failed");
+ }))
+ .throwException(new IOException("Forced"));
+
+ from("direct:start2")
+
.errorHandler(deadLetterChannel("mock:dead").deadLetterHandleNewException(false)
+ .onPrepareFailure(e -> {
+ throw new IllegalStateException("prepare
failed");
+ }))
+ .throwException(new IOException("Forced"));
+ }
+ });
+ context.start();
+
+ // the exchange is not delivered to the dead letter channel
+ getMockEndpoint("mock:dead").expectedMessageCount(0);
+
+ // the new exception is handled (deadLetterHandleNewException=true)
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertNull(out.getException());
+
+ // the new exception is not handled
+ out = template.send("direct:start2", e ->
e.getMessage().setBody("Hello"));
+ assertNotNull(out.getException());
+ assertEquals(IllegalStateException.class,
out.getException().getClass());
+ assertEquals(IOException.class,
out.getException().getSuppressed()[0].getClass());
+
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testErrorRegistryNotHandled() throws Exception {
+ context.getErrorRegistry().setEnabled(true);
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ onException(IllegalArgumentException.class).to("mock:error");
+
+ from("direct:start").throwException(new
IllegalArgumentException("Forced"));
+ }
+ });
+ context.start();
+
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertNotNull(out.getException());
+
+ Collection<BacklogErrorEventMessage> entries =
context.getErrorRegistry().browse();
+ assertEquals(1, entries.size());
+ assertFalse(entries.iterator().next().isHandled());
+ }
+
+ @Test
+ public void testOnRedeliveryThrows() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
errorHandler(defaultErrorHandler().maximumRedeliveries(2).redeliveryDelay(0).onRedelivery(e
-> {
+ throw new IllegalStateException("onRedelivery failed");
+ }));
+
+ from("direct:start").to("mock:target");
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:target").whenAnyExchangeReceived(e -> {
+ throw new IOException("Forced");
+ });
+
+ // the target is not called again with the exception from the on
redelivery processor
+ getMockEndpoint("mock:target").expectedMessageCount(1);
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertMockEndpointsSatisfied();
+ assertEquals(IllegalStateException.class,
out.getException().getClass());
+ }
+
+ @Test
+ public void testContinuedThenUnrelatedFailure() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ onException(IllegalArgumentException.class).continued(true);
+
+ from("direct:start")
+ .throwException(new IllegalArgumentException("first"))
+ .throwException(new IllegalStateException("second"));
+ }
+ });
+ context.start();
+
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertEquals(IllegalStateException.class,
out.getException().getClass());
+ // the continued exception is not a previous exception of the new
failure
+ assertEquals(0, out.getException().getSuppressed().length);
+ }
+
+ @Test
+ public void testOnCompletionFailureHandledByOnException() throws Exception
{
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
onException(IllegalStateException.class).handled(true).to("mock:handled");
+ onCompletion().onFailureOnly().modeBeforeConsumer().process(e
-> {
+ throw new IllegalStateException("onCompletion failed");
+ });
+
+ from("direct:start").throwException(new
IllegalArgumentException("Forced"));
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:handled").expectedMessageCount(1);
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertMockEndpointsSatisfied();
+ assertEquals(IllegalArgumentException.class,
out.getException().getClass());
+ }
+
+ @Test
+ public void testOnCompletionHandledFailureNotLeaked() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
onException(IllegalStateException.class).handled(true).to("mock:handled");
+ onCompletion().onCompleteOnly().modeBeforeConsumer().process(e
-> {
+ throw new IllegalStateException("onCompletion failed");
+ });
+
+ from("direct:start").to("mock:result");
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:handled").expectedMessageCount(1);
+ Exchange out = template.send("direct:start", e ->
e.getMessage().setBody("Hello"));
+ assertMockEndpointsSatisfied();
+ assertNull(out.getException());
+ // the handled failure of the onCompletion is not left on the
(successful) exchange
+ assertNull(out.getProperty(ExchangePropertyKey.EXCEPTION_CAUGHT));
+ assertFalse(out.getExchangeExtension().isErrorHandlerHandledSet());
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/onexception/OnExceptionHandledThrowsExceptionTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/onexception/OnExceptionHandledThrowsExceptionTest.java
index 0ad0417a9514..3f0e0b832fe5 100644
---
a/core/camel-core/src/test/java/org/apache/camel/processor/onexception/OnExceptionHandledThrowsExceptionTest.java
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/onexception/OnExceptionHandledThrowsExceptionTest.java
@@ -29,13 +29,18 @@ public class OnExceptionHandledThrowsExceptionTest extends
ContextTestSupport {
@Test
public void testHandled() throws Exception {
- getMockEndpoint("mock:handled").expectedMessageCount(0);
+ // the handled predicate fails, which is regarded as not handled, so
the onException is still processed
+ // and the exchange fails with the original exception (with the
exception from the predicate as suppressed)
+ getMockEndpoint("mock:handled").expectedMessageCount(1);
try {
template.sendBody("direct:start", "Hello World");
fail("Should have thrown exception");
} catch (Exception e) {
- IllegalArgumentException iae =
assertIsInstanceOf(IllegalArgumentException.class, e.getCause());
+ IOException io = assertIsInstanceOf(IOException.class,
e.getCause());
+ assertEquals("Forced", io.getMessage());
+ assertEquals(1, io.getSuppressed().length);
+ IllegalArgumentException iae =
assertIsInstanceOf(IllegalArgumentException.class, io.getSuppressed()[0]);
assertEquals("Another Forced", iae.getMessage());
}
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 92d26282ea77..c4719e95dcf1 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
@@ -3053,6 +3053,35 @@ always ran into the graceful shutdown timeout and was
then forced (or aborted).
(or purged with `purgeWhenStopping=true`), and are processed if the route is
started again while the queue still
exists.
+=== camel-core - error handler predicates that throw an exception
+
+When an `onWhen`, `handled`, `continued` or `retryWhile` predicate of
`onException` (or `retryWhile` of the error
+handler) throws an exception during error handling, then the error handler now
logs this at `WARN` level and regards
+the predicate as `false`:
+
+* A failing `onWhen` means the `onException` does not match, so another
`onException` (or the error handler itself)
+handles the exception instead. Previously the exception from the predicate
escaped from the error handler.
+* A failing `handled`, `continued` or `retryWhile` means the exchange is not
handled, not continued and not redelivered.
+The exchange is still processed by the `onException` outputs (or sent to the
dead letter channel), and the original
+exception is kept as the exception, with the exception from the predicate
attached as a suppressed exception.
+Previously the exception from the predicate replaced the original exception,
and the `onException` outputs
+were skipped.
+
+=== camel-core - error handler onPrepareFailure that throws an exception
+
+When the `onPrepareFailure` processor of the Dead Letter Channel or the
Default Error Handler throws an exception,
+then the exchange is no longer sent to the dead letter queue (or failure
processor) with that exception set on it.
+Instead, the exchange is not delivered, and with the Dead Letter Channel the
new exception is handled
+according to `deadLetterHandleNewException` (by default `true`, so the
exchange completes and a `WARN` is logged).
+With `deadLetterHandleNewException=false`, or with the Default Error Handler,
the exchange fails with the new exception,
+which has the original exception attached as a suppressed exception.
+
+=== camel-core - error handler error registry handled flag
+
+The error registry (`ErrorRegistry`) now records whether an error was handled
from the exchange itself. An exception
+that is processed by an `onException` without `handled(true)` is now recorded
as not handled; previously any exception
+that was processed by an `onException` or a failure processor was recorded as
handled.
+
=== camel-sql, camel-sql-stored - the query/template override headers are gated
A message header could override the endpoint-configured SQL by default,
letting an incoming message choose the