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

Reply via email to