This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch camel-4.22.x
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/camel-4.22.x by this push:
     new d2a3ce43318b CAMEL-24408: Fix SEDA discardIfNoConsumers after consumer 
route removal (4.22.x backport)
d2a3ce43318b is described below

commit d2a3ce43318b7fe1fd9b2d9cbeedbc8fbe14d487
Author: Claus Ibsen <[email protected]>
AuthorDate: Thu Aug 27 16:50:19 2026 +0200

    CAMEL-24408: Fix SEDA discardIfNoConsumers after consumer route removal 
(4.22.x backport)
    
    Backport of #25558 to camel-4.22.x.
    
    Straight cherry-pick of the squashed commit 
d19b4f96fadc28c6ecbff5f3efeb16f9f2c56f7f from #25558, already reviewed on main. 
Replaces #25560, which was missing the 
SedaDiscardIfNoConsumersProducerRouteRestartTest regression test.
    
    Closes #25811
    Closes #25560
    
    Co-authored-by: Cursor Agent <[email protected]>
    Co-authored-by: Omar Atie <[email protected]>
---
 .../camel/component/seda/QueueReference.java       |  21 ++++
 .../apache/camel/component/seda/SedaComponent.java |  23 ++--
 .../apache/camel/component/seda/SedaEndpoint.java  |  23 ++--
 .../SedaDiscardIfNoConsumerAfterRemovalTest.java   | 137 +++++++++++++++++++++
 ...scardIfNoConsumersProducerRouteRestartTest.java |  56 +++++++++
 5 files changed, 244 insertions(+), 16 deletions(-)

diff --git 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/QueueReference.java
 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/QueueReference.java
index 4efa264ba340..28e97fff894a 100644
--- 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/QueueReference.java
+++ 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/QueueReference.java
@@ -101,6 +101,9 @@ public final class QueueReference {
         return queue;
     }
 
+    /**
+     * Whether any of the endpoints sharing this queue reference still have 
active consumers.
+     */
     public boolean hasConsumers() {
         lock.lock();
         try {
@@ -115,4 +118,22 @@ public final class QueueReference {
             lock.unlock();
         }
     }
+
+    /**
+     * Whether any of the endpoints sharing this queue reference still have 
active producers.
+     */
+    public boolean hasProducers() {
+        lock.lock();
+        try {
+            for (SedaEndpoint endpoint : endpoints) {
+                if (!endpoint.getProducers().isEmpty()) {
+                    return true;
+                }
+            }
+
+            return false;
+        } finally {
+            lock.unlock();
+        }
+    }
 }
diff --git 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaComponent.java
 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaComponent.java
