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");
+            }
+        };
+    }
+}

Reply via email to