This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-25067 in repository https://gitbox.apache.org/repos/asf/camel.git
commit 0cbbb90c649d87881e9896204a49abd39a29099d Author: Claus Ibsen <[email protected]> AuthorDate: Tue Sep 29 12:09:50 2026 +0200 CAMEL-25067: camel-core - Reifiers: fix the remaining follow-ups from the deep review - interceptSendToEndpoint: the intercepted route id and route endpoint are those of the route that sends to the endpoint, not of the first route that registered the intercept. - The load balancer and its outputs get their ids and route id, so context.getProcessor(id) finds them. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../reifier/InterceptSendToEndpointReifier.java | 11 +++- .../apache/camel/reifier/LoadBalanceReifier.java | 4 ++ .../org/apache/camel/reifier/ProcessorReifier.java | 55 +++++++---------- .../camel/processor/InterceptPropertiesTest.java | 26 ++++++++ .../apache/camel/processor/LoadBalanceIdTest.java | 72 ++++++++++++++++++++++ .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 6 ++ 6 files changed, 140 insertions(+), 34 deletions(-) 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..500c6935a09d 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 @@ -30,6 +30,7 @@ import org.apache.camel.model.RouteDefinition; import org.apache.camel.model.ToDefinition; import org.apache.camel.processor.InterceptSendToEndpointCallback; import org.apache.camel.processor.Pipeline; +import org.apache.camel.support.ExchangeHelper; import org.apache.camel.support.PluginHelper; public class InterceptSendToEndpointReifier extends ProcessorReifier<InterceptSendToEndpointDefinition> { @@ -64,10 +65,16 @@ 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()); + // the endpoint is decorated once (by the first route of the intercept), 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 diff --git a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/LoadBalanceReifier.java b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/LoadBalanceReifier.java index 5883256fa328..ae7ea103821e 100644 --- a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/LoadBalanceReifier.java +++ b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/LoadBalanceReifier.java @@ -70,10 +70,14 @@ public class LoadBalanceReifier extends ProcessorReifier<LoadBalanceDefinition> + processorType); } Processor processor = createProcessor(processorType); + // the children are not created via createOutputsProcessor, so inject their ids here + injectIds(processor, processorType); Channel channel = wrapChannel(processor, processorType, childInherit); loadBalancer.addProcessor(channel); } + // the load balancer is returned wrapped in a channel, so its id must be injected here + injectIds(loadBalancer, definition); return wrapChannel(loadBalancer, definition, inherit); } diff --git a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java index 71e8fbb5c47a..45396721e88e 100644 --- a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java +++ b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ProcessorReifier.java @@ -760,6 +760,27 @@ public abstract class ProcessorReifier<T extends ProcessorDefinition<?>> extends return errorHandler; } + /** + * Injects the id, route id and step id of the definition into the processor (when it is aware of them). + */ + protected void injectIds(Processor processor, ProcessorDefinition<?> output) { + if (processor instanceof IdAware idAware) { + String id = getId(output); + idAware.setId(id); + } + if (processor instanceof RouteIdAware routeIdAware) { + routeIdAware.setRouteId(route.getRouteId()); + } + if (processor instanceof StepIdAware stepIdAware) { + StepDefinition step = ProcessorDefinitionHelper.findFirstParentOfType( + StepDefinition.class, output, true); + if (step != null) { + stepIdAware.setStepId(step.idOrCreate( + camelContext.getCamelContextExtension().getContextPlugin(NodeIdFactory.class))); + } + } + } + /** * Creates a new instance of some kind of composite processor which defaults to using a {@link Pipeline} but derived * classes could change the behaviour @@ -782,22 +803,7 @@ public abstract class ProcessorReifier<T extends ProcessorDefinition<?>> extends Processor processor = createProcessor(output); - // inject id - if (processor instanceof IdAware idAware) { - String id = getId(output); - idAware.setId(id); - } - if (processor instanceof RouteIdAware routeIdAware) { - routeIdAware.setRouteId(route.getRouteId()); - } - if (processor instanceof StepIdAware stepIdAware) { - StepDefinition step = ProcessorDefinitionHelper.findFirstParentOfType( - StepDefinition.class, output, true); - if (step != null) { - stepIdAware.setStepId(step.idOrCreate( - camelContext.getCamelContextExtension().getContextPlugin(NodeIdFactory.class))); - } - } + injectIds(processor, output); if (output instanceof Channel && processor == null) { continue; @@ -868,22 +874,7 @@ public abstract class ProcessorReifier<T extends ProcessorDefinition<?>> extends processor = createProcessor(); } - // inject id - if (processor instanceof IdAware idAware) { - String id = getId(definition); - idAware.setId(id); - } - if (processor instanceof RouteIdAware routeIdAware) { - routeIdAware.setRouteId(route.getRouteId()); - } - if (processor instanceof StepIdAware stepIdAware) { - StepDefinition step = ProcessorDefinitionHelper.findFirstParentOfType( - StepDefinition.class, definition, true); - if (step != null) { - stepIdAware.setStepId(step.idOrCreate( - camelContext.getCamelContextExtension().getContextPlugin(NodeIdFactory.class))); - } - } + injectIds(processor, definition); if (processor == null) { // no processor to make diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/InterceptPropertiesTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/InterceptPropertiesTest.java index 954b6dd7b386..ec040eda8ef9 100644 --- a/core/camel-core/src/test/java/org/apache/camel/processor/InterceptPropertiesTest.java +++ b/core/camel-core/src/test/java/org/apache/camel/processor/InterceptPropertiesTest.java @@ -108,4 +108,30 @@ public class InterceptPropertiesTest extends ContextTestSupport { assertMockEndpointsSatisfied(); } + @Test + public void testInterceptSendToEndpointPropertiesTwoRoutes() throws Exception { + context.addRoutes(new RouteBuilder() { + @Override + public void configure() throws Exception { + interceptSendToEndpoint("mock:target") + .to("mock:interceptSendToEndpoint"); + + from("direct:a").routeId("a").to("mock:target"); + from("direct:b").routeId("b").to("mock:target"); + } + }); + // the intercepted route is the route that sends to the endpoint + getMockEndpoint("mock:interceptSendToEndpoint").expectedMessageCount(2); + getMockEndpoint("mock:interceptSendToEndpoint") + .expectedPropertyValuesReceivedInAnyOrder(ExchangePropertyKey.INTERCEPTED_ROUTE_ID.getName(), "a", "b"); + getMockEndpoint("mock:interceptSendToEndpoint") + .expectedPropertyValuesReceivedInAnyOrder(ExchangePropertyKey.INTERCEPTED_ROUTE_ENDPOINT_URI.getName(), + "direct://a", "direct://b"); + + template.sendBody("direct:a", "A"); + template.sendBody("direct:b", "B"); + + assertMockEndpointsSatisfied(); + } + } diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/LoadBalanceIdTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/LoadBalanceIdTest.java new file mode 100644 index 000000000000..dbb076777d7d --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/processor/LoadBalanceIdTest.java @@ -0,0 +1,72 @@ +/* + * 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.ContextTestSupport; +import org.apache.camel.Processor; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.processor.loadbalancer.LoadBalancerSupport; +import org.apache.camel.spi.IdAware; +import org.apache.camel.spi.RouteIdAware; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +public class LoadBalanceIdTest extends ContextTestSupport { + + @Test + public void testLoadBalancerAndChildrenIds() { + LoadBalancerSupport lb = context.getProcessor("myBalancer", LoadBalancerSupport.class); + assertNotNull(lb, "the load balancer should be found by its id"); + assertEquals("myRoute", lb.getRouteId()); + + for (String id : new String[] { "toA", "toB", "setC" }) { + Processor child = context.getProcessor(id); + assertNotNull(child, "the load balancer output " + id + " should be found by its id"); + assertEquals(id, assertInstanceOf(IdAware.class, child).getId()); + assertEquals("myRoute", assertInstanceOf(RouteIdAware.class, child).getRouteId()); + } + } + + @Test + public void testLoadBalancerRouting() throws Exception { + getMockEndpoint("mock:a").expectedBodiesReceived("Hello"); + getMockEndpoint("mock:b").expectedBodiesReceived("World"); + + template.sendBody("direct:start", "Hello"); + template.sendBody("direct:start", "World"); + + assertMockEndpointsSatisfied(); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("direct:start").routeId("myRoute") + .loadBalance().roundRobin().id("myBalancer") + .to("mock:a").id("toA") + .to("mock:b").id("toB") + .setHeader("c", constant("C")).id("setC") + .end(); + } + }; + } +} 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..f7d9ded8cf78 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,12 @@ 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 `CamelInterceptedRouteEndpoint` 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`. + === Route templates - When both the route template (`configure`) and the `TemplatedRouteBuilder` (`configure`) have a configurer, then