index da1d72b9e86c..224629ac2a3e 100644
--- 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaComponent.java
+++ 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaComponent.java
@@ -323,16 +323,21 @@ public class SedaComponent extends DefaultComponent {
      * @param endpoint the endpoint
      */
     void onShutdownEndpoint(SedaEndpoint endpoint) {
-        // we need to remove the endpoint from the reference counter
-        String key = getQueueKey(endpoint.getEndpointUri());
-        QueueReference ref = getQueues().get(key);
-        if (ref != null && endpoint.getConsumers().isEmpty()) {
-            // only remove the endpoint when the consumers are removed
-            ref.removeReference(endpoint);
-            if (ref.getCount() <= 0) {
-                // reference no longer needed so remove from queues
-                getQueues().remove(key);
+        lock.lock();
+        try {
+            // we need to remove the endpoint from the reference counter
+            String key = getQueueKey(endpoint.getEndpointUri());
+            QueueReference ref = getQueues().get(key);
+            if (ref != null && endpoint.getConsumers().isEmpty() && 
endpoint.getProducers().isEmpty()) {
+                // only remove the endpoint when both consumers and producers 
are removed
+                ref.removeReference(endpoint);
+                if (ref.getCount() <= 0 && !ref.hasConsumers() && 
!ref.hasProducers()) {
+                    // reference no longer needed so remove from queues
+                    getQueues().remove(key);
+                }
             }
+        } finally {
+            lock.unlock();
         }
     }
 
diff --git 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaEndpoint.java
 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaEndpoint.java
index 42c0534db793..2dc89a4023cb 100644
--- 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaEndpoint.java
+++ 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaEndpoint.java
@@ -605,10 +605,20 @@ public class SedaEndpoint extends DefaultEndpoint 
implements AsyncEndpoint, Brow
 
     void onStarted(SedaProducer producer) {
         producers.add(producer);
+        if (getComponent() != null && (ref == null || queue == null)) {
+            // re-register queue reference when producer restarts after queue 
was released on stop
+            Integer size = (getSize() == Integer.MAX_VALUE || getSize() == 
SedaConstants.QUEUE_SIZE) ? null : getSize();
+            ref = getComponent().getOrCreateQueue(this, size, 
isMultipleConsumers(), queueFactory);
+            queue = ref.getQueue();
+        }
     }
 
     void onStopped(SedaProducer producer) {
         producers.remove(producer);
+        if (getConsumers().isEmpty() && getProducers().isEmpty() && 
getComponent() != null) {
+            // may also be invoked from shutdown(); onShutdownEndpoint is 
idempotent
+            getComponent().onShutdownEndpoint(this);
+        }
     }
 
     void onStarted(SedaConsumer consumer) throws Exception {
@@ -660,13 +670,12 @@ public class SedaEndpoint extends DefaultEndpoint 
implements AsyncEndpoint, Brow
 
     @Override
     public void stop() {
-        if (getConsumers().isEmpty()) {
+        if (getConsumers().isEmpty() && getProducers().isEmpty()) {
             super.stop();
+            ref = null;
         } else {
-            LOG.debug("There is still active consumers.");
+            LOG.debug("There are still active consumers or producers.");
         }
-
-        ref = null;
     }
 
     @Override
@@ -676,15 +685,15 @@ public class SedaEndpoint extends DefaultEndpoint 
implements AsyncEndpoint, Brow
             return;
         }
 
-        // notify component we are shutting down this endpoint
+        // notify component we are shutting down this endpoint (onStopped may 
invoke this too; safe to call twice)
         if (getComponent() != null) {
             getComponent().onShutdownEndpoint(this);
         }
 
-        if (getConsumers().isEmpty()) {
+        if (getConsumers().isEmpty() && getProducers().isEmpty()) {
             super.shutdown();
         } else {
-            LOG.debug("There is still active consumers.");
+            LOG.debug("There are still active consumers or producers.");
         }
     }
 
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaDiscardIfNoConsumerAfterRemovalTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaDiscardIfNoConsumerAfterRemovalTest.java
new file mode 100644
index 000000000000..0391cae10b00
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaDiscardIfNoConsumerAfterRemovalTest.java
@@ -0,0 +1,137 @@
+/*
+ * 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.seda;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.support.service.ServiceHelper;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+class SedaDiscardIfNoConsumerAfterRemovalTest extends ContextTestSupport {
+
+    @Test
+    void testDiscardAfterConsumerRouteRemoved() throws Exception {
+        SedaComponent seda = context.getComponent("seda", SedaComponent.class);
+        SedaEndpoint bar = 
getMandatoryEndpoint("seda:bar?discardIfNoConsumers=true", SedaEndpoint.class);
+        String key = seda.getQueueKey(bar.getEndpointUri());
+        assertThat(bar.getCurrentQueueSize()).isZero();
+
+        MockEndpoint mock = getMockEndpoint("mock:result");
+        mock.expectedBodiesReceived("Hello World");
+
+        template.sendBody("direct:start", "Hello World");
+
+        mock.assertIsSatisfied();
+
+        context.getRouteController().stopRoute("consumer");
+        context.removeRoute("consumer");
+
+        assertThat(ServiceHelper.isStarted(bar)).isTrue();
+        assertThat(bar.getQueueReference()).isNotNull();
+        assertThat(seda.getQueues().get(key)).isNotNull();
+        assertThat(bar.getQueueReference().hasConsumers()).isFalse();
+
+        template.sendBody("direct:start", "Should be discarded");
+
+        assertThat(bar.getCurrentQueueSize()).isZero();
+    }
+
+    @Test
+    void testReAddConsumerAfterRemoval() throws Exception {
+        template.sendBody("direct:start", "Hello World");
+        getMockEndpoint("mock:result").assertIsSatisfied();
+
+        context.getRouteController().stopRoute("consumer");
+        context.removeRoute("consumer");
+
+        context.addRoutes(new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("seda:bar?discardIfNoConsumers=true").routeId("consumer").to("mock:result");
+            }
+        });
+
+        MockEndpoint mock = getMockEndpoint("mock:result");
+        mock.reset();
+        mock.expectedBodiesReceived("After re-add");
+
+        template.sendBody("direct:start", "After re-add");
+
+        mock.assertIsSatisfied();
+    }
+
+    @Test
+    void testFailIfNoConsumersAfterConsumerRouteRemoved() throws Exception {
+        context.addRoutes(new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("direct:fail").routeId("failProducer").to("seda:fail?failIfNoConsumers=true");
+                
from("seda:fail?failIfNoConsumers=true").routeId("failConsumer").to("mock:fail");
+            }
+        });
+
+        SedaComponent seda = context.getComponent("seda", SedaComponent.class);
+        SedaEndpoint fail = 
getMandatoryEndpoint("seda:fail?failIfNoConsumers=true", SedaEndpoint.class);
+        String key = seda.getQueueKey(fail.getEndpointUri());
+
+        context.getRouteController().stopRoute("failConsumer");
+        context.removeRoute("failConsumer");
+
+        assertThat(fail.getQueueReference()).isNotNull();
+        assertThat(seda.getQueues().get(key)).isNotNull();
+
+        assertThatThrownBy(() -> template.sendBody("direct:fail", "Should 
fail"))
+                .cause()
+                .isInstanceOf(SedaConsumerNotAvailableException.class)
+                .hasMessageContaining("No consumers available");
+    }
+
+    @Test
+    void testQueueRemovedAfterProducerRouteRemoved() throws Exception {
+        SedaComponent seda = context.getComponent("seda", SedaComponent.class);
+        SedaEndpoint bar = 
getMandatoryEndpoint("seda:bar?discardIfNoConsumers=true", SedaEndpoint.class);
+        String key = seda.getQueueKey(bar.getEndpointUri());
+
+        template.sendBody("direct:start", "Hello World");
+        getMockEndpoint("mock:result").assertIsSatisfied();
+
+        context.getRouteController().stopRoute("consumer");
+        context.removeRoute("consumer");
+
+        assertThat(seda.getQueues().get(key)).isNotNull();
+
+        context.getRouteController().stopRoute("producer");
+        context.removeRoute("producer");
+
+        assertThat(seda.getQueues().get(key)).isNull();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("direct:start").routeId("producer").to("seda:bar?discardIfNoConsumers=true");
+                
from("seda:bar?discardIfNoConsumers=true").routeId("consumer").to("mock:result");
+            }
+        };
+    }
+}
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaDiscardIfNoConsumersProducerRouteRestartTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaDiscardIfNoConsumersProducerRouteRestartTest.java
new file mode 100644
index 000000000000..b2da39646a83
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaDiscardIfNoConsumersProducerRouteRestartTest.java
@@ -0,0 +1,56 @@
+/*
+ * 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.seda;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatCode;
+
+class SedaDiscardIfNoConsumersProducerRouteRestartTest extends 
ContextTestSupport {
+
+    @Test
+    void testSendStillWorksAfterProducerRouteRestart() throws Exception {
+        SedaEndpoint bar = 
getMandatoryEndpoint("seda:bar?discardIfNoConsumers=true", SedaEndpoint.class);
+
+        template.sendBody("direct:start", "discarded-1");
+        assertThat(bar.getQueueReference()).isNotNull();
+
+        context.getRouteController().stopRoute("producer");
+        assertThat(bar.getQueueReference()).isNull();
+
+        context.getRouteController().startRoute("producer");
+
+        assertThatCode(() -> template.sendBody("direct:start", "discarded-2"))
+                .as("send after producer route restart should discard 
silently, not fail")
+                .doesNotThrowAnyException();
+        assertThat(bar.getQueueReference()).isNotNull();
+        assertThat(bar.getCurrentQueueSize()).isZero();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("direct:start").routeId("producer").to("seda:bar?discardIfNoConsumers=true");
+            }
+        };
+    }
+}

Reply via email to