This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-25140 in repository https://gitbox.apache.org/repos/asf/camel.git
commit cf3ba08500629ef0580ef60d66a812be83e6d090 Author: Claus Ibsen <[email protected]> AuthorDate: Tue Sep 29 13:06:30 2026 +0200 CAMEL-25140: camel-core - interceptSendToEndpoint: the interceptors are registered per route on an endpoint wrapped once An endpoint that is intercepted is wrapped once, and each route registers its interceptor on it while the route is running (InterceptSendToEndpointService, InterceptSendToEndpointManager). When sending, the interceptors of the sending route are used (or of the first route that registered one, as before), followed by the interceptor of the endpoint itself (such as a mock). - Removing a route no longer breaks the other routes that send to the intercepted endpoint. - The interception is kept when CamelContext is restarted (the definition stays in the route). - Interceptors from two RouteBuilders are each used by their own routes, and mockEndpoints together with an interceptSendToEndpoint both run. - A service added with Route.addService is kept when the route gathers its services again. - InterceptSendToEndpointCallback is deprecated as it is no longer used. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../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 | 48 ++--- .../InterceptSendToEndpointRouteLifecycleTest.java | 194 +++++++++++++++++++ .../support/DefaultInterceptSendToEndpoint.java | 43 +++++ .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 19 ++ 11 files changed, 722 insertions(+), 57 deletions(-) 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 c657dd41913f..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,13 @@ 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; public class InterceptSendToEndpointReifier extends ProcessorReifier<InterceptSendToEndpointDefinition> { @@ -64,30 +65,29 @@ public class InterceptSendToEndpointReifier extends ProcessorReifier<InterceptSe when = new OnWhenPredicate(createPredicate(definition.getOnWhen().getExpression())); } + final Route registeringRoute = route; Processor p = exchange -> { - exchange.setProperty(ExchangePropertyKey.INTERCEPTED_ROUTE_ID, route.getId()); + // 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; + } + exchange.setProperty(ExchangePropertyKey.INTERCEPTED_ROUTE_ID, current.getId()); exchange.setProperty(ExchangePropertyKey.INTERCEPTED_NODE_ID, definition.getId()); - exchange.setProperty(ExchangePropertyKey.INTERCEPTED_ROUTE_ENDPOINT_URI, route.getEndpoint().getEndpointUri()); + 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 04d6027f2ded..287796e11c6e 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 @@ -356,6 +356,25 @@ 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 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. + +The `org.apache.camel.processor.InterceptSendToEndpointCallback` class is deprecated, as it is no longer used. + === Route templates - When both the route template (`configure`) and the `TemplatedRouteBuilder` (`configure`) have a configurer, then
