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 730e8715eb67 CAMEL-25297: camel-direct, camel-kamelet - do not send to
a suspended or stopped route while another exchange waits for its consumer
(#27338)
730e8715eb67 is described below
commit 730e8715eb67d2fefe4cbca8527817fce75b2e36
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Oct 5 14:10:28 2026 +0530
CAMEL-25297: camel-direct, camel-kamelet - do not send to a suspended or
stopped route while another exchange waits for its consumer (#27338)
DirectProducer cached the consumer in two fields, stateCounter and consumer,
and on a change it set stateCounter first and then looked up the consumer
(which waits for the consumer with block=true, the default). While an
exchange
waited in that lookup, every other exchange sent with the same producer saw
the new stateCounter and the old consumer, and was processed by the route
that
was just suspended or stopped, instead of waiting for its consumer.
KameletProducer has the same code.
Both producers now keep the consumer and the state counter read before the
lookup together in one immutable holder, so another exchange either sees the
old counter (and looks up the consumer too) or the result of the lookup.
Found with a TLA+ model of the consumer cache (invariant: an exchange whose
check starts after the consumer was removed is not sent to that consumer).
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../camel/component/direct/DirectProducer.java | 33 ++++--
.../camel/component/kamelet/KameletProducer.java | 28 ++++-
.../KameletProducerSuspendedConsumerTest.java | 123 +++++++++++++++++++++
.../DirectProducerSuspendedConsumerTest.java | 119 ++++++++++++++++++++
4 files changed, 287 insertions(+), 16 deletions(-)
diff --git
a/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
b/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
index ecd3a9cf70aa..3ff25eeacf7c 100644
---
a/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
+++
b/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
@@ -30,8 +30,9 @@ public class DirectProducer extends DefaultAsyncProducer {
private static final Logger LOG =
LoggerFactory.getLogger(DirectProducer.class);
- private volatile DirectConsumer consumer;
- private int stateCounter;
+ // the consumer and the state counter of the component when the consumer
was looked up, kept together as several
+ // threads can send with this producer at the same time
+ private volatile CachedConsumer cachedConsumer;
private final DirectEndpoint endpoint;
private final DirectComponent component;
@@ -50,10 +51,7 @@ public class DirectProducer extends DefaultAsyncProducer {
@Override
public void process(Exchange exchange) throws Exception {
- if (consumer == null || stateCounter != component.getStateCounter()) {
- stateCounter = component.getStateCounter();
- consumer = component.getConsumer(key, block, timeout);
- }
+ DirectConsumer consumer = getConsumer();
if (consumer == null) {
if (endpoint.isFailIfNoConsumers()) {
throw new DirectConsumerNotAvailableException("No consumers
available on endpoint: " + endpoint, exchange);
@@ -74,10 +72,7 @@ public class DirectProducer extends DefaultAsyncProducer {
callback.done(true);
return true;
}
- if (consumer == null || stateCounter !=
component.getStateCounter()) {
- stateCounter = component.getStateCounter();
- consumer = component.getConsumer(key, block, timeout);
- }
+ DirectConsumer consumer = getConsumer();
if (consumer == null) {
if (endpoint.isFailIfNoConsumers()) {
exchange.setException(new
DirectConsumerNotAvailableException(
@@ -116,4 +111,22 @@ public class DirectProducer extends DefaultAsyncProducer {
}
}
+ /**
+ * Gets the consumer, which is looked up again when it has been added or
removed (such as when its route is
+ * suspended or stopped) since it was looked up last.
+ */
+ private DirectConsumer getConsumer() throws InterruptedException {
+ CachedConsumer cached = cachedConsumer;
+ // read the counter before the lookup, so a change during the lookup
makes the next exchange look up again
+ int stateCounter = component.getStateCounter();
+ if (cached == null || cached.consumer() == null ||
cached.stateCounter() != stateCounter) {
+ DirectConsumer consumer = component.getConsumer(key, block,
timeout);
+ cachedConsumer = new CachedConsumer(consumer, stateCounter);
+ return consumer;
+ }
+ return cached.consumer();
+ }
+
+ private record CachedConsumer(DirectConsumer consumer, int stateCounter) {
+ }
}
diff --git
a/components/camel-kamelet/src/main/java/org/apache/camel/component/kamelet/KameletProducer.java
b/components/camel-kamelet/src/main/java/org/apache/camel/component/kamelet/KameletProducer.java
index 1d219d524d34..3715297bc09f 100644
---
a/components/camel-kamelet/src/main/java/org/apache/camel/component/kamelet/KameletProducer.java
+++
b/components/camel-kamelet/src/main/java/org/apache/camel/component/kamelet/KameletProducer.java
@@ -31,8 +31,9 @@ final class KameletProducer extends DefaultAsyncProducer
implements RouteIdAware
private static final Logger LOG =
LoggerFactory.getLogger(KameletProducer.class);
- private volatile KameletConsumer consumer;
- private int stateCounter;
+ // the consumer and the state counter of the component when the consumer
was looked up, kept together as several
+ // threads can send with this producer at the same time
+ private volatile CachedConsumer cachedConsumer;
private final KameletEndpoint endpoint;
private final KameletComponent component;
@@ -56,10 +57,7 @@ final class KameletProducer extends DefaultAsyncProducer
implements RouteIdAware
@Override
public boolean process(Exchange exchange, AsyncCallback callback) {
try {
- if (consumer == null || stateCounter !=
component.getStateCounter()) {
- stateCounter = component.getStateCounter();
- consumer = component.getConsumer(key, block, timeout);
- }
+ final KameletConsumer consumer = getConsumer();
if (consumer == null) {
if (endpoint.isFailIfNoConsumers()) {
exchange.setException(new
KameletConsumerNotAvailableException(
@@ -148,4 +146,22 @@ final class KameletProducer extends DefaultAsyncProducer
implements RouteIdAware
}
}
+ /**
+ * Gets the consumer, which is looked up again when it has been added or
removed (such as when its route is
+ * suspended or stopped) since it was looked up last.
+ */
+ private KameletConsumer getConsumer() throws InterruptedException {
+ CachedConsumer cached = cachedConsumer;
+ // read the counter before the lookup, so a change during the lookup
makes the next exchange look up again
+ int stateCounter = component.getStateCounter();
+ if (cached == null || cached.consumer() == null ||
cached.stateCounter() != stateCounter) {
+ KameletConsumer consumer = component.getConsumer(key, block,
timeout);
+ cachedConsumer = new CachedConsumer(consumer, stateCounter);
+ return consumer;
+ }
+ return cached.consumer();
+ }
+
+ private record CachedConsumer(KameletConsumer consumer, int stateCounter) {
+ }
}
diff --git
a/components/camel-kamelet/src/test/java/org/apache/camel/component/kamelet/KameletProducerSuspendedConsumerTest.java
b/components/camel-kamelet/src/test/java/org/apache/camel/component/kamelet/KameletProducerSuspendedConsumerTest.java
new file mode 100644
index 000000000000..bab44936cdd3
--- /dev/null
+++
b/components/camel-kamelet/src/test/java/org/apache/camel/component/kamelet/KameletProducerSuspendedConsumerTest.java
@@ -0,0 +1,123 @@
+/*
+ * 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.kamelet;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * While one exchange waits for the consumer of a suspended or stopped kamelet
route, other exchanges sent with the same
+ * producer must wait too, and not be sent to the suspended or stopped
consumer.
+ */
+@Timeout(30)
+public class KameletProducerSuspendedConsumerTest extends CamelTestSupport {
+
+ // counted down by each exchange that looks up the consumer of the kamelet
route
+ private final AtomicReference<CountDownLatch> lookups = new
AtomicReference<>(new CountDownLatch(0));
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext context = super.createCamelContext();
+ context.addComponent("kamelet", new KameletComponent() {
+ @Override
+ protected KameletConsumer getConsumer(String key, boolean block,
long timeout) throws InterruptedException {
+ lookups.get().countDown();
+ return super.getConsumer(key, block, timeout);
+ }
+ });
+ return context;
+ }
+
+ @Test
+ public void testExchangesWaitForSuspendedConsumer() throws Exception {
+ testExchangesWaitForConsumer(() ->
context.getRouteController().suspendRoute("echo"),
+ () -> context.getRouteController().resumeRoute("echo"));
+ }
+
+ @Test
+ public void testExchangesWaitForStoppedConsumer() throws Exception {
+ testExchangesWaitForConsumer(() ->
context.getRouteController().stopRoute("echo"),
+ () -> context.getRouteController().startRoute("echo"));
+ }
+
+ private void testExchangesWaitForConsumer(RouteAction removeConsumer,
RouteAction addConsumer) throws Exception {
+ MockEndpoint mock = getMockEndpoint("mock:kamelet");
+ mock.expectedBodiesReceived("warm");
+ // the producer of the kamelet now caches the consumer of the kamelet
route
+ template.sendBody("direct:start", "warm");
+ mock.assertIsSatisfied();
+
+ removeConsumer.run();
+
+ CountDownLatch firstWaiting = new CountDownLatch(1);
+ lookups.set(firstWaiting);
+ CompletableFuture<Object> first =
template.asyncSendBody("direct:start", "first");
+ // the first exchange found that the consumer changed and waits for
the consumer
+ assertThat(firstWaiting.await(10, TimeUnit.SECONDS)).isTrue();
+
+ CountDownLatch secondWaiting = new CountDownLatch(1);
+ lookups.set(secondWaiting);
+ CompletableFuture<Object> second =
template.asyncSendBody("direct:start", "second");
+
+ // the second exchange must wait for the consumer too, and not reach
the kamelet route while its consumer is gone
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
secondWaiting.getCount() == 0 || second.isDone());
+ assertThat(mock.getReceivedCounter()).as("No exchange should reach the
kamelet route while its consumer is gone")
+ .isEqualTo(1);
+ assertThat(second).as("The second exchange should wait for the
consumer of the kamelet route").isNotDone();
+
+ mock.reset();
+ mock.expectedBodiesReceivedInAnyOrder("first", "second");
+ addConsumer.run();
+
+ first.get(10, TimeUnit.SECONDS);
+ second.get(10, TimeUnit.SECONDS);
+ mock.assertIsSatisfied();
+ }
+
+ @FunctionalInterface
+ private interface RouteAction {
+ void run() throws Exception;
+ }
+
+ @Override
+ protected RoutesBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ routeTemplate("echo")
+ .from("kamelet:source")
+ .to("mock:kamelet");
+
+ from("direct:start")
+ .to("kamelet:echo/echo?timeout=20000");
+ }
+ };
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/component/direct/DirectProducerSuspendedConsumerTest.java
b/core/camel-core/src/test/java/org/apache/camel/component/direct/DirectProducerSuspendedConsumerTest.java
new file mode 100644
index 000000000000..b49b4c02c83e
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/component/direct/DirectProducerSuspendedConsumerTest.java
@@ -0,0 +1,119 @@
+/*
+ * 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.direct;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * While one exchange waits for the consumer of a suspended or stopped direct
route, other exchanges sent with the same
+ * producer must wait too, and not be sent to the suspended or stopped
consumer.
+ */
+@Timeout(30)
+public class DirectProducerSuspendedConsumerTest extends ContextTestSupport {
+
+ // counted down by each exchange that looks up the consumer of the
suspended route
+ private final AtomicReference<CountDownLatch> lookups = new
AtomicReference<>(new CountDownLatch(0));
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext context = super.createCamelContext();
+ context.addComponent("direct", new DirectComponent() {
+ @Override
+ protected DirectConsumer getConsumer(String key, boolean block,
long timeout) throws InterruptedException {
+ lookups.get().countDown();
+ return super.getConsumer(key, block, timeout);
+ }
+ });
+ return context;
+ }
+
+ @Test
+ public void testExchangesWaitForSuspendedConsumer() throws Exception {
+ testExchangesWaitForConsumer(() ->
context.getRouteController().suspendRoute("b"),
+ () -> context.getRouteController().resumeRoute("b"));
+ }
+
+ @Test
+ public void testExchangesWaitForStoppedConsumer() throws Exception {
+ testExchangesWaitForConsumer(() ->
context.getRouteController().stopRoute("b"),
+ () -> context.getRouteController().startRoute("b"));
+ }
+
+ private void testExchangesWaitForConsumer(RouteAction removeConsumer,
RouteAction addConsumer) throws Exception {
+ MockEndpoint mock = getMockEndpoint("mock:b");
+ mock.expectedBodiesReceived("warm");
+ // the producer of direct:b now caches the consumer of route b
+ template.sendBody("direct:start", "warm");
+ mock.assertIsSatisfied();
+
+ removeConsumer.run();
+
+ CountDownLatch firstWaiting = new CountDownLatch(1);
+ lookups.set(firstWaiting);
+ CompletableFuture<Object> first =
template.asyncSendBody("direct:start", "first");
+ // the first exchange found that the consumer changed and waits for
the consumer
+ assertTrue(firstWaiting.await(10, TimeUnit.SECONDS));
+
+ CountDownLatch secondWaiting = new CountDownLatch(1);
+ lookups.set(secondWaiting);
+ CompletableFuture<Object> second =
template.asyncSendBody("direct:start", "second");
+
+ // the second exchange must wait for the consumer too, and not reach
route b while its consumer is gone
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
secondWaiting.getCount() == 0 || second.isDone());
+ assertEquals(1, mock.getReceivedCounter(), "No exchange should reach
route b while its consumer is gone");
+ assertFalse(second.isDone(), "The second exchange should wait for the
consumer of route b");
+
+ mock.reset();
+ mock.expectedBodiesReceivedInAnyOrder("first", "second");
+ addConsumer.run();
+
+ first.get(10, TimeUnit.SECONDS);
+ second.get(10, TimeUnit.SECONDS);
+ mock.assertIsSatisfied();
+ }
+
+ @FunctionalInterface
+ private interface RouteAction {
+ void run() throws Exception;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ public void configure() {
+ from("direct:start").to("direct:b?timeout=20000");
+
+ from("direct:b").routeId("b").to("mock:b");
+ }
+ };
+ }
+}