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