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 8ac66d3d61e7 CAMEL-25276: camel-pubnub - a consumer must remove its 
listener on stop and only receive its own channel (#27305)
8ac66d3d61e7 is described below

commit 8ac66d3d61e7f7d79b8d665152e87b21142aef7d
Author: allthingssecurity <[email protected]>
AuthorDate: Sat Oct 3 13:46:36 2026 +0530

    CAMEL-25276: camel-pubnub - a consumer must remove its listener on stop and 
only receive its own channel (#27305)
    
    A PubNub client hands every message and presence event to all its 
listeners, whatever the channel, and since CAMEL-16142 one client can be shared 
by all `pubnub:` endpoints (autowired, or `pubnub=#bean`):
    - **Cross-talk.** The consumer's `SubscribeCallback` did not look at the 
channel, so with a shared client each route also received the messages of the 
other routes' channels.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../camel/component/pubnub/PubNubConsumer.java     |  31 ++-
 .../camel/component/pubnub/PubNubEndpoint.java     |  11 +-
 .../pubnub/PubNubConsumerListenerTest.java         | 254 +++++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  12 +
 4 files changed, 305 insertions(+), 3 deletions(-)

diff --git 
a/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubConsumer.java
 
b/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubConsumer.java
index 07b0a91c55c5..d99b36dea714 100644
--- 
a/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubConsumer.java
+++ 
b/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubConsumer.java
@@ -48,6 +48,7 @@ public class PubNubConsumer extends DefaultConsumer {
 
     private final PubNubEndpoint endpoint;
     private final PubNubConfiguration pubNubConfiguration;
+    private PubNubCallback callback;
 
     public PubNubConsumer(PubNubEndpoint endpoint, Processor processor, 
PubNubConfiguration pubNubConfiguration) {
         super(endpoint, processor);
@@ -56,7 +57,10 @@ public class PubNubConsumer extends DefaultConsumer {
     }
 
     private void initCommunication() {
-        endpoint.getPubnub().addListener(new PubNubCallback());
+        // the listener receives the events of all the channels the PubNub 
client subscribes to (the client may be
+        // shared by several endpoints), so it filters on the channel, and it 
is removed when this consumer stops
+        callback = new PubNubCallback();
+        endpoint.getPubnub().addListener(callback);
         if (pubNubConfiguration.isWithPresence()) {
             
endpoint.getPubnub().subscribe().channels(Arrays.asList(pubNubConfiguration.getChannel())).withPresence().execute();
         } else {
@@ -65,6 +69,10 @@ public class PubNubConsumer extends DefaultConsumer {
     }
 
     private void terminateCommunication() {
+        if (callback != null) {
+            endpoint.getPubnub().removeListener(callback);
+            callback = null;
+        }
         try {
             
endpoint.getPubnub().unsubscribe().channels(Arrays.asList(pubNubConfiguration.getChannel())).execute();
         } catch (Exception e) {
@@ -72,6 +80,21 @@ public class PubNubConsumer extends DefaultConsumer {
         }
     }
 
+    /**
+     * Whether the event is for a channel of this consumer: the channel of the 
event, or the subscription it was
+     * received through (a wildcard channel such as news.*). The channel 
option may list several channels separated by
+     * commas, which the PubNub client subscribes to as separate channels.
+     */
+    private boolean isForThisChannel(String channel, String subscription) {
+        for (String ours : pubNubConfiguration.getChannel().split(",")) {
+            ours = ours.trim();
+            if (ours.equals(channel) || ours.equals(subscription)) {
+                return true;
+            }
+        }
+        return false;
+    }
+
     @Override
     protected void doStart() throws Exception {
         super.doStart();
@@ -110,6 +133,9 @@ public class PubNubConsumer extends DefaultConsumer {
 
         @Override
         public void message(PubNub pubnub, PNMessageResult message) {
+            if (!isForThisChannel(message.getChannel(), 
message.getSubscription())) {
+                return;
+            }
             Exchange exchange = createExchange(true);
             Message inmessage = exchange.getIn();
             inmessage.setBody(message);
@@ -129,6 +155,9 @@ public class PubNubConsumer extends DefaultConsumer {
 
         @Override
         public void presence(PubNub pubnub, PNPresenceEventResult presence) {
+            if (!isForThisChannel(presence.getChannel(), 
presence.getSubscription())) {
+                return;
+            }
             Exchange exchange = createExchange(true);
             Message inmessage = exchange.getIn();
             inmessage.setBody(presence);
diff --git 
a/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubEndpoint.java
 
b/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubEndpoint.java
index f157394a3b0c..64d9038a55c7 100644
--- 
a/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubEndpoint.java
+++ 
b/components/camel-pubnub/src/main/java/org/apache/camel/component/pubnub/PubNubEndpoint.java
@@ -43,6 +43,9 @@ public class PubNubEndpoint extends DefaultEndpoint {
     @UriParam
     private PubNubConfiguration configuration;
 
+    // whether the endpoint created the client (a client from the registry may 
be shared, and is not ours to destroy)
+    private boolean createdClient;
+
     public PubNubEndpoint(String uri, PubNubComponent component, 
PubNubConfiguration configuration) {
         super(uri, component);
         this.configuration = configuration;
@@ -76,15 +79,19 @@ public class PubNubEndpoint extends DefaultEndpoint {
     @Override
     protected void doStop() throws Exception {
         super.doStop();
-        if (pubnub != null) {
+        if (pubnub != null && createdClient) {
             pubnub.destroy();
             pubnub = null;
+            createdClient = false;
         }
     }
 
     @Override
     protected void doStart() throws Exception {
-        this.pubnub = getPubnub() != null ? getPubnub() : getInstance();
+        if (pubnub == null) {
+            pubnub = getInstance();
+            createdClient = true;
+        }
         super.doStart();
     }
 
diff --git 
a/components/camel-pubnub/src/test/java/org/apache/camel/component/pubnub/PubNubConsumerListenerTest.java
 
b/components/camel-pubnub/src/test/java/org/apache/camel/component/pubnub/PubNubConsumerListenerTest.java
new file mode 100644
index 000000000000..14a81b529b6f
--- /dev/null
+++ 
b/components/camel-pubnub/src/test/java/org/apache/camel/component/pubnub/PubNubConsumerListenerTest.java
@@ -0,0 +1,254 @@
+/*
+ * 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.pubnub;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+import com.google.gson.JsonPrimitive;
+import com.pubnub.api.UserId;
+import com.pubnub.api.enums.PNLogVerbosity;
+import com.pubnub.api.java.v2.PNConfiguration;
+import com.pubnub.api.models.consumer.pubsub.BasePubSubResult;
+import com.pubnub.api.models.consumer.pubsub.PNMessageResult;
+import com.pubnub.api.models.consumer.pubsub.PNPresenceEventResult;
+import com.pubnub.internal.java.PubNubForJavaImpl;
+import org.apache.camel.BindToRegistry;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static com.github.tomakehurst.wiremock.client.WireMock.aResponse;
+import static com.github.tomakehurst.wiremock.client.WireMock.get;
+import static com.github.tomakehurst.wiremock.client.WireMock.stubFor;
+import static com.github.tomakehurst.wiremock.client.WireMock.urlPathMatching;
+import static com.pubnub.api.enums.PNHeartbeatNotificationOptions.NONE;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+
+/**
+ * The consumers of several endpoints share one PubNub client: each consumer 
only gets the messages of its own channel,
+ * only once after a restart, none after it stopped, and removing one route 
does not destroy the shared client.
+ * <p>
+ * The PubNub client hands every message it receives to all its listeners, so 
the tests announce messages through the
+ * listeners of the real client instead of going through a subscribe response.
+ */
+public class PubNubConsumerListenerTest extends PubNubTestBase {
+
+    @BindToRegistry("sharedPubnub")
+    private final SharedPubNub sharedPubnub = new SharedPubNub(port.getPort());
+
+    @BeforeEach
+    public void stubSubscribe() {
+        stubFor(get(urlPathMatching("/v2/subscribe/mySubscribeKey/.*"))
+                
.willReturn(aResponse().withBody("{\"t\":{\"t\":\"14607577960932487\",\"r\":1},\"m\":[]}")
+                        .withFixedDelay(100)));
+        stubFor(get(urlPathMatching("/v2/presence/.*"))
+                .willReturn(aResponse().withBody("{\"status\": 200, 
\"message\": \"OK\", \"service\": \"Presence\"}")));
+    }
+
+    @Override
+    protected void cleanupResources() {
+        super.cleanupResources();
+        sharedPubnub.destroy();
+    }
+
+    @Test
+    public void testConsumerOnlyReceivesItsOwnChannel() throws Exception {
+        context.getRouteController().startRoute("alpha");
+        context.getRouteController().startRoute("beta");
+
+        MockEndpoint alpha = getMockEndpoint("mock:alpha");
+        alpha.expectedMessageCount(1);
+        alpha.expectedHeaderReceived(PubNubConstants.CHANNEL, "alpha");
+        MockEndpoint beta = getMockEndpoint("mock:beta");
+        beta.expectedMessageCount(0);
+
+        sharedPubnub.announce("alpha");
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    public void testWildcardSubscriptionAndPresenceEvents() throws Exception {
+        context.getRouteController().startRoute("alpha");
+        context.getRouteController().startRoute("news");
+
+        MockEndpoint alpha = getMockEndpoint("mock:alpha");
+        alpha.expectedMessageCount(1);
+        alpha.message(0).body().isInstanceOf(PNPresenceEventResult.class);
+        MockEndpoint news = getMockEndpoint("mock:news");
+        news.expectedMessageCount(1);
+        news.expectedHeaderReceived(PubNubConstants.CHANNEL, "news.sports");
+
+        // a message of a channel matched by the wildcard subscription, and a 
presence event (the client removes the
+        // -pnpres suffix of the presence channel)
+        sharedPubnub.announce("news.sports", "news.*");
+        sharedPubnub.announcePresence("alpha");
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    public void testCommaSeparatedChannels() throws Exception {
+        context.getRouteController().startRoute("multi");
+
+        MockEndpoint multi = getMockEndpoint("mock:multi");
+        multi.expectedMessageCount(2);
+
+        // the PubNub client subscribes to each channel of a comma separated 
list
+        sharedPubnub.announce("delta");
+        sharedPubnub.announce("epsilon");
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    public void testMessageReceivedOnceAfterRestart() throws Exception {
+        context.getRouteController().startRoute("alpha");
+        context.getRouteController().stopRoute("alpha");
+        context.getRouteController().startRoute("alpha");
+
+        MockEndpoint alpha = getMockEndpoint("mock:alpha");
+        alpha.expectedMessageCount(1);
+
+        sharedPubnub.announce("alpha");
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    public void testNoMessageAfterStop() throws Exception {
+        context.getRouteController().startRoute("alpha");
+        context.getRouteController().stopRoute("alpha");
+
+        MockEndpoint alpha = getMockEndpoint("mock:alpha");
+        alpha.expectedMessageCount(0);
+
+        sharedPubnub.announce("alpha");
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    public void testRemovingRouteKeepsSharedClient() throws Exception {
+        context.getRouteController().startRoute("alpha");
+        context.getRouteController().startRoute("beta");
+
+        // alpha is the only route that uses its endpoint, so removing the 
route stops and removes the endpoint
+        context.getRouteController().stopRoute("alpha");
+        context.removeRoute("alpha");
+
+        assertEquals(0, sharedPubnub.destroyed.get(), "The shared PubNub 
client must not be destroyed");
+        assertSame(sharedPubnub, 
context.getEndpoint("pubnub:beta?pubnub=#sharedPubnub", 
PubNubEndpoint.class).getPubnub());
+
+        MockEndpoint beta = getMockEndpoint("mock:beta");
+        beta.expectedMessageCount(1);
+
+        sharedPubnub.announce("beta");
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    public void testEndpointDestroysItsOwnClient() throws Exception {
+        PubNubEndpoint endpoint = context.getEndpoint(
+                
"pubnub:gamma?subscribeKey=mySubscribeKey&publishKey=myPublishKey&secretKey=mySecretKey&authKey=myAuthKey"
+                                                      + "&uuid=myUUID",
+                PubNubEndpoint.class);
+        assertNotNull(endpoint.getPubnub());
+
+        endpoint.stop();
+        assertNull(endpoint.getPubnub());
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            public void configure() {
+                
from("pubnub:alpha?pubnub=#sharedPubnub").id("alpha").autoStartup(false)
+                        .to("mock:alpha");
+
+                
from("pubnub:beta?pubnub=#sharedPubnub").id("beta").autoStartup(false)
+                        .to("mock:beta");
+
+                
from("pubnub:news.*?pubnub=#sharedPubnub").id("news").autoStartup(false)
+                        .to("mock:news");
+
+                
from("pubnub:delta,epsilon?pubnub=#sharedPubnub").id("multi").autoStartup(false)
+                        .to("mock:multi");
+            }
+        };
+    }
+
+    /**
+     * A real PubNub client that counts how often it is destroyed and lets the 
test hand a message to its listeners.
+     * <p>
+     * It subclasses {@code PubNubForJavaImpl}, which is internal to the 
PubNub SDK ({@code com.pubnub.internal}),
+     * because that is the class {@code PubNub.create(...)} instantiates, and 
only the implementation exposes
+     * {@code getListenerManager()}: the public {@code PubNub} interface 
offers no way to hand an event to the listeners
+     * without a subscribe response from the server. {@code PubNubTestBase} 
subclasses it for the same reason. If a
+     * PubNub upgrade renames or moves this class, follow the class that 
{@code PubNub.create(...)} instantiates.
+     */
+    static class SharedPubNub extends PubNubForJavaImpl {
+
+        final AtomicInteger destroyed = new AtomicInteger();
+
+        SharedPubNub(int port) {
+            super(createConfiguration(port));
+        }
+
+        private static PNConfiguration createConfiguration(int port) {
+            try {
+                return PNConfiguration.builder(new UserId("myUUID"), 
"mySubscribeKey")
+                        .publishKey("myPublishKey")
+                        .secure(false)
+                        .origin("localhost:" + port)
+                        .logVerbosity(PNLogVerbosity.NONE)
+                        .heartbeatNotificationOptions(NONE)
+                        .build();
+            } catch (Exception e) {
+                throw new RuntimeException(e);
+            }
+        }
+
+        void announce(String channel) {
+            // the client gives no subscription when the message came through 
the channel itself
+            announce(channel, null);
+        }
+
+        void announce(String channel, String subscription) {
+            getListenerManager().announce(new PNMessageResult(
+                    new BasePubSubResult(channel, subscription, 
14607577960932488L, null, "publisher"),
+                    new JsonPrimitive("Hello " + channel), null, null));
+        }
+
+        void announcePresence(String channel) {
+            getListenerManager().announce(new PNPresenceEventResult(
+                    "join", "someone", 1460757796L, 1, null, channel, null, 
14607577960932488L, null, null, null, false,
+                    null));
+        }
+
+        @Override
+        public void destroy() {
+            destroyed.incrementAndGet();
+            super.destroy();
+        }
+    }
+}
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 38f4fc054357..de469765b4f5 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -1551,6 +1551,18 @@ such a configuration that declares no 
`JavaSerializationFilterConfig`, unless a
 is set. To remove the warning, declare a `java-serialization-filter` in the 
configuration, as described in
 the xref:components::hazelcast-summary.adoc[Hazelcast component] documentation.
 
+=== camel-pubnub - consumers of a shared PubNub client
+
+A PubNub client hands every message it receives to all its listeners. When 
several endpoints use the same client
+(a `pubnub` bean that is autowired or referenced), each consumer received the 
messages of the channels of the other
+consumers too, and every stop and start of a route added one more listener, so 
its messages were then received once
+more. A consumer now only receives the messages and presence events of its own 
channel, and removes its listener
+when it stops.
+
+The endpoint now destroys the PubNub client only when it created it. A client 
from the registry is no longer
+destroyed when an endpoint stops, for example when a route that uses the 
client is removed while other routes still
+use it: the application that created the client is responsible for destroying 
it.
+
 === camel-netty - NettyConverter.toByteArray returns a copy of the buffer's 
readable bytes
 
 The `ByteBuf` to `byte[]` type converter (`NettyConverter.toByteArray`, also 
used when converting a Netty

Reply via email to