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 cf526fbe7625 CAMEL-25108: camel-core - Component, endpoint and
consumer base classes: fix bugs found in a deep review (#27014)
cf526fbe7625 is described below
commit cf526fbe76257ce2a57700528436193fc9a470db
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 22:44:17 2026 +0200
CAMEL-25108: camel-core - Component, endpoint and consumer base classes:
fix bugs found in a deep review (#27014)
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../camel/support/BaseClassesEdgeCasesTest.java | 188 +++++++++++++++++++++
.../BridgeExceptionHandlerToErrorHandler.java | 16 ++
.../org/apache/camel/support/DefaultComponent.java | 23 ++-
.../org/apache/camel/support/DefaultEndpoint.java | 9 +-
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 14 ++
5 files changed, 240 insertions(+), 10 deletions(-)
diff --git
a/core/camel-core/src/test/java/org/apache/camel/support/BaseClassesEdgeCasesTest.java
b/core/camel-core/src/test/java/org/apache/camel/support/BaseClassesEdgeCasesTest.java
new file mode 100644
index 000000000000..aa64c353541c
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/support/BaseClassesEdgeCasesTest.java
@@ -0,0 +1,188 @@
+/*
+ * 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.support;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePattern;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.ExceptionHandler;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class BaseClassesEdgeCasesTest extends ContextTestSupport {
+
+ private final AtomicInteger onException = new AtomicInteger();
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ private static class MyComponent extends DefaultComponent {
+ private Map<String, Object> parameters;
+
+ @Override
+ protected Endpoint createEndpoint(String uri, String remaining,
Map<String, Object> parameters) {
+ this.parameters = new HashMap<>(parameters);
+ parameters.clear();
+ return new MyEndpoint(uri, this);
+ }
+ }
+
+ private static class MyEndpoint extends DefaultEndpoint {
+ MyEndpoint(String uri, DefaultComponent component) {
+ super(uri, component);
+ }
+
+ MyEndpoint() {
+ }
+
+ @Override
+ public Producer createProducer() {
+ return null;
+ }
+
+ @Override
+ public Consumer createConsumer(Processor processor) {
+ return null;
+ }
+
+ @Override
+ protected String createEndpointUri() {
+ return "my:endpoint";
+ }
+ }
+
+ @Test
+ public void testBridgeErrorHandlerRouteFailureHandledOnce() throws
Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ onException(IllegalStateException.class).process(e ->
onException.incrementAndGet());
+
+ from("timer:foo?repeatCount=1&delay=1&bridgeErrorHandler=true")
+ .throwException(new IllegalStateException("Forced"));
+ }
+ });
+ context.start();
+
+ await().atMost(5, TimeUnit.SECONDS).until(() -> onException.get() >=
1);
+ // give a potential second (bridged) error handling time to happen
+ await().pollDelay(500, TimeUnit.MILLISECONDS).atMost(2,
TimeUnit.SECONDS).until(() -> true);
+ assertEquals(1, onException.get());
+ }
+
+ @Test
+ public void testBridgeErrorHandlerLaterErrorOfHandledExchange() throws
Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
onException(IllegalStateException.class).handled(true).process(e ->
onException.incrementAndGet());
+
+
from("direct:start?bridgeErrorHandler=true").routeId("foo").to("mock:result");
+ }
+ });
+ context.start();
+
+ DefaultConsumer consumer = (DefaultConsumer)
context.getRoute("foo").getConsumer();
+ ExceptionHandler bridge = consumer.getExceptionHandler();
+ assertInstanceOf(BridgeExceptionHandlerToErrorHandler.class, bridge);
+
+ // a later error (such as when committing) of an exchange that was
handled by the error handler is bridged
+ Exchange handled = consumer.getEndpoint().createExchange();
+ handled.getExchangeExtension().setErrorHandlerHandled(true);
+ bridge.handleException("Error committing", handled, new
IllegalStateException("Commit failed"));
+ assertEquals(1, onException.get());
+
+ // an exchange that failed while routed (not handled by the error
handler) is not bridged again
+ Exchange failed = consumer.getEndpoint().createExchange();
+ failed.getExchangeExtension().setErrorHandlerHandled(false);
+ bridge.handleException("Error processing exchange", failed, new
IllegalStateException("Forced"));
+ assertEquals(1, onException.get());
+ }
+
+ @Test
+ public void testHashParameterIsKept() throws Exception {
+ MyComponent component = new MyComponent();
+ context.addComponent("my", component);
+ context.start();
+
+ context.getEndpoint("my:foo?hash=abc&x=1");
+ assertEquals("abc", component.parameters.get("hash"));
+ assertEquals("1", component.parameters.get("x"));
+ }
+
+ @Test
+ public void testBridgeErrorHandlerHasPrecedenceOverExceptionHandler()
throws Exception {
+ ExceptionHandler handler = new LoggingExceptionHandler(context,
getClass());
+ context.getRegistry().bind("myHandler", handler);
+ context.start();
+
+ Endpoint endpoint =
context.getEndpoint("timer:bar?bridgeErrorHandler=true&exceptionHandler=#myHandler");
+ DefaultConsumer consumer = (DefaultConsumer) endpoint.createConsumer(e
-> {
+ });
+ assertInstanceOf(BridgeExceptionHandlerToErrorHandler.class,
consumer.getExceptionHandler());
+ }
+
+ @Test
+ public void testEndpointWithoutComponent() throws Exception {
+ MyEndpoint endpoint = new MyEndpoint();
+ endpoint.setCamelContext(context);
+ Map<String, Object> props = new HashMap<>();
+ props.put("exchangePattern", "InOut");
+ assertDoesNotThrow(() -> endpoint.configureProperties(props));
+ assertEquals(ExchangePattern.InOut, endpoint.getExchangePattern());
+ }
+
+ @Test
+ public void testReferenceParameterDefaultValue() throws Exception {
+ MyComponent component = new MyComponent();
+ component.setCamelContext(context);
+ Map<String, Object> parameters = new HashMap<>();
+ parameters.put("foo", new Object());
+ assertEquals(5,
component.getAndRemoveOrResolveReferenceParameter(parameters, "foo",
Integer.class, 5));
+ }
+
+ @Test
+ public void testReferenceListParameterAsList() throws Exception {
+ context.getRegistry().bind("a", "A");
+ context.getRegistry().bind("b", "B");
+ MyComponent component = new MyComponent();
+ component.setCamelContext(context);
+ Map<String, Object> parameters = new HashMap<>();
+ parameters.put("foo", List.of("#a", "#b"));
+ List<String> list =
component.resolveAndRemoveReferenceListParameter(parameters, "foo",
String.class, null);
+ assertEquals(List.of("A", "B"), list);
+ assertTrue(parameters.isEmpty());
+ }
+}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/BridgeExceptionHandlerToErrorHandler.java
b/core/camel-support/src/main/java/org/apache/camel/support/BridgeExceptionHandlerToErrorHandler.java
index 210732e29fe7..de3681eced31 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/BridgeExceptionHandlerToErrorHandler.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/BridgeExceptionHandlerToErrorHandler.java
@@ -58,6 +58,14 @@ public class BridgeExceptionHandlerToErrorHandler implements
ExceptionHandler {
@Override
public void handleException(String message, Exchange exchange, Throwable
exception) {
+ if (exchange != null && isFailedByErrorHandler(exchange)) {
+ // the exchange failed while being routed and has already been
processed by the error handler (which did not
+ // handle it), so it must not be processed by the error handler
again (the bridge is only for errors when the
+ // consumer picks up messages, or for later errors such as when
committing a message that was handled)
+ fallback.handleException(message, exchange, exception);
+ return;
+ }
+
Exchange copy;
if (exchange == null) {
copy = consumer.getEndpoint().createExchange();
@@ -90,4 +98,12 @@ public class BridgeExceptionHandlerToErrorHandler implements
ExceptionHandler {
UnitOfWorkHelper.doneUow(uow, copy);
}
}
+
+ private static boolean isFailedByErrorHandler(Exchange exchange) {
+ if (exchange.getExchangeExtension().isRedeliveryExhausted()) {
+ return true;
+ }
+ return exchange.getExchangeExtension().isErrorHandlerHandledSet()
+ && !exchange.getExchangeExtension().isErrorHandlerHandled();
+ }
}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultComponent.java
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultComponent.java
index cec14989670e..daa8d95e015b 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultComponent.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultComponent.java
@@ -123,11 +123,12 @@ public abstract class DefaultComponent extends
ServiceSupport implements Compone
// and use method parseParameters
parameters = URISupport.parseParameters(u);
}
- if (properties != null) {
+ if (properties != null && !properties.isEmpty()) {
parameters.putAll(properties);
+ // This special property (added by endpoint-dsl together with the
properties) is only to identify
+ // endpoints in a unique manner (a hash parameter in the uri
itself is a regular parameter)
+ parameters.remove("hash");
}
- // This special property is only to identify endpoints in a unique
manner
- parameters.remove("hash");
if (resolveRawParameterValues()) {
// parameters using raw syntax: RAW(value)
@@ -564,7 +565,8 @@ public abstract class DefaultComponent extends
ServiceSupport implements Compone
if (EndpointHelper.isReferenceParameter(str)) {
return
EndpointHelper.resolveReferenceParameter(getCamelContext(), str, type);
} else {
- return getCamelContext().getTypeConverter().convertTo(type,
value);
+ T answer =
getCamelContext().getTypeConverter().convertTo(type, value);
+ return answer != null ? answer : defaultValue;
}
}
}
@@ -653,8 +655,17 @@ public abstract class DefaultComponent extends
ServiceSupport implements Compone
Map<String, Object> parameters, String key, Class<T> elementType,
List<T> defaultValue) {
// the value may already be a list such as when using endpoint-dsl
Object value = getAndRemoveParameter(parameters, key, Object.class);
- if (value instanceof List) {
- return (List<T>) value;
+ if (value instanceof List<?> list) {
+ // the elements may be references (such as #myBean) to resolve
+ List<T> answer = new ArrayList<>(list.size());
+ for (Object element : list) {
+ if (element instanceof String str &&
EndpointHelper.isReferenceParameter(str)) {
+
answer.add(EndpointHelper.resolveReferenceParameter(getCamelContext(), str,
elementType));
+ } else {
+ answer.add((T) element);
+ }
+ }
+ return answer;
}
if (value == null) {
return defaultValue;
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultEndpoint.java
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultEndpoint.java
index 4058eb13aaba..b7994aa758d1 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultEndpoint.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultEndpoint.java
@@ -408,9 +408,10 @@ public abstract class DefaultEndpoint extends
ServiceSupport implements Endpoint
}
PropertyConfigurer configurer = null;
- if (bean instanceof Component) {
+ // an endpoint may be created without a component (such as a bean)
+ if (bean instanceof Component && getComponent() != null) {
configurer = getComponent().getComponentPropertyConfigurer();
- } else if (bean instanceof Endpoint) {
+ } else if (bean instanceof Endpoint && getComponent() != null) {
configurer = getComponent().getEndpointPropertyConfigurer();
} else if (bean instanceof PropertyConfigurerAware
propertyConfigurerAware) {
configurer = propertyConfigurerAware.getPropertyConfigurer(bean);
@@ -475,8 +476,8 @@ public abstract class DefaultEndpoint extends
ServiceSupport implements Endpoint
+ " having their consumer
extend DefaultConsumer. The consumer is a "
+
consumer.getClass().getName() + " class.");
}
- }
- if (exceptionHandler != null) {
+ } else if (exceptionHandler != null) {
+ // not in use when bridgeErrorHandler is enabled
if (consumer instanceof DefaultConsumer defaultConsumer) {
defaultConsumer.setExceptionHandler(exceptionHandler);
}
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 3396da29d08b..3e50f5aa2cba 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
@@ -3453,6 +3453,20 @@ use the same queue name with different consumer options,
such as `seda:foo?multi
`seda:foo?multipleConsumers=true&concurrentConsumers=5`. Previously each
distinct endpoint uri only multicast to its own
consumers, so such consumers silently competed for the messages instead of
each receiving a copy.
+=== camel-core - bridgeErrorHandler no longer handles a failed exchange twice
+
+With `bridgeErrorHandler=true`, an exchange that failed while being routed
(and was already handled by the error
+handler of the route, such as an `onException` that does not mark it as
handled) is no longer bridged to the error
+handler a second time. Previously the `onException` (and dead letter channel)
of the route was invoked twice for the
+same failure, the second time with the message body replaced by the error
message. The bridge error handler is only
+for errors that happen when the consumer picks up messages (as documented),
and such errors are bridged as before.
+
+A query parameter named `hash` in an endpoint uri (such as
`http://host/path?hash=abc`) is now kept as a regular
+parameter. Previously it was always removed (it is only used internally by the
endpoint DSL to identify endpoints).
+
+When both `bridgeErrorHandler=true` and `exceptionHandler` are configured on
an endpoint, then `bridgeErrorHandler` is
+now used, as documented. Previously the `exceptionHandler` was used.
+
=== camel-core - Recipient List releases the producers of recipients it did
not send to
The Recipient List acquires a producer for every recipient before it starts
sending. When it completes before