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 5f96f71dfada CAMEL-25278: camel-cometd - a stopped consumer must stop 
listening to its channel, and the server extensions and listeners must be 
registered once (#27307)
5f96f71dfada is described below

commit 5f96f71dfada648a73577d1dc494fbfeb7659dc8
Author: allthingssecurity <[email protected]>
AuthorDate: Sat Oct 3 12:18:44 2026 +0530

    CAMEL-25278: camel-cometd - a stopped consumer must stop listening to its 
channel, and the server extensions and listeners must be registered once 
(#27307)
    
    The producers and consumers of a host and port share one CometD server, 
which is only stopped with the last of them:
    - `CometdConsumer` added its service (a listener of the channel) on start 
and never removed it. When the server outlived the consumer (a producer, or 
another consumer, on the same host and port), the stopped consumer kept calling 
the processor of the route: after a restart of the route each message was 
processed twice, once more per restart.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../camel/component/cometd/CometdComponent.java    |  34 +++---
 .../camel/component/cometd/CometdConsumer.java     |   7 ++
 .../cometd/CometdConsumerRestartTest.java          | 114 +++++++++++++++++++++
 3 files changed, 139 insertions(+), 16 deletions(-)

diff --git 
a/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdComponent.java
 
b/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdComponent.java
index b51a0ef064a9..3a1d5fdf12af 100644
--- 
a/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdComponent.java
+++ 
b/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdComponent.java
@@ -144,26 +144,28 @@ public class CometdComponent extends DefaultComponent 
implements SSLContextParam
                 server.start();
 
                 connectors.put(connectorKey, connectorRef);
-            } else {
-                connectorRef.increment();
-            }
 
-            BayeuxServerImpl bayeux = (BayeuxServerImpl) 
connectorRef.servlet.getBayeuxServer();
-
-            if (securityPolicy != null) {
-                bayeux.setSecurityPolicy(securityPolicy);
-            }
-            if (extensions != null) {
-                for (BayeuxServer.Extension extension : extensions) {
-                    bayeux.addExtension(extension);
+                // the server is shared by the producers and consumers of this 
host and port: configure it once,
+                // or the extensions and listeners are called once per 
producer and consumer
+                BayeuxServerImpl bayeux = (BayeuxServerImpl) 
connectorRef.servlet.getBayeuxServer();
+                if (securityPolicy != null) {
+                    bayeux.setSecurityPolicy(securityPolicy);
                 }
-            }
-            if (serverListeners != null) {
-                for (BayeuxServer.BayeuxServerListener serverListener : 
serverListeners) {
-                    bayeux.addListener(serverListener);
+                if (extensions != null) {
+                    for (BayeuxServer.Extension extension : extensions) {
+                        bayeux.addExtension(extension);
+                    }
+                }
+                if (serverListeners != null) {
+                    for (BayeuxServer.BayeuxServerListener serverListener : 
serverListeners) {
+                        bayeux.addListener(serverListener);
+                    }
                 }
+            } else {
+                connectorRef.increment();
             }
-            prodcon.setBayeux(bayeux);
+
+            prodcon.setBayeux((BayeuxServerImpl) 
connectorRef.servlet.getBayeuxServer());
         } finally {
             connectorsLock.unlock();
         }
diff --git 
a/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdConsumer.java
 
b/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdConsumer.java
index 406324965ec2..6417c141d540 100644
--- 
a/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdConsumer.java
+++ 
b/components/camel-cometd/src/main/java/org/apache/camel/component/cometd/CometdConsumer.java
@@ -55,6 +55,13 @@ public class CometdConsumer extends DefaultConsumer 
implements CometdProducerCon
 
     @Override
     public void doStop() throws Exception {
+        if (service != null) {
+            // the server outlives this consumer when other producers or 
consumers use its host and port: stop
+            // listening to the channel, or the messages are still processed 
by this consumer
+            service.removeService(endpoint.getPath());
+            service.getLocalSession().disconnect();
+            service = null;
+        }
         endpoint.disconnect(this);
         super.doStop();
     }
diff --git 
a/components/camel-cometd/src/test/java/org/apache/camel/component/cometd/CometdConsumerRestartTest.java
 
b/components/camel-cometd/src/test/java/org/apache/camel/component/cometd/CometdConsumerRestartTest.java
new file mode 100644
index 000000000000..fba16416010e
--- /dev/null
+++ 
b/components/camel-cometd/src/test/java/org/apache/camel/component/cometd/CometdConsumerRestartTest.java
@@ -0,0 +1,114 @@
+/*
+ * 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.cometd;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.AvailablePortFinder;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.cometd.bayeux.server.BayeuxServer;
+import org.cometd.bayeux.server.LocalSession;
+import org.cometd.bayeux.server.ServerMessage;
+import org.cometd.bayeux.server.ServerSession;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * A producer and a consumer share the CometD server of their host and port: 
the server is only stopped with the last of
+ * them, so what a consumer registers on the server must be removed when it 
stops, and what the component registers must
+ * be registered once.
+ */
+public class CometdConsumerRestartTest extends CamelTestSupport {
+
+    @RegisterExtension
+    AvailablePortFinder.Port port = AvailablePortFinder.find();
+
+    private final CountingListener listener = new CountingListener();
+    private String uri;
+
+    @Test
+    void testRestartedConsumerReceivesAMessageOnce() throws Exception {
+        context.getRouteController().stopRoute("consumer");
+        context.getRouteController().startRoute("consumer");
+
+        MockEndpoint mock = getMockEndpoint("mock:test");
+        mock.expectedBodiesReceived("Hello");
+
+        template.sendBody("direct:input", "Hello");
+
+        mock.assertIsSatisfied();
+        // the message is delivered synchronously: the consumer of before the 
restart must not receive it too
+        assertEquals(1, mock.getReceivedCounter());
+    }
+
+    @Test
+    void testStoppedConsumerDoesNotReceive() throws Exception {
+        context.getRouteController().stopRoute("consumer");
+
+        template.sendBody("direct:input", "Hello");
+
+        assertEquals(0, getMockEndpoint("mock:test").getReceivedCounter());
+    }
+
+    @Test
+    void testServerListenerIsCalledOncePerSession() {
+        // the producer and the consumer are both connected to the server
+        CometdConsumer consumer = (CometdConsumer) 
context.getRoute("consumer").getConsumer();
+        LocalSession session = 
consumer.getConsumerService().getBayeux().newLocalSession("probe");
+        listener.added.set(0);
+        session.handshake();
+        try {
+            assertEquals(1, listener.added.get());
+        } finally {
+            session.disconnect();
+        }
+    }
+
+    @Override
+    public void doPreSetup() {
+        uri = "cometd://127.0.0.1:" + port.getPort() + 
"/service/test?baseResource=file:./target/test-classes/webapp&"
+              + 
"timeout=240000&interval=0&maxInterval=30000&multiFrameInterval=1500&jsonCommented=true&logLevel=2";
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                context.getComponent("cometd", 
CometdComponent.class).addServerListener(listener);
+
+                from("direct:input").to(uri);
+
+                from(uri).routeId("consumer").to("mock:test");
+            }
+        };
+    }
+
+    public static final class CountingListener implements 
BayeuxServer.SessionListener {
+
+        private final AtomicInteger added = new AtomicInteger();
+
+        @Override
+        public void sessionAdded(ServerSession session, ServerMessage message) 
{
+            added.incrementAndGet();
+        }
+    }
+}

Reply via email to