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 0e0be00e8bfb CAMEL-24993: camel-core - toD, Routing Slip and producer
cache: fix bugs found in a deep review (#26848)
0e0be00e8bfb is described below
commit 0e0be00e8bfb4c9a3f6a9aad6f2573895177b3f3
Author: Claus Ibsen <[email protected]>
AuthorDate: Sun Sep 27 10:10:16 2026 +0200
CAMEL-24993: camel-core - toD, Routing Slip and producer cache: fix bugs
found in a deep review (#26848)
- toD sent the exchange after the pre-processor of a SendDynamicAware
component failed, and the callback was called twice.
- A routing slip (and dynamic router) did not stop a prototype endpoint
(cacheSize=-1) whose step completed asynchronously.
- After DefaultProducerCache.purge(), the next send to the last used
endpoint reused the stopped producer.
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../org/apache/camel/processor/RoutingSlip.java | 9 +-
.../camel/processor/SendDynamicProcessor.java | 7 +-
.../component/bar/KaboomSendDynamicAware.java | 42 ++++++++
.../camel/impl/DefaultProducerCachePurgeTest.java | 82 +++++++++++++++
.../processor/RoutingSlipPrototypeAsyncTest.java | 111 +++++++++++++++++++++
...micSendDynamicAwarePreProcessorFailureTest.java | 107 ++++++++++++++++++++
.../services/org/apache/camel/send-dynamic/kaboom | 18 ++++
.../camel/support/cache/DefaultProducerCache.java | 2 +
8 files changed, 373 insertions(+), 5 deletions(-)
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/RoutingSlip.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/RoutingSlip.java
index 619bd1850365..f4c896037aeb 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/RoutingSlip.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/RoutingSlip.java
@@ -428,12 +428,13 @@ public class RoutingSlip extends BaseProcessorSupport
implements Traceable, IdAw
// cleanup producer after usage
ex.removeProperty(ExchangePropertyKey.SLIP_PRODUCER);
+ // and stop prototype endpoints (also when this step
completed asynchronously)
+ if (prototype) {
+ ServiceHelper.stopAndShutdownService(endpoint);
+ }
+
// we only have to handle async completion of the routing
slip
if (doneSync) {
- // and stop prototype endpoints
- if (prototype) {
- ServiceHelper.stopAndShutdownService(endpoint);
- }
cb.done(true);
return;
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/SendDynamicProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/SendDynamicProcessor.java
index b1e48e7f78c3..bf6dcdd23a12 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/SendDynamicProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/SendDynamicProcessor.java
@@ -248,8 +248,13 @@ public class SendDynamicProcessor extends
BaseProcessorSupport
e.setException(t);
// restore previous MEP
target.setPattern(existingPattern);
- // we failed
+ // stop endpoint if prototype as it is not used
+ if (stopEndpoint) {
+ ServiceHelper.stopAndShutdownService(endpoint);
+ }
+ // we failed, so do not send
c.done(true);
+ return true;
}
LOG.debug(">>>> {} {}", endpoint, e);
diff --git
a/core/camel-core/src/test/java/org/apache/camel/component/bar/KaboomSendDynamicAware.java
b/core/camel-core/src/test/java/org/apache/camel/component/bar/KaboomSendDynamicAware.java
new file mode 100644
index 000000000000..7f9396a4863b
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/component/bar/KaboomSendDynamicAware.java
@@ -0,0 +1,42 @@
+/*
+ * 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.component.bar;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+
+/**
+ * A {@link BarSendDynamicAware} whose pre-processor fails, and which counts
the sends done by its post-processor.
+ */
+public class KaboomSendDynamicAware extends BarSendDynamicAware {
+
+ public static final AtomicInteger SENT = new AtomicInteger();
+
+ @Override
+ public Processor createPreProcessor(Exchange exchange, DynamicAwareEntry
entry) {
+ return e -> {
+ throw new IllegalArgumentException("Forced pre-processor failure");
+ };
+ }
+
+ @Override
+ public Processor createPostProcessor(Exchange exchange, DynamicAwareEntry
entry) {
+ return e -> SENT.incrementAndGet();
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/impl/DefaultProducerCachePurgeTest.java
b/core/camel-core/src/test/java/org/apache/camel/impl/DefaultProducerCachePurgeTest.java
new file mode 100644
index 000000000000..14a8b7a34f9a
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/impl/DefaultProducerCachePurgeTest.java
@@ -0,0 +1,82 @@
+/*
+ * 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.impl;
+
+import java.util.Map;
+
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.support.DefaultProducer;
+import org.apache.camel.support.cache.DefaultProducerCache;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+/**
+ * After the producer cache is purged, a send does not use the (stopped)
producer that was used last.
+ */
+public class DefaultProducerCachePurgeTest extends ContextTestSupport {
+
+ @Test
+ public void testSendAfterPurge() throws Exception {
+ context.addComponent("stoppable", new DefaultComponent() {
+ @Override
+ protected Endpoint createEndpoint(String uri, String remaining,
Map<String, Object> parameters) {
+ return new DefaultEndpoint(uri, this) {
+ @Override
+ public Producer createProducer() {
+ return new DefaultProducer(this) {
+ @Override
+ public void process(Exchange exchange) {
+ if (!isStarted()) {
+ throw new IllegalStateException("Producer
is stopped");
+ }
+ }
+ };
+ }
+
+ @Override
+ public Consumer createConsumer(Processor processor) {
+ throw new UnsupportedOperationException();
+ }
+ };
+ }
+ });
+ Endpoint endpoint = context.getEndpoint("stoppable:a");
+
+ DefaultProducerCache cache = new DefaultProducerCache(this, context,
0);
+ cache.start();
+ try {
+ Exchange first = cache.send(endpoint, new
DefaultExchange(context), null);
+ assertNull(first.getException());
+
+ cache.purge();
+
+ Exchange second = cache.send(endpoint, new
DefaultExchange(context), null);
+ assertNull(second.getException(), "the stopped producer should not
be used after the purge");
+ } finally {
+ cache.stop();
+ }
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/RoutingSlipPrototypeAsyncTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/RoutingSlipPrototypeAsyncTest.java
new file mode 100644
index 000000000000..152049867a7e
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/RoutingSlipPrototypeAsyncTest.java
@@ -0,0 +1,111 @@
+/*
+ * 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.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.AsyncCallback;
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.support.DefaultAsyncProducer;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * A prototype endpoint (cacheSize=-1) of a routing slip, and its producer,
are stopped after use, also when the step
+ * completes asynchronously.
+ */
+public class RoutingSlipPrototypeAsyncTest extends ContextTestSupport {
+
+ private final AtomicInteger started = new AtomicInteger();
+ private final AtomicInteger stopped = new AtomicInteger();
+ private final AtomicInteger endpointStopped = new AtomicInteger();
+
+ @Test
+ public void testPrototypeEndpointStoppedAfterAsyncStep() throws Exception {
+ getMockEndpoint("mock:result").expectedMessageCount(1);
+
+ template.sendBodyAndHeader("direct:start", "Hello", "slip", "async:a");
+
+ assertMockEndpointsSatisfied();
+ assertEquals(1, started.get());
+ await().atMost(5, TimeUnit.SECONDS).untilAsserted(() ->
assertEquals(1, stopped.get()));
+ // the prototype endpoint itself must also be stopped, as it was only
used once
+ await().atMost(5, TimeUnit.SECONDS).untilAsserted(() ->
assertEquals(1, endpointStopped.get()));
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ context.addComponent("async", new DefaultComponent() {
+ @Override
+ protected Endpoint createEndpoint(String uri, String remaining,
Map<String, Object> parameters) {
+ return new DefaultEndpoint(uri, this) {
+ @Override
+ public Producer createProducer() {
+ return new DefaultAsyncProducer(this) {
+ @Override
+ public boolean process(Exchange exchange,
AsyncCallback callback) {
+ CompletableFuture.runAsync(() ->
callback.done(false));
+ return false;
+ }
+
+ @Override
+ protected void doStart() {
+ started.incrementAndGet();
+ }
+
+ @Override
+ protected void doStop() {
+ stopped.incrementAndGet();
+ }
+ };
+ }
+
+ @Override
+ protected void doStop() throws Exception {
+ endpointStopped.incrementAndGet();
+ super.doStop();
+ }
+
+ @Override
+ public Consumer createConsumer(Processor processor) {
+ throw new UnsupportedOperationException();
+ }
+ };
+ }
+ });
+
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("direct:start").routingSlip(header("slip")).cacheSize(-1).end().to("mock:result");
+ }
+ };
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/ToDynamicSendDynamicAwarePreProcessorFailureTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/ToDynamicSendDynamicAwarePreProcessorFailureTest.java
new file mode 100644
index 000000000000..e04c7fd77ccf
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/ToDynamicSendDynamicAwarePreProcessorFailureTest.java
@@ -0,0 +1,107 @@
+/*
+ * 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.Map;
+import java.util.Set;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.CamelExecutionException;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.bar.BarComponent;
+import org.apache.camel.component.bar.KaboomSendDynamicAware;
+import org.apache.camel.support.component.EndpointUriFactorySupport;
+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.assertThrows;
+
+/**
+ * When the pre-processor of a SendDynamicAware fails, toD fails the exchange
and does not send it.
+ */
+public class ToDynamicSendDynamicAwarePreProcessorFailureTest extends
ContextTestSupport {
+
+ private final AtomicInteger after = new AtomicInteger();
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext context = super.createCamelContext();
+ context.getRegistry().bind("kaboomFactory", new
KaboomEndpointUriFactory());
+ return context;
+ }
+
+ @Test
+ public void testPreProcessorFailureDoesNotSend() {
+ KaboomSendDynamicAware.SENT.set(0);
+
+ CamelExecutionException e = assertThrows(CamelExecutionException.class,
+ () -> template.sendBodyAndHeader("direct:start", "Hello
Camel", "drink", "beer"));
+ assertInstanceOf(IllegalArgumentException.class, e.getCause());
+
+ assertEquals(0, KaboomSendDynamicAware.SENT.get(), "the exchange
should not be sent");
+ assertEquals(0, after.get(), "the route should not continue after the
failure");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ context.addComponent("kaboom", new BarComponent());
+
+
from("direct:start").toD("kaboom:order?drink=${header.drink}").process(e ->
after.incrementAndGet());
+ }
+ };
+ }
+
+ private static class KaboomEndpointUriFactory extends
EndpointUriFactorySupport {
+
+ @Override
+ public boolean isEnabled(String scheme) {
+ return "kaboom".equals(scheme);
+ }
+
+ @Override
+ public String buildUri(String scheme, Map<String, Object> properties,
boolean encode) {
+ // not in use for this test
+ return null;
+ }
+
+ @Override
+ public Set<String> propertyNames() {
+ return Set.of("name", "drink");
+ }
+
+ @Override
+ public Set<String> secretPropertyNames() {
+ return null;
+ }
+
+ @Override
+ public Map<String, String> multiValuePrefixes() {
+ return null;
+ }
+
+ @Override
+ public boolean isLenientProperties() {
+ return false;
+ }
+ }
+}
diff --git
a/core/camel-core/src/test/resources/META-INF/services/org/apache/camel/send-dynamic/kaboom
b/core/camel-core/src/test/resources/META-INF/services/org/apache/camel/send-dynamic/kaboom
new file mode 100644
index 000000000000..33acc39e51a7
--- /dev/null
+++
b/core/camel-core/src/test/resources/META-INF/services/org/apache/camel/send-dynamic/kaboom
@@ -0,0 +1,18 @@
+#
+# 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.
+#
+
+class=org.apache.camel.component.bar.KaboomSendDynamicAware
\ No newline at end of file
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/cache/DefaultProducerCache.java
b/core/camel-support/src/main/java/org/apache/camel/support/cache/DefaultProducerCache.java
index f61b39591bc6..912afadeecac 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/cache/DefaultProducerCache.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/cache/DefaultProducerCache.java
@@ -381,6 +381,8 @@ public class DefaultProducerCache extends ServiceSupport
implements ProducerCach
try {
if (producers != null) {
producers.stop();
+ // the last used producer has been stopped, so it must not
be reused
+ lastUsedProducer = null;
producers.start();
}
} catch (Exception e) {