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 c10aa0cc60dd CAMEL-25047: camel-core - Reifiers and error handler: fix
bugs found in a deep review (#26928)
c10aa0cc60dd is described below
commit c10aa0cc60dd8758441caea86ad58915c9ad9f6a
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 09:10:20 2026 +0200
CAMEL-25047: camel-core - Reifiers and error handler: fix bugs found in a
deep review (#26928)
- onException useOriginalBody turns on allowUseOriginalMessage (as
useOriginalMessage does)
- the redelivery options of onException that are not set are inherited from
the error handler
instead of being reset to their defaults
- logName of defaultErrorHandler/deadLetterChannel is used
- a disabled multicast or pipeline with a single output is skipped
- an unknown executorServiceRef on an error handler fails with a clear
message instead of NPE
- placeholders are resolved in delay/throttle executorService,
aggregationStrategyMethodName
and the deadLetterChannel retryWhileRef/executorServiceRef
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../apache/camel/reifier/ClaimCheckReifier.java | 3 +-
.../org/apache/camel/reifier/MulticastReifier.java | 7 +-
.../apache/camel/reifier/OnExceptionReifier.java | 6 +-
.../org/apache/camel/reifier/PipelineReifier.java | 9 +
.../org/apache/camel/reifier/ProcessorReifier.java | 4 +-
.../org/apache/camel/reifier/SplitReifier.java | 3 +-
.../errorhandler/DeadLetterChannelReifier.java | 10 +-
.../errorhandler/DefaultErrorHandlerReifier.java | 7 +-
.../reifier/errorhandler/ErrorHandlerReifier.java | 32 ++-
.../apache/camel/reifier/ReifierEdgeCasesTest.java | 226 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 7 +
11 files changed, 283 insertions(+), 31 deletions(-)
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ClaimCheckReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ClaimCheckReifier.java
index 2504fa737666..663e25812fea 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ClaimCheckReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ClaimCheckReifier.java
@@ -124,7 +124,8 @@ public class ClaimCheckReifier extends
ProcessorReifier<ClaimCheckDefinition> {
} else if (aggStrategy instanceof BiFunction biFunction) {
strategy = new
AggregationStrategyBiFunctionAdapter(biFunction);
} else if (aggStrategy != null) {
- strategy = new AggregationStrategyBeanAdapter(aggStrategy,
definition.getAggregationStrategyMethodName());
+ strategy = new AggregationStrategyBeanAdapter(
+ aggStrategy,
parseString(definition.getAggregationStrategyMethodName()));
} else {
throw new IllegalArgumentException(
"Cannot find AggregationStrategy in Registry with
name: " + definition.getAggregationStrategy());
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/MulticastReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/MulticastReifier.java
index ef1df16ddcec..e23f0894755f 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/MulticastReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/MulticastReifier.java
@@ -42,9 +42,6 @@ public class MulticastReifier extends
ProcessorReifier<MulticastDefinition> {
@Override
public Processor createProcessor() throws Exception {
Processor answer = this.createChildProcessor(true);
- if (answer instanceof DisabledAware da) {
- da.setDisabled(isDisabled(camelContext, definition));
- }
// force the answer as a multicast processor even if there is only one
// child processor in the multicast
@@ -53,6 +50,10 @@ public class MulticastReifier extends
ProcessorReifier<MulticastDefinition> {
list.add(answer);
answer = createCompositeProcessor(list);
}
+ // set disabled on the multicast (and not on a single output, which it
wraps)
+ if (answer instanceof DisabledAware da) {
+ da.setDisabled(isDisabled(camelContext, definition));
+ }
return answer;
}
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/OnExceptionReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/OnExceptionReifier.java
index d1f8764174b2..da109a0abf7f 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/OnExceptionReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/OnExceptionReifier.java
@@ -49,7 +49,8 @@ public class OnExceptionReifier extends
ProcessorReifier<OnExceptionDefinition>
definition.validateConfiguration();
}
- if (parseBoolean(definition.getUseOriginalMessage(), false)) {
+ if (parseBoolean(definition.getUseOriginalMessage(), false)
+ || parseBoolean(definition.getUseOriginalBody(), false)) {
// ensure allow original is turned on
route.setAllowUseOriginalMessage(true);
}
@@ -82,7 +83,8 @@ public class OnExceptionReifier extends
ProcessorReifier<OnExceptionDefinition>
classes = createExceptionClasses(camelContext.getClassResolver());
}
- if (parseBoolean(definition.getUseOriginalMessage(), false)) {
+ if (parseBoolean(definition.getUseOriginalMessage(), false)
+ || parseBoolean(definition.getUseOriginalBody(), false)) {
// ensure allow original is turned on
route.setAllowUseOriginalMessage(true);
}
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/PipelineReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/PipelineReifier.java
index d7267101df90..0f600e44fb35 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/PipelineReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/PipelineReifier.java
@@ -16,11 +16,15 @@
*/
package org.apache.camel.reifier;
+import java.util.List;
+
+import org.apache.camel.Channel;
import org.apache.camel.DisabledAware;
import org.apache.camel.Processor;
import org.apache.camel.Route;
import org.apache.camel.model.PipelineDefinition;
import org.apache.camel.model.ProcessorDefinition;
+import org.apache.camel.processor.Pipeline;
public class PipelineReifier extends ProcessorReifier<PipelineDefinition> {
@@ -31,6 +35,11 @@ public class PipelineReifier extends
ProcessorReifier<PipelineDefinition> {
@Override
public Processor createProcessor() throws Exception {
Processor answer = this.createChildProcessor(true);
+ if (answer instanceof Channel) {
+ // a single output is its channel, which is not where disabled is
checked at runtime,
+ // so use a pipeline of the single output
+ answer = new Pipeline(camelContext, List.of(answer));
+ }
if (answer instanceof DisabledAware da) {
da.setDisabled(isDisabled(camelContext, definition));
}
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java
index 30e93a906380..594045c7389b 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java
@@ -472,7 +472,7 @@ public abstract class ProcessorReifier<T extends
ProcessorDefinition<?>> extends
+ " is not an
ScheduledExecutorService instance");
} else if (definition.getExecutorServiceRef() != null) {
ScheduledExecutorService answer =
lookupScheduledExecutorServiceRef(name, definition,
- definition.getExecutorServiceRef());
+ parseString(definition.getExecutorServiceRef()));
if (answer == null) {
throw new IllegalArgumentException(
"ExecutorServiceRef " +
definition.getExecutorServiceRef()
@@ -955,7 +955,7 @@ public abstract class ProcessorReifier<T extends
ProcessorDefinition<?>> extends
// closure.
AggregationStrategyBeanAdapter adapter = new
AggregationStrategyBeanAdapter(
aggStrategy,
- definition.getAggregationStrategyMethodName());
+
parseString(definition.getAggregationStrategyMethodName()));
if (definition.getAggregationStrategyMethodAllowNull() !=
null) {
adapter.setAllowNullNewExchange(
parseBoolean(definition.getAggregationStrategyMethodAllowNull(), false));
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SplitReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SplitReifier.java
index a9abe990aebd..348af5aa5293 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SplitReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SplitReifier.java
@@ -168,7 +168,8 @@ public class SplitReifier extends
ExpressionReifier<SplitDefinition> {
@SuppressWarnings("resource")
// NOTE: the adapter holds no leaking resources, so we can
safely ignore its closure.
AggregationStrategyBeanAdapter adapter
- = new AggregationStrategyBeanAdapter(aggStrategy,
definition.getAggregationStrategyMethodName());
+ = new AggregationStrategyBeanAdapter(
+ aggStrategy,
parseString(definition.getAggregationStrategyMethodName()));
if (definition.getAggregationStrategyMethodAllowNull() !=
null) {
adapter.setAllowNullNewExchange(parseBoolean(definition.getAggregationStrategyMethodAllowNull(),
false));
adapter.setAllowNullOldExchange(parseBoolean(definition.getAggregationStrategyMethodAllowNull(),
false));
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DeadLetterChannelReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DeadLetterChannelReifier.java
index bbcbf8bfeb8e..b2e881850be1 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DeadLetterChannelReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DeadLetterChannelReifier.java
@@ -76,7 +76,7 @@ public class DeadLetterChannelReifier extends
ErrorHandlerReifier<DeadLetterChan
if (answer == null && definition.getRetryWhileRef() != null) {
// it is a bean expression
Language bean = camelContext.resolveLanguage("bean");
- answer = bean.createPredicate(definition.getRetryWhileRef());
+ answer =
bean.createPredicate(parseString(definition.getRetryWhileRef()));
answer.initPredicate(camelContext);
}
@@ -88,6 +88,9 @@ public class DeadLetterChannelReifier extends
ErrorHandlerReifier<DeadLetterChan
if (answer == null && definition.getLoggerRef() != null) {
answer = mandatoryLookup(definition.getLoggerRef(),
CamelLogger.class);
}
+ if (answer == null && definition.getLogName() != null) {
+ answer = new
CamelLogger(LoggerFactory.getLogger(parseString(definition.getLogName())),
LoggingLevel.ERROR);
+ }
if (answer == null) {
answer = new
CamelLogger(LoggerFactory.getLogger(DeadLetterChannel.class),
LoggingLevel.ERROR);
}
@@ -134,6 +137,7 @@ public class DeadLetterChannelReifier extends
ErrorHandlerReifier<DeadLetterChan
ScheduledExecutorService executorService, String
executorServiceRef) {
lock.lock();
try {
+ executorServiceRef = parseString(executorServiceRef);
if (executorService == null || executorService.isShutdown()) {
// camel context will shutdown the executor when it shutdown
so no
// need to shut it down when stopping
@@ -142,7 +146,9 @@ public class DeadLetterChannelReifier extends
ErrorHandlerReifier<DeadLetterChan
if (executorService == null) {
ExecutorServiceManager manager =
camelContext.getExecutorServiceManager();
ThreadPoolProfile profile =
manager.getThreadPoolProfile(executorServiceRef);
- executorService = manager.newScheduledThreadPool(this,
executorServiceRef, profile);
+ if (profile != null) {
+ executorService =
manager.newScheduledThreadPool(this, executorServiceRef, profile);
+ }
}
if (executorService == null) {
throw new IllegalArgumentException("ExecutorService "
+ executorServiceRef + " not found in registry.");
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DefaultErrorHandlerReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DefaultErrorHandlerReifier.java
index 198c839e8d31..ca934917d5bc 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DefaultErrorHandlerReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/DefaultErrorHandlerReifier.java
@@ -62,6 +62,9 @@ public class DefaultErrorHandlerReifier extends
ErrorHandlerReifier<DefaultError
if (answer == null && definition.getLoggerRef() != null) {
answer = mandatoryLookup(definition.getLoggerRef(),
CamelLogger.class);
}
+ if (answer == null && definition.getLogName() != null) {
+ answer = new
CamelLogger(LoggerFactory.getLogger(parseString(definition.getLogName())),
LoggingLevel.ERROR);
+ }
if (answer == null) {
answer = new
CamelLogger(LoggerFactory.getLogger(DefaultErrorHandler.class),
LoggingLevel.ERROR);
}
@@ -105,7 +108,9 @@ public class DefaultErrorHandlerReifier extends
ErrorHandlerReifier<DefaultError
if (executorService == null) {
ExecutorServiceManager manager =
camelContext.getExecutorServiceManager();
ThreadPoolProfile profile =
manager.getThreadPoolProfile(executorServiceRef);
- executorService = manager.newScheduledThreadPool(this,
executorServiceRef, profile);
+ if (profile != null) {
+ executorService =
manager.newScheduledThreadPool(this, executorServiceRef, profile);
+ }
}
if (executorService == null) {
throw new IllegalArgumentException("ExecutorService "
+ executorServiceRef + " not found in registry.");
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/ErrorHandlerReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/ErrorHandlerReifier.java
index efec07c9af87..fe8e307a267c 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/ErrorHandlerReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/errorhandler/ErrorHandlerReifier.java
@@ -251,44 +251,38 @@ public abstract class ErrorHandlerReifier<T extends
ErrorHandlerFactory> extends
if (definition == null) {
return null;
}
+ // only the options that are set, so the others are inherited from the
error handler (when used by onException)
Map<RedeliveryOption, String> policy = new
EnumMap<>(RedeliveryOption.class);
setOption(policy, RedeliveryOption.maximumRedeliveries,
definition.getMaximumRedeliveries());
- setOption(policy, RedeliveryOption.redeliveryDelay,
definition.getRedeliveryDelay(), "1000");
+ setOption(policy, RedeliveryOption.redeliveryDelay,
definition.getRedeliveryDelay());
setOption(policy, RedeliveryOption.asyncDelayedRedelivery,
definition.getAsyncDelayedRedelivery());
- setOption(policy, RedeliveryOption.backOffMultiplier,
definition.getBackOffMultiplier(), "2");
+ setOption(policy, RedeliveryOption.backOffMultiplier,
definition.getBackOffMultiplier());
setOption(policy, RedeliveryOption.useExponentialBackOff,
definition.getUseExponentialBackOff());
- setOption(policy, RedeliveryOption.collisionAvoidanceFactor,
definition.getCollisionAvoidanceFactor(), "0.15");
+ setOption(policy, RedeliveryOption.collisionAvoidanceFactor,
definition.getCollisionAvoidanceFactor());
setOption(policy, RedeliveryOption.useCollisionAvoidance,
definition.getUseCollisionAvoidance());
- setOption(policy, RedeliveryOption.maximumRedeliveryDelay,
definition.getMaximumRedeliveryDelay(), "60000");
- setOption(policy, RedeliveryOption.retriesExhaustedLogLevel,
definition.getRetriesExhaustedLogLevel(), "ERROR");
- setOption(policy, RedeliveryOption.retryAttemptedLogLevel,
definition.getRetryAttemptedLogLevel(), "DEBUG");
- setOption(policy, RedeliveryOption.retryAttemptedLogInterval,
definition.getRetryAttemptedLogInterval(), "1");
- setOption(policy, RedeliveryOption.logRetryAttempted,
definition.getLogRetryAttempted(), "true");
- setOption(policy, RedeliveryOption.logStackTrace,
definition.getLogStackTrace(), "true");
+ setOption(policy, RedeliveryOption.maximumRedeliveryDelay,
definition.getMaximumRedeliveryDelay());
+ setOption(policy, RedeliveryOption.retriesExhaustedLogLevel,
definition.getRetriesExhaustedLogLevel());
+ setOption(policy, RedeliveryOption.retryAttemptedLogLevel,
definition.getRetryAttemptedLogLevel());
+ setOption(policy, RedeliveryOption.retryAttemptedLogInterval,
definition.getRetryAttemptedLogInterval());
+ setOption(policy, RedeliveryOption.logRetryAttempted,
definition.getLogRetryAttempted());
+ setOption(policy, RedeliveryOption.logStackTrace,
definition.getLogStackTrace());
setOption(policy, RedeliveryOption.logRetryStackTrace,
definition.getLogRetryStackTrace());
setOption(policy, RedeliveryOption.logHandled,
definition.getLogHandled());
- setOption(policy, RedeliveryOption.logNewException,
definition.getLogNewException(), "true");
+ setOption(policy, RedeliveryOption.logNewException,
definition.getLogNewException());
setOption(policy, RedeliveryOption.logContinued,
definition.getLogContinued());
- setOption(policy, RedeliveryOption.logExhausted,
definition.getLogExhausted(), "true");
+ setOption(policy, RedeliveryOption.logExhausted,
definition.getLogExhausted());
setOption(policy, RedeliveryOption.logExhaustedMessageHistory,
definition.getLogExhaustedMessageHistory());
setOption(policy, RedeliveryOption.logExhaustedMessageBody,
definition.getLogExhaustedMessageBody());
setOption(policy, RedeliveryOption.disableRedelivery,
definition.getDisableRedelivery());
setOption(policy, RedeliveryOption.delayPattern,
definition.getDelayPattern());
- setOption(policy, RedeliveryOption.allowRedeliveryWhileStopping,
definition.getAllowRedeliveryWhileStopping(), "true");
+ setOption(policy, RedeliveryOption.allowRedeliveryWhileStopping,
definition.getAllowRedeliveryWhileStopping());
setOption(policy, RedeliveryOption.exchangeFormatterRef,
definition.getExchangeFormatterRef());
return policy;
}
private void setOption(Map<RedeliveryOption, String> policy,
RedeliveryOption option, String value) {
- setOption(policy, option, value, null);
- }
-
- private void setOption(
- Map<RedeliveryOption, String> policy, RedeliveryOption option,
String value, String defaultValue) {
if (value != null) {
policy.put(option, parseString(value));
- } else if (defaultValue != null) {
- policy.put(option, defaultValue);
}
}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/reifier/ReifierEdgeCasesTest.java
b/core/camel-core/src/test/java/org/apache/camel/reifier/ReifierEdgeCasesTest.java
new file mode 100644
index 000000000000..9cb44d1b2e69
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/reifier/ReifierEdgeCasesTest.java
@@ -0,0 +1,226 @@
+/*
+ * 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.reifier;
+
+import java.io.IOException;
+import java.util.List;
+import java.util.Properties;
+
+import org.apache.camel.Channel;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.FailedToCreateRouteException;
+import org.apache.camel.Processor;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.model.errorhandler.DeadLetterChannelDefinition;
+import org.apache.camel.model.errorhandler.DefaultErrorHandlerDefinition;
+import org.apache.camel.processor.errorhandler.RedeliveryErrorHandler;
+import org.apache.camel.spi.ThreadPoolProfile;
+import org.apache.camel.util.StopWatch;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+public class ReifierEdgeCasesTest extends ContextTestSupport {
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ private void start(RouteBuilder builder) throws Exception {
+ context.addRoutes(builder);
+ context.start();
+ }
+
+ private RedeliveryErrorHandler firstErrorHandler() {
+ List<Processor> nodes = context.getRoutes().get(0).navigate().next();
+ Channel channel = unwrapChannel(nodes.get(0));
+ return (RedeliveryErrorHandler) channel.getErrorHandler();
+ }
+
+ @Test
+ public void testOnExceptionUseOriginalBody() throws Exception {
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+
onException(IllegalArgumentException.class).useOriginalBody().handled(true).to("mock:dead");
+
+
from("direct:start").setBody(constant("Changed")).throwException(new
IllegalArgumentException("Forced"));
+ }
+ });
+ getMockEndpoint("mock:dead").expectedBodiesReceived("Hello");
+ template.sendBody("direct:start", "Hello");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testDefaultErrorHandlerLogName() throws Exception {
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+ DefaultErrorHandlerDefinition eh = new
DefaultErrorHandlerDefinition();
+ eh.setLogName("com.foo.Errors");
+ errorHandler(eh);
+
+ from("direct:start").to("mock:result");
+ }
+ });
+
assertThat(firstErrorHandler().getLogger().getLog().getName()).isEqualTo("com.foo.Errors");
+ }
+
+ @Test
+ public void testDeadLetterChannelLogName() throws Exception {
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+ DeadLetterChannelDefinition eh = new
DeadLetterChannelDefinition("mock:dead");
+ eh.setLogName("com.foo.Dead");
+ errorHandler(eh);
+
+ from("direct:start").to("mock:result");
+ }
+ });
+
assertThat(firstErrorHandler().getLogger().getLog().getName()).isEqualTo("com.foo.Dead");
+ }
+
+ @Test
+ public void testDisabledMulticastWithSingleOutput() throws Exception {
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("direct:start").multicast().disabled().to("mock:a").end().to("mock:result");
+ }
+ });
+ getMockEndpoint("mock:a").expectedMessageCount(0);
+ getMockEndpoint("mock:result").expectedMessageCount(1);
+ template.sendBody("direct:start", "Hello");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testDisabledPipelineWithSingleOutput() throws Exception {
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("direct:start").pipeline().disabled().to("mock:a").end().to("mock:result");
+ }
+ });
+ getMockEndpoint("mock:a").expectedMessageCount(0);
+ getMockEndpoint("mock:result").expectedMessageCount(1);
+ template.sendBody("direct:start", "Hello");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testPipelineWithSingleOutput() throws Exception {
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("direct:start").pipeline().to("mock:a").end().to("mock:result");
+ }
+ });
+ getMockEndpoint("mock:a").expectedMessageCount(1);
+ getMockEndpoint("mock:result").expectedMessageCount(1);
+ template.sendBody("direct:start", "Hello");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testErrorHandlerWithUnknownExecutorService() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
errorHandler(defaultErrorHandler().executorServiceRef("unknown"));
+
+ from("direct:start").to("mock:result");
+ }
+ });
+ assertThatThrownBy(() -> context.start())
+ .isInstanceOf(FailedToCreateRouteException.class)
+ .rootCause().isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("ExecutorService unknown not found in
registry");
+ }
+
+ @Test
+ public void testDelayWithExecutorServiceProfileFromPlaceholder() throws
Exception {
+ ThreadPoolProfile profile = new ThreadPoolProfile("myPool");
+ profile.setPoolSize(1);
+ context.getExecutorServiceManager().registerThreadPoolProfile(profile);
+ Properties prop = new Properties();
+ prop.put("pool", "myPool");
+ context.getPropertiesComponent().setInitialProperties(prop);
+
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("direct:start").delay(10).asyncDelayed().executorService("{{pool}}").to("mock:result");
+ }
+ });
+ getMockEndpoint("mock:result").expectedMessageCount(1);
+ template.sendBody("direct:start", "Hello");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testOnExceptionRedeliveryInheritsErrorHandlerDelay() throws
Exception {
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+ errorHandler(defaultErrorHandler().redeliveryDelay(0));
+
onException(IOException.class).maximumRedeliveries(3).handled(true).to("mock:dead");
+
+ from("direct:start").throwException(new IOException("Forced"));
+ }
+ });
+ getMockEndpoint("mock:dead").expectedMessageCount(1);
+ StopWatch watch = new StopWatch();
+ template.sendBody("direct:start", "Hello");
+ assertMockEndpointsSatisfied();
+ // 3 redeliveries with the delay of the error handler (0) and not the
default (1000)
+ assertThat(watch.taken()).isLessThan(1000);
+ }
+
+ @Test
+ public void testSplitAggregationStrategyMethodNameWithPlaceholder() throws
Exception {
+ Properties prop = new Properties();
+ prop.put("method", "join");
+ context.getPropertiesComponent().setInitialProperties(prop);
+ context.getRegistry().bind("myJoiner", new MyJoiner());
+
+ start(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+
.split(body().tokenize(",")).aggregationStrategy("myJoiner")
+ .aggregationStrategyMethodName("{{method}}")
+ .to("mock:split")
+ .end()
+ .to("mock:result");
+ }
+ });
+ getMockEndpoint("mock:result").expectedBodiesReceived("A+B");
+ template.sendBody("direct:start", "A,B");
+ assertMockEndpointsSatisfied();
+ }
+
+ public static class MyJoiner {
+ public String join(String existing, String next) {
+ return existing == null ? next : existing + "+" + next;
+ }
+ }
+}
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 ef160f30c70a..d3aedd21274e 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
@@ -301,6 +301,13 @@ The `inflight` and `blocked` developer consoles now carry
a `nodeSource` field p
exchange sits at is in the source (such as `orders.camel.yaml:18`). It is
`null` when message history or source
location is disabled. Existing fields are unchanged.
+=== Error handler - redelivery options of onException
+
+An `onException` that sets some of the redelivery options (such as
`maximumRedeliveries`) now inherits the options it
+does not set from the error handler, as documented. Prior to Camel 4.23 (since
Camel 3.17) the options it did not set
+were reset to their defaults: for example with
`errorHandler(defaultErrorHandler().redeliveryDelay(0))` and
+`onException(IOException.class).maximumRedeliveries(3)`, the redeliveries were
delayed 1 second each and not 0.
+
=== camel-management
The JMX `browse` operation of the `DefaultInflightRepository` MBean and the
`listAwaitThreads` data of the