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 db15e9eaf3ca CAMEL-25140: camel-core - interceptSendToEndpoint: the
interceptors are registered per route on an endpoint wrapped once (#27081)
db15e9eaf3ca is described below
commit db15e9eaf3ca03b8a1c536880996ef0b875e1127
Author: Claus Ibsen <[email protected]>
AuthorDate: Wed Sep 30 12:42:28 2026 +0200
CAMEL-25140: camel-core - interceptSendToEndpoint: the interceptors are
registered per route on an endpoint wrapped once (#27081)
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
---
.../org/apache/camel/catalog/docs/intercept.adoc | 12 ++
.../mock/InterceptSendToMockEndpointStrategy.java | 20 +-
.../org/apache/camel/impl/engine/DefaultRoute.java | 11 ++
.../main/docs/modules/eips/pages/intercept.adoc | 12 ++
.../processor/InterceptSendToEndpointCallback.java | 4 +
.../processor/InterceptSendToEndpointManager.java | 214 +++++++++++++++++++++
.../InterceptSendToEndpointProcessor.java | 144 +++++++++++---
.../processor/InterceptSendToEndpointService.java | 70 +++++++
.../reifier/InterceptSendToEndpointReifier.java | 39 ++--
.../InterceptSendToEndpointRouteLifecycleTest.java | 194 +++++++++++++++++++
.../support/DefaultInterceptSendToEndpoint.java | 43 +++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 26 ++-
12 files changed, 728 insertions(+), 61 deletions(-)
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/intercept.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/intercept.adoc
index 02bb30045552..8211ded75f3c 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/intercept.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/intercept.adoc
@@ -693,6 +693,18 @@ YAML::
----
====
+=== Several routes and interceptors on the same endpoint
+
+An `interceptSendToEndpoint` in a `RouteBuilder` (or a route configuration)
applies to each of its routes, and each route
+uses its own interceptor when it sends to the endpoint. The interceptor of a
route is active while the route is running,
+so stopping or removing a route does not affect the other routes. When the
endpoint is sent to from somewhere that has no
+interceptor of its own (such as another route or a `ProducerTemplate`), the
interceptor of the first route that
+registered one is used.
+
+When several interceptors apply to the same endpoint, they run in a fixed
order: first the interceptors of the route
+that is sending, and then the interceptor of the endpoint itself, such as when
mocking endpoints with
+`mockEndpoints`.
+
== Intercepting endpoints using pattern matching
The `interceptFrom` and `interceptSendToEndpoint` support endpoint pattern
diff --git
a/components/camel-mock/src/main/java/org/apache/camel/component/mock/InterceptSendToMockEndpointStrategy.java
b/components/camel-mock/src/main/java/org/apache/camel/component/mock/InterceptSendToMockEndpointStrategy.java
index 3a092da3f189..5ad05ac62320 100644
---
a/components/camel-mock/src/main/java/org/apache/camel/component/mock/InterceptSendToMockEndpointStrategy.java
+++
b/components/camel-mock/src/main/java/org/apache/camel/component/mock/InterceptSendToMockEndpointStrategy.java
@@ -24,6 +24,7 @@ import org.apache.camel.spi.EndpointStrategy;
import org.apache.camel.spi.InterceptSendToEndpoint;
import org.apache.camel.support.DefaultInterceptSendToEndpoint;
import org.apache.camel.support.EndpointHelper;
+import org.apache.camel.support.service.ServiceHelper;
import org.apache.camel.util.StringHelper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -71,7 +72,24 @@ public class InterceptSendToMockEndpointStrategy implements
EndpointStrategy {
@Override
public Endpoint registerEndpoint(String uri, Endpoint endpoint) {
- if (endpoint instanceof InterceptSendToEndpoint) {
+ if (endpoint instanceof DefaultInterceptSendToEndpoint dise &&
dise.getBefore() == null
+ && dise.getAfter() == null &&
!dise.getEndpointUri().startsWith("mock:")
+ && matchPattern(uri, dise.getOriginalEndpoint(), pattern)) {
+ // endpoint decorated for the interceptors of routes (intercept
send to endpoint EIP), which the mock
+ // interceptor is added to (and runs after the interceptors of the
routes)
+ dise.setSkip(skip);
+ try {
+ Producer producer = createProducer(endpoint.getCamelContext(),
uri, dise);
+ // allow custom logic
+ producer = onInterceptEndpoint(uri,
dise.getOriginalEndpoint(), producer.getEndpoint(), producer);
+ dise.setBefore(producer);
+ // the endpoint may already be started
+ ServiceHelper.startService(producer);
+ } catch (Exception e) {
+ throw new RuntimeCamelException(e);
+ }
+ return dise;
+ } else if (endpoint instanceof InterceptSendToEndpoint) {
// endpoint already decorated
return endpoint;
} else if (endpoint.getEndpointUri().startsWith("mock:")) {
diff --git
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultRoute.java
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultRoute.java
index 4d218f606174..0f7dcf35fd7d 100644
---
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultRoute.java
+++
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultRoute.java
@@ -111,6 +111,8 @@ public class DefaultRoute extends ServiceSupport implements
Route {
private final Map<String, Object> properties = new HashMap<>();
private final List<Service> services = new ArrayList<>();
private final List<Service> servicesToStop = new ArrayList<>();
+ // the services added with addService, which are kept when the services
are gathered again
+ private final List<Service> addedServices = new ArrayList<>();
private final StopWatch stopWatch = new StopWatch(false);
private RouteError routeError;
private Integer startupOrder;
@@ -234,6 +236,12 @@ public class DefaultRoute extends ServiceSupport
implements Route {
services.clear();
// gather all the services for this route
gatherServices(services);
+ // and the services added to the route (such as by the
interceptSendToEndpoint EIP)
+ for (Service service : addedServices) {
+ if (!services.contains(service)) {
+ services.add(service);
+ }
+ }
}
@Override
@@ -246,6 +254,9 @@ public class DefaultRoute extends ServiceSupport implements
Route {
if (!services.contains(service)) {
services.add(service);
}
+ if (!addedServices.contains(service)) {
+ addedServices.add(service);
+ }
}
@Override
diff --git
a/core/camel-core-engine/src/main/docs/modules/eips/pages/intercept.adoc
b/core/camel-core-engine/src/main/docs/modules/eips/pages/intercept.adoc
index 02bb30045552..8211ded75f3c 100644
--- a/core/camel-core-engine/src/main/docs/modules/eips/pages/intercept.adoc
+++ b/core/camel-core-engine/src/main/docs/modules/eips/pages/intercept.adoc
@@ -693,6 +693,18 @@ YAML::
----
====
+=== Several routes and interceptors on the same endpoint
+
+An `interceptSendToEndpoint` in a `RouteBuilder` (or a route configuration)
applies to each of its routes, and each route
+uses its own interceptor when it sends to the endpoint. The interceptor of a
route is active while the route is running,
+so stopping or removing a route does not affect the other routes. When the
endpoint is sent to from somewhere that has no
+interceptor of its own (such as another route or a `ProducerTemplate`), the
interceptor of the first route that
+registered one is used.
+
+When several interceptors apply to the same endpoint, they run in a fixed
order: first the interceptors of the route
+that is sending, and then the interceptor of the endpoint itself, such as when
mocking endpoints with
+`mockEndpoints`.
+
== Intercepting endpoints using pattern matching
The `interceptFrom` and `interceptSendToEndpoint` support endpoint pattern
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointCallback.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointCallback.java
index 536c8e390743..15599144649a 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointCallback.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointCallback.java
@@ -28,7 +28,11 @@ import org.apache.camel.util.URISupport;
/**
* Endpoint strategy used by intercept send to endpoint.
+ *
+ * @deprecated the intercept send to endpoint EIP no longer uses this, but
registers the interceptors of routes with
+ * {@link InterceptSendToEndpointManager}
*/
+@Deprecated(since = "4.23.0")
public class InterceptSendToEndpointCallback implements EndpointStrategy {
private final CamelContext camelContext;
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointManager.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointManager.java
new file mode 100644
index 000000000000..044fea716fbd
--- /dev/null
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointManager.java
@@ -0,0 +1,214 @@
+/*
+ * 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;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Endpoint;
+import org.apache.camel.spi.EndpointStrategy;
+import org.apache.camel.spi.InterceptSendToEndpoint;
+import org.apache.camel.spi.NormalizedEndpointUri;
+import org.apache.camel.support.DefaultInterceptSendToEndpoint;
+import org.apache.camel.support.DefaultInterceptSendToEndpoint.Interceptor;
+import org.apache.camel.support.EndpointHelper;
+import org.apache.camel.support.PluginHelper;
+import org.apache.camel.util.URISupport;
+
+/**
+ * Manages the interceptors of the intercept send to endpoint EIP for a {@link
CamelContext}.
+ * <p/>
+ * Each endpoint that is intercepted is wrapped once in a {@link
DefaultInterceptSendToEndpoint}, and the routes
+ * register and unregister their interceptors on it (when the route starts and
stops). The wrapped endpoint (and the
+ * producers created from it) stay the same, so removing a route does not
affect the other routes that send to the
+ * endpoint.
+ */
+public final class InterceptSendToEndpointManager implements EndpointStrategy {
+
+ private record Registration(String matchUri, Interceptor interceptor) {
+ }
+
+ private final CamelContext camelContext;
+ private final List<Registration> registrations = new
CopyOnWriteArrayList<>();
+ // the uri patterns of the endpoints to wrap (null to wrap all)
+ private final List<String> patterns = new CopyOnWriteArrayList<>();
+ private volatile boolean wrapAll;
+ private final Lock lock = new ReentrantLock();
+
+ private InterceptSendToEndpointManager(CamelContext camelContext) {
+ this.camelContext = camelContext;
+ }
+
+ /**
+ * Gets (or creates) the manager of the CamelContext
+ */
+ public static InterceptSendToEndpointManager getOrCreate(CamelContext
camelContext) {
+ InterceptSendToEndpointManager answer
+ =
camelContext.getCamelContextExtension().getContextPlugin(InterceptSendToEndpointManager.class);
+ if (answer == null) {
+ answer = new InterceptSendToEndpointManager(camelContext);
+
camelContext.getCamelContextExtension().addContextPlugin(InterceptSendToEndpointManager.class,
answer);
+
camelContext.getCamelContextExtension().registerEndpointCallback(answer);
+ }
+ return answer;
+ }
+
+ /**
+ * Adds the uri pattern (null to match all) of the endpoints to wrap, so
the interceptors of routes can be
+ * registered on them. This must be done when the route is created (before
its endpoints are resolved), as the
+ * producers of the route are created from the (wrapped) endpoints.
+ */
+ public void addPattern(String matchUri) {
+ lock.lock();
+ try {
+ if (matchUri == null) {
+ if (wrapAll) {
+ return;
+ }
+ wrapAll = true;
+ } else if (patterns.contains(matchUri)) {
+ return;
+ } else {
+ patterns.add(matchUri);
+ }
+ wrapExistingEndpoints();
+ } finally {
+ lock.unlock();
+ }
+ }
+
+ /**
+ * Registers the interceptor for the endpoints that match the uri (null to
match all)
+ */
+ public void register(String matchUri, Interceptor interceptor) {
+ lock.lock();
+ try {
+ Registration registration = new Registration(matchUri,
interceptor);
+ if (registrations.contains(registration)) {
+ return;
+ }
+ registrations.add(registration);
+ wrapExistingEndpoints();
+ } finally {
+ lock.unlock();
+ }
+ }
+
+ private void wrapExistingEndpoints() {
+ List<Map.Entry<NormalizedEndpointUri, Endpoint>> replace = new
ArrayList<>();
+ for (Map.Entry<NormalizedEndpointUri, Endpoint> entry :
camelContext.getEndpointRegistry().entrySet()) {
+ Endpoint endpoint = entry.getValue();
+ Endpoint answer = registerEndpoint(endpoint.getEndpointUri(),
endpoint);
+ if (answer != endpoint) {
+ replace.add(Map.entry(entry.getKey(), answer));
+ }
+ }
+ for (Map.Entry<NormalizedEndpointUri, Endpoint> entry : replace) {
+ camelContext.getEndpointRegistry().put(entry.getKey(),
entry.getValue());
+ }
+ }
+
+ /**
+ * Unregisters the interceptor. The endpoints stay wrapped (as producers
may use them), but no longer use the
+ * interceptor.
+ */
+ public void unregister(Interceptor interceptor) {
+ lock.lock();
+ try {
+ registrations.removeIf(r -> r.interceptor().equals(interceptor));
+ for (Endpoint endpoint :
camelContext.getEndpointRegistry().values()) {
+ if (endpoint instanceof DefaultInterceptSendToEndpoint dise) {
+ dise.removeInterceptor(interceptor);
+ }
+ }
+ } finally {
+ lock.unlock();
+ }
+ }
+
+ @Override
+ public Endpoint registerEndpoint(String uri, Endpoint endpoint) {
+ if (!wrapAll && patterns.isEmpty()) {
+ return endpoint;
+ }
+ if (endpoint instanceof InterceptSendToEndpoint && !(endpoint
instanceof DefaultInterceptSendToEndpoint)) {
+ // endpoint decorated by a custom interceptor
+ return endpoint;
+ }
+ DefaultInterceptSendToEndpoint wrapped = endpoint instanceof
DefaultInterceptSendToEndpoint dise ? dise : null;
+ if (wrapped == null) {
+ if (!matchesAnyPattern(uri)) {
+ return endpoint;
+ }
+ InterceptSendToEndpoint answer =
PluginHelper.getInterceptEndpointFactory(camelContext)
+ .createInterceptSendToEndpoint(camelContext, endpoint,
false, null, null, null);
+ if (!(answer instanceof DefaultInterceptSendToEndpoint dise)) {
+ // a custom factory that does not support the interceptors of
routes
+ return endpoint;
+ }
+ wrapped = dise;
+ }
+ // add the interceptors of the running routes that match
+ for (Registration registration : registrations) {
+ if (registration.matchUri() == null || matchPattern(uri,
registration.matchUri())) {
+ wrapped.addInterceptor(registration.interceptor());
+ }
+ }
+ return wrapped;
+ }
+
+ private boolean matchesAnyPattern(String uri) {
+ if (wrapAll) {
+ return true;
+ }
+ for (String pattern : patterns) {
+ if (matchPattern(uri, pattern)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ /**
+ * Does the uri match the pattern.
+ *
+ * @param uri the uri
+ * @param pattern the pattern, which can be an endpoint uri as well
+ * @return <tt>true</tt> if matched and we should intercept,
<tt>false</tt> if not matched, and not
+ * intercept.
+ */
+ private boolean matchPattern(String uri, String pattern) {
+ // match using the pattern as-is
+ boolean match = EndpointHelper.matchEndpoint(camelContext, uri,
pattern);
+ if (!match) {
+ try {
+ // the pattern could be an uri, so we need to normalize it
+ // before matching again
+ pattern = URISupport.normalizeUri(pattern);
+ match = EndpointHelper.matchEndpoint(camelContext, uri,
pattern);
+ } catch (Exception e) {
+ // ignore
+ }
+ }
+ return match;
+ }
+}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointProcessor.java
index 1cb59a713bf4..635183ed5072 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointProcessor.java
@@ -16,6 +16,9 @@
*/
package org.apache.camel.processor;
+import java.util.ArrayList;
+import java.util.List;
+
import org.apache.camel.AsyncCallback;
import org.apache.camel.AsyncProcessor;
import org.apache.camel.AsyncProducer;
@@ -24,10 +27,12 @@ import org.apache.camel.Endpoint;
import org.apache.camel.Exchange;
import org.apache.camel.ExchangePropertyKey;
import org.apache.camel.Predicate;
+import org.apache.camel.Route;
import org.apache.camel.spi.InterceptSendToEndpoint;
import org.apache.camel.support.AsyncProcessorConverterHelper;
import org.apache.camel.support.DefaultAsyncProducer;
import org.apache.camel.support.DefaultInterceptSendToEndpoint;
+import org.apache.camel.support.DefaultInterceptSendToEndpoint.Interceptor;
import org.apache.camel.support.ExchangeHelper;
import org.apache.camel.support.service.ServiceHelper;
import org.slf4j.Logger;
@@ -50,6 +55,7 @@ public class InterceptSendToEndpointProcessor extends
DefaultAsyncProducer {
private final Predicate onWhen;
private AsyncProcessor pipeline;
private AsyncProcessor after;
+ private Interceptor interceptor;
public InterceptSendToEndpointProcessor(InterceptSendToEndpoint endpoint,
Endpoint delegate, AsyncProducer producer,
boolean skip, Predicate onWhen) {
@@ -74,50 +80,120 @@ public class InterceptSendToEndpointProcessor extends
DefaultAsyncProducer {
endpoint.getBefore(), exchange);
}
exchange.setProperty(ExchangePropertyKey.INTERCEPTED_ENDPOINT,
delegate.getEndpointUri());
- return pipeline.process(exchange, doneSync -> callback(exchange,
callback, doneSync));
+
+ List<Interceptor> chain = chain(exchange);
+ if (chain.isEmpty()) {
+ // no interceptor (anymore) so send to the endpoint
+ return producer.process(exchange, callback);
+ }
+ return process(exchange, chain, 0, callback);
}
- private boolean callback(Exchange exchange, AsyncCallback callback,
boolean doneSync) {
- // Decide whether to continue or not; similar logic to the Pipeline
- // check for error if so we should break out
- if (!continueProcessing(exchange, "skip sending to original intended
destination: " + getEndpoint(), LOG)) {
- callback.done(doneSync);
- return doneSync;
+ /**
+ * The interceptors to use in their fixed order: the interceptors of the
route that is sending (or when the route
+ * has none, the interceptors of the first route that registered one), and
then the interceptor of the endpoint
+ * itself (such as a mock).
+ */
+ private List<Interceptor> chain(Exchange exchange) {
+ List<Interceptor> routes = routeInterceptors(exchange);
+ if (routes.isEmpty()) {
+ return this.interceptor != null ? List.of(this.interceptor) :
List.of();
+ }
+ if (this.interceptor == null) {
+ return routes;
}
+ List<Interceptor> answer = new ArrayList<>(routes.size() + 1);
+ answer.addAll(routes);
+ answer.add(this.interceptor);
+ return answer;
+ }
- // determine if we should skip or not
- boolean shouldSkip = skip;
+ private List<Interceptor> routeInterceptors(Exchange exchange) {
+ if (!(endpoint instanceof DefaultInterceptSendToEndpoint dise)) {
+ return List.of();
+ }
+ List<Interceptor> all = dise.getInterceptors();
+ if (all.isEmpty()) {
+ return List.of();
+ }
+ Route route = ExchangeHelper.getRoute(exchange);
+ String routeId = route != null ? route.getRouteId() : null;
+ List<Interceptor> answer = interceptorsOfRoute(all, routeId);
+ if (answer.isEmpty()) {
+ // the route that is sending has no interceptor (or it is not sent
from a route)
+ // so use the interceptors of the first route that registered one
+ answer = interceptorsOfRoute(all, all.get(0).routeId());
+ }
+ return answer;
+ }
- // if then interceptor has predicate, then we should only skip if
matched
- Boolean whenMatches = (Boolean)
exchange.getProperty(ExchangePropertyKey.INTERCEPT_SEND_TO_ENDPOINT_WHEN_MATCHED);
- if (whenMatches != null) {
- shouldSkip = skip && whenMatches;
+ private static List<Interceptor> interceptorsOfRoute(List<Interceptor>
all, String routeId) {
+ if (routeId == null) {
+ return List.of();
+ }
+ List<Interceptor> answer = null;
+ for (Interceptor i : all) {
+ if (routeId.equals(i.routeId())) {
+ if (answer == null) {
+ answer = new ArrayList<>(1);
+ }
+ answer.add(i);
+ }
}
+ return answer != null ? answer : List.of();
+ }
- if (!shouldSkip) {
- ExchangeHelper.prepareOutToIn(exchange);
+ private boolean process(Exchange exchange, List<Interceptor> chain, int
index, AsyncCallback callback) {
+ if (index == chain.size()) {
+ // route to original destination
+ return producer.process(exchange, callback);
+ }
+ Interceptor current = chain.get(index);
+ if (index > 0) {
+ // each interceptor only sees whether its own onWhen predicate
matched
+
exchange.removeProperty(ExchangePropertyKey.INTERCEPT_SEND_TO_ENDPOINT_WHEN_MATCHED);
+ }
+ AsyncProcessor before =
AsyncProcessorConverterHelper.convert(current.before());
+ return before.process(exchange, doneSync -> afterBefore(exchange,
chain, index, current, callback, doneSync));
+ }
- AsyncCallback ac1 = doneSync1 -> {
-
exchange.removeProperty(ExchangePropertyKey.INTERCEPT_SEND_TO_ENDPOINT_WHEN_MATCHED);
- callback.done(doneSync1);
- };
- AsyncCallback ac2 = null;
- if (after != null && (whenMatches == null || whenMatches)) {
- ac2 = doneSync2 -> after.process(exchange, ac1);
- }
+ private void afterBefore(
+ Exchange exchange, List<Interceptor> chain, int index, Interceptor
current, AsyncCallback callback,
+ boolean doneSync) {
+ // Decide whether to continue or not; similar logic to the Pipeline
+ // check for error if so we should break out
+ if (!continueProcessing(exchange, "skip sending to original intended
destination: " + getEndpoint(), LOG)) {
+ callback.done(doneSync);
+ return;
+ }
- // route to original destination (using producer) and when done,
then
- // optional route to the after processor
- boolean s = producer.process(exchange, ac2 != null ? ac2 : ac1);
- return doneSync && s;
- } else {
+ // if the interceptor has predicate, then we should only skip if
matched
+ Boolean whenMatches = (Boolean)
exchange.getProperty(ExchangePropertyKey.INTERCEPT_SEND_TO_ENDPOINT_WHEN_MATCHED);
+ boolean matched = whenMatches == null || whenMatches;
+ if (current.skip() && matched) {
if (LOG.isDebugEnabled()) {
LOG.debug("Skip sending exchange to original intended
destination: {} for exchange: {}",
getEndpoint(), exchange);
}
callback.done(doneSync);
- return doneSync;
+ return;
+ }
+
+ ExchangeHelper.prepareOutToIn(exchange);
+
+ AsyncCallback ac1 = doneSync1 -> {
+
exchange.removeProperty(ExchangePropertyKey.INTERCEPT_SEND_TO_ENDPOINT_WHEN_MATCHED);
+ callback.done(doneSync1);
+ };
+ AsyncCallback ac2 = null;
+ if (current.after() != null && matched) {
+ AsyncProcessor after =
AsyncProcessorConverterHelper.convert(current.after());
+ ac2 = doneSync2 -> after.process(exchange, ac1);
}
+
+ // route to the next interceptor (or the original destination) and
when done, then
+ // optional route to the after processor
+ process(exchange, chain, index + 1, ac2 != null ? ac2 : ac1);
}
@Override
@@ -129,9 +205,13 @@ public class InterceptSendToEndpointProcessor extends
DefaultAsyncProducer {
protected void doBuild() throws Exception {
CamelContextAware.trySetCamelContext(producer,
endpoint.getCamelContext());
- pipeline = new FilterProcessor(getEndpoint().getCamelContext(),
onWhen, endpoint.getBefore());
- if (endpoint.getAfter() != null) {
- after = AsyncProcessorConverterHelper.convert(endpoint.getAfter());
+ // the interceptor of the endpoint itself (such as a mock)
+ if (endpoint.getBefore() != null || endpoint.getAfter() != null) {
+ pipeline = new FilterProcessor(getEndpoint().getCamelContext(),
onWhen, endpoint.getBefore());
+ if (endpoint.getAfter() != null) {
+ after =
AsyncProcessorConverterHelper.convert(endpoint.getAfter());
+ }
+ interceptor = new Interceptor(null, pipeline, after, skip);
}
ServiceHelper.buildService(producer, pipeline, after);
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointService.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointService.java
new file mode 100644
index 000000000000..65c1046f8237
--- /dev/null
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/InterceptSendToEndpointService.java
@@ -0,0 +1,70 @@
+/*
+ * 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;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.support.DefaultInterceptSendToEndpoint.Interceptor;
+import org.apache.camel.support.service.ServiceHelper;
+import org.apache.camel.support.service.ServiceSupport;
+
+/**
+ * A route service for the intercept send to endpoint EIP, which registers the
interceptor of the route when the route
+ * starts, and unregisters it when the route stops (or is removed).
+ */
+public class InterceptSendToEndpointService extends ServiceSupport {
+
+ private final CamelContext camelContext;
+ private final String matchUri;
+ private final Interceptor interceptor;
+
+ public InterceptSendToEndpointService(CamelContext camelContext, String
matchUri, Interceptor interceptor) {
+ this.camelContext = camelContext;
+ this.matchUri = matchUri;
+ this.interceptor = interceptor;
+ }
+
+ public Interceptor getInterceptor() {
+ return interceptor;
+ }
+
+ @Override
+ protected void doBuild() throws Exception {
+ ServiceHelper.buildService(interceptor.before(), interceptor.after());
+ }
+
+ @Override
+ protected void doInit() throws Exception {
+ ServiceHelper.initService(interceptor.before(), interceptor.after());
+ }
+
+ @Override
+ protected void doStart() throws Exception {
+ ServiceHelper.startService(interceptor.before(), interceptor.after());
+
InterceptSendToEndpointManager.getOrCreate(camelContext).register(matchUri,
interceptor);
+ }
+
+ @Override
+ protected void doStop() throws Exception {
+
InterceptSendToEndpointManager.getOrCreate(camelContext).unregister(interceptor);
+ ServiceHelper.stopService(interceptor.before(), interceptor.after());
+ }
+
+ @Override
+ protected void doShutdown() throws Exception {
+ ServiceHelper.stopAndShutdownServices(interceptor.before(),
interceptor.after());
+ }
+}
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/InterceptSendToEndpointReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/InterceptSendToEndpointReifier.java
index 500c6935a09d..6b7a22dae543 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/InterceptSendToEndpointReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/InterceptSendToEndpointReifier.java
@@ -16,8 +16,6 @@
*/
package org.apache.camel.reifier;
-import java.util.List;
-
import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
import org.apache.camel.ExchangePropertyKey;
@@ -26,10 +24,12 @@ import org.apache.camel.Processor;
import org.apache.camel.Route;
import org.apache.camel.model.InterceptSendToEndpointDefinition;
import org.apache.camel.model.ProcessorDefinition;
-import org.apache.camel.model.RouteDefinition;
import org.apache.camel.model.ToDefinition;
-import org.apache.camel.processor.InterceptSendToEndpointCallback;
+import org.apache.camel.processor.FilterProcessor;
+import org.apache.camel.processor.InterceptSendToEndpointManager;
+import org.apache.camel.processor.InterceptSendToEndpointService;
import org.apache.camel.processor.Pipeline;
+import org.apache.camel.support.DefaultInterceptSendToEndpoint.Interceptor;
import org.apache.camel.support.ExchangeHelper;
import org.apache.camel.support.PluginHelper;
@@ -67,7 +67,7 @@ public class InterceptSendToEndpointReifier extends
ProcessorReifier<InterceptSe
final Route registeringRoute = route;
Processor p = exchange -> {
- // the endpoint is decorated once (by the first route of the
intercept), so use the route that is sending
+ // other routes may use this interceptor (such as when they have
none), so use the route that is sending
Route current = ExchangeHelper.getRoute(exchange);
if (current == null) {
current = registeringRoute;
@@ -77,24 +77,17 @@ public class InterceptSendToEndpointReifier extends
ProcessorReifier<InterceptSe
exchange.setProperty(ExchangePropertyKey.INTERCEPTED_ROUTE_ENDPOINT_URI,
current.getEndpoint().getEndpointUri());
};
- // register endpoint callback so we can proxy the endpoint
- camelContext.getCamelContextExtension()
- .registerEndpointCallback(
- new InterceptSendToEndpointCallback(
- camelContext,
- Pipeline.newInstance(camelContext, p, before),
- after,
- matchURI, skip, when));
-
- // remove the original intercepted route from the outputs as we do not
- // intercept as the regular interceptor
- // instead we use the proxy endpoints producer do the triggering. That
- // is we trigger when someone sends
- // an exchange to the endpoint, see InterceptSendToEndpoint for
details.
- RouteDefinition route = (RouteDefinition) this.route.getRoute();
- List<ProcessorDefinition<?>> outputs = route.getOutputs();
- outputs.remove(definition);
-
+ // the interceptor of this route, which it registers when it starts
and unregisters when it stops, so the
+ // endpoints are intercepted by the routes that are running (see
InterceptSendToEndpointManager)
+ Predicate predicate = when != null ? when : exchange -> true;
+ Processor pipeline = new FilterProcessor(camelContext, predicate,
Pipeline.newInstance(camelContext, p, before));
+ Interceptor interceptor = new Interceptor(route.getRouteId(),
pipeline, after, skip);
+ // the matching endpoints must be wrapped now (before the route
resolves its endpoints)
+
InterceptSendToEndpointManager.getOrCreate(camelContext).addPattern(matchURI);
+ route.addService(new InterceptSendToEndpointService(camelContext,
matchURI, interceptor));
+
+ // the interceptor is not a processor in the route (the definition is
abstract, and is kept in the route, so the
+ // interceptor is created again when the route is created again, such
as when CamelContext is restarted)
// and return no processor to invoke next from me
return null;
}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/intercept/InterceptSendToEndpointRouteLifecycleTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/intercept/InterceptSendToEndpointRouteLifecycleTest.java
new file mode 100644
index 000000000000..bdc6b4ed5007
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/intercept/InterceptSendToEndpointRouteLifecycleTest.java
@@ -0,0 +1,194 @@
+/*
+ * 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.intercept;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.ExchangePropertyKey;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.InterceptSendToMockEndpointStrategy;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * The interceptSendToEndpoint of routes when routes are stopped, removed and
added again, and when CamelContext is
+ * restarted.
+ */
+public class InterceptSendToEndpointRouteLifecycleTest extends
ContextTestSupport {
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ private static RouteBuilder twoRoutes() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ interceptSendToEndpoint("mock:target").to("mock:intercepted");
+
+ from("direct:a").routeId("a").to("mock:target");
+ from("direct:b").routeId("b").to("mock:target");
+ }
+ };
+ }
+
+ private void assertIntercepted(String uri, String expectedRouteId) throws
Exception {
+ MockEndpoint intercepted = getMockEndpoint("mock:intercepted");
+ MockEndpoint target = getMockEndpoint("mock:target");
+ intercepted.reset();
+ target.reset();
+ intercepted.expectedMessageCount(1);
+ if (expectedRouteId != null) {
+
intercepted.expectedPropertyReceived(ExchangePropertyKey.INTERCEPTED_ROUTE_ID.getName(),
expectedRouteId);
+ }
+ target.expectedMessageCount(1);
+
+ template.sendBody(uri, "Hello");
+
+ MockEndpoint.assertIsSatisfied(intercepted, target);
+ }
+
+ @Test
+ public void testRemoveRoute() throws Exception {
+ context.addRoutes(twoRoutes());
+ context.start();
+ assertIntercepted("direct:a", "a");
+ assertIntercepted("direct:b", "b");
+
+ // removing a route does not affect the other route
+ context.getRouteController().stopRoute("a");
+ context.removeRoute("a");
+ assertIntercepted("direct:b", "b");
+
+ context.getRouteController().stopRoute("b");
+ context.removeRoute("b");
+ context.addRoutes(twoRoutes());
+ assertIntercepted("direct:a", "a");
+ assertIntercepted("direct:b", "b");
+ }
+
+ @Test
+ public void testStopRoute() throws Exception {
+ context.addRoutes(twoRoutes());
+ context.start();
+
+ context.getRouteController().stopRoute("a");
+ assertIntercepted("direct:b", "b");
+
+ context.getRouteController().startRoute("a");
+ assertIntercepted("direct:a", "a");
+ assertIntercepted("direct:b", "b");
+ }
+
+ @Test
+ public void testNoInterceptorWhenAllRoutesRemoved() throws Exception {
+ context.addRoutes(twoRoutes());
+ context.start();
+ assertIntercepted("direct:a", "a");
+
+ context.getRouteController().stopRoute("a");
+ context.getRouteController().stopRoute("b");
+ context.removeRoute("a");
+ context.removeRoute("b");
+
+ // the interceptors of the removed routes are no longer used
+ MockEndpoint intercepted = getMockEndpoint("mock:intercepted");
+ intercepted.reset();
+ intercepted.expectedMessageCount(0);
+ template.sendBody("mock:target", "Hello");
+ intercepted.assertIsSatisfied();
+ }
+
+ @Test
+ public void testProducerTemplateUsesRegisteredInterceptor() throws
Exception {
+ context.addRoutes(twoRoutes());
+ context.start();
+
+ // not sent from a route, so the interceptor of the first route is used
+ assertIntercepted("mock:target", "a");
+ }
+
+ @Test
+ public void testRestartCamelContext() throws Exception {
+ context.addRoutes(twoRoutes());
+ context.start();
+ assertIntercepted("direct:a", "a");
+
+ context.stop();
+ context.start();
+ // the producer template caches producers of the endpoints from before
the restart
+ template.stop();
+ template = context.createProducerTemplate();
+ assertIntercepted("direct:a", "a");
+ assertIntercepted("direct:b", "b");
+ }
+
+ @Test
+ public void testTwoRouteBuilders() throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ interceptSendToEndpoint("mock:target").setHeader("by",
constant("x")).to("mock:intercepted");
+ from("direct:x").routeId("x").to("mock:target");
+ }
+ });
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ interceptSendToEndpoint("mock:target").setHeader("by",
constant("y")).to("mock:intercepted");
+ from("direct:y").routeId("y").to("mock:target");
+ }
+ });
+ context.start();
+
+ // each route uses the interceptor of its own RouteBuilder
+ MockEndpoint intercepted = getMockEndpoint("mock:intercepted");
+ intercepted.expectedHeaderValuesReceivedInAnyOrder("by", "x", "y");
+ template.sendBody("direct:x", "Hello");
+ template.sendBody("direct:y", "Hello");
+ intercepted.assertIsSatisfied();
+ assertEquals("x",
intercepted.getReceivedExchanges().get(0).getMessage().getHeader("by"));
+ }
+
+ @Test
+ public void testMockEndpointsAndIntercept() throws Exception {
+ // mock the endpoint as well (the mock runs after the interceptor of
the route)
+ context.getCamelContextExtension().registerEndpointCallback(
+ new InterceptSendToMockEndpointStrategy("direct:target"));
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
interceptSendToEndpoint("direct:target").setHeader("intercepted",
constant(true)).to("mock:intercepted");
+
+ from("direct:a").routeId("a").to("direct:target");
+ from("direct:target").routeId("target").to("mock:result");
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:intercepted").expectedMessageCount(1);
+ getMockEndpoint("mock:direct:target").expectedMessageCount(1);
+
getMockEndpoint("mock:direct:target").expectedHeaderReceived("intercepted",
true);
+ getMockEndpoint("mock:result").expectedMessageCount(1);
+
+ template.sendBody("direct:a", "Hello");
+
+ assertMockEndpointsSatisfied();
+ }
+}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultInterceptSendToEndpoint.java
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultInterceptSendToEndpoint.java
index e59f7b19d653..9d324dbffdcd 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultInterceptSendToEndpoint.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultInterceptSendToEndpoint.java
@@ -16,7 +16,9 @@
*/
package org.apache.camel.support;
+import java.util.List;
import java.util.Map;
+import java.util.concurrent.CopyOnWriteArrayList;
import org.apache.camel.AsyncProducer;
import org.apache.camel.CamelContext;
@@ -35,9 +37,29 @@ import org.apache.camel.support.service.ServiceHelper;
/**
* This is an endpoint when sending to it, is intercepted and is routed in a
detour (before and optionally after).
+ * <p/>
+ * The endpoint has one interceptor set by {@link #setBefore(Processor)},
{@link #setAfter(Processor)},
+ * {@link #setSkip(boolean)} and {@link #setOnWhen(Predicate)} (such as when
mocking endpoints), and it can have
+ * interceptors that routes register and unregister with {@link
#addInterceptor(Interceptor)} and
+ * {@link #removeInterceptor(Interceptor)} (the interceptSendToEndpoint EIP).
*/
public class DefaultInterceptSendToEndpoint implements
InterceptSendToEndpoint, ShutdownableService {
+ /**
+ * An interceptor that a route has registered on the endpoint.
+ *
+ * @param routeId the id of the route the interceptor belongs to
+ * @param before the processor to route to before sending to the
endpoint, which takes care of the onWhen predicate
+ * (and sets the {@link
org.apache.camel.ExchangePropertyKey#INTERCEPT_SEND_TO_ENDPOINT_WHEN_MATCHED}
+ * property when it has one)
+ * @param after the optional processor to route to after sending to the
endpoint
+ * @param skip whether to skip sending to the endpoint (when the onWhen
predicate matched, if any)
+ */
+ public record Interceptor(String routeId, Processor before, Processor
after, boolean skip) {
+ }
+
+ private final CopyOnWriteArrayList<Interceptor> interceptors = new
CopyOnWriteArrayList<>();
+
private final CamelContext camelContext;
private final Endpoint delegate;
private Predicate onWhen;
@@ -77,6 +99,27 @@ public class DefaultInterceptSendToEndpoint implements
InterceptSendToEndpoint,
this.skip = skip;
}
+ /**
+ * Adds an interceptor that a route registers (if not already added)
+ */
+ public void addInterceptor(Interceptor interceptor) {
+ interceptors.addIfAbsent(interceptor);
+ }
+
+ /**
+ * Removes an interceptor that a route has registered
+ */
+ public void removeInterceptor(Interceptor interceptor) {
+ interceptors.remove(interceptor);
+ }
+
+ /**
+ * The interceptors that routes have registered, in the order they were
added
+ */
+ public List<Interceptor> getInterceptors() {
+ return interceptors;
+ }
+
@Override
public Processor getBefore() {
return before;
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 0e6cac6bcd92..a9d0ab1df4f3 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
@@ -389,11 +389,27 @@ The Event developer console now exposes the full
structured JSON payload in the
each event entry, while keeping the existing flat `type`, `timestamp`,
`exchangeId`, and
`message` fields for backwards compatibility.
-=== Intercept Send To Endpoint - intercepted route
-
-When an `interceptSendToEndpoint` applies to several routes (such as one
defined in a `RouteBuilder` with several
-routes), the `CamelInterceptedRouteId` and `CamelInterceptedParentEndpointUri`
exchange properties are now those of the route
-that sends to the endpoint. Before, they were always those of the first route
of the `RouteBuilder`.
+=== Intercept Send To Endpoint EIP
+
+The interceptors of `interceptSendToEndpoint` are now registered by each route
while it is running, on an endpoint that
+is wrapped once:
+
+- Each route uses its own interceptor. Before, the endpoint was wrapped by the
interceptor of one of the routes (which
+ one depended on the order of an unordered set), and all the routes used that
one.
+- Removing a route no longer breaks the other routes that send to the
intercepted endpoint (they failed with
+ `RejectedExecutionException` when the route whose interceptor wrapped the
endpoint was removed).
+- The interception is kept when the `CamelContext` is restarted.
+- The interceptor of a stopped route is no longer used.
+- Interceptors from two `RouteBuilder`s for the same endpoint are both used,
each by the routes of its own
+ `RouteBuilder`. Before, only the one that wrapped the endpoint first was
used.
+- `mockEndpoints` together with an `interceptSendToEndpoint` on the same
endpoint now both run, in a fixed order: the
+ interceptors of the route that is sending first, and then the mock. Before,
only the one that wrapped the endpoint
+ first was used.
+- When an interceptor applies to several routes (such as one defined in a
`RouteBuilder` with several routes), the
+ `CamelInterceptedRouteId` and `CamelInterceptedParentEndpointUri` exchange
properties are those of the route that
+ sends to the endpoint. Before, they were always those of the first route of
the `RouteBuilder`.
+
+The `org.apache.camel.processor.InterceptSendToEndpointCallback` class is
deprecated, as it is no longer used.
=== Route templates