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 325a115d5412 CAMEL-24788: camel-vertx-websocket - keep the shared host
until its last consumer stops
325a115d5412 is described below
commit 325a115d5412af4bbb27870a7436fa52a13908dd
Author: Andrea Cosentino <[email protected]>
AuthorDate: Sat Sep 19 08:36:02 2026 +0200
CAMEL-24788: camel-vertx-websocket - keep the shared host until its last
consumer stops
Every consumer on the same host and port shares one VertxWebsocketHost,
but disconnectConsumer removed the host from the registry as soon as the
first consumer stopped. The remaining consumers then found no host on
their own stop, so their route was never unregistered and the Vert.x
HTTP server outlived every Camel route. Restarting a consumer on that
port built a second host against a port the first server still held.
The host is now looked up rather than removed, and only leaves the
registry once it serves no consumer, which is the same moment it stops
its own server. The host lifecycle (connect/disconnect/start/stop) is
guarded by the host's monitor so concurrent route-controller stops
cannot double-stop the server, and disconnect no longer dereferences a
null route for a path it does not serve.
VertxWebsocketMultiConsumerLifecycleTest covers a shared host surviving
the first stop, going away with the last, stop/start/stop of one
consumer, and concurrent stops of both.
Closes #26533
Co-Authored-By: Claude Opus 5 <[email protected]>
Signed-off-by: Andrea Cosentino <[email protected]>
---
.../vertx/websocket/VertxWebsocketComponent.java | 8 +-
.../vertx/websocket/VertxWebsocketHost.java | 23 +++-
.../VertxWebsocketMultiConsumerLifecycleTest.java | 132 +++++++++++++++++++++
3 files changed, 156 insertions(+), 7 deletions(-)
diff --git
a/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketComponent.java
b/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketComponent.java
index 9b057cb8c8a8..a0ea300240f2 100644
---
a/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketComponent.java
+++
b/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketComponent.java
@@ -202,10 +202,16 @@ public class VertxWebsocketComponent extends
DefaultComponent implements SSLCont
VertxWebsocketEndpoint endpoint = consumer.getEndpoint();
VertxWebsocketConfiguration configuration =
endpoint.getConfiguration();
VertxWebsocketHostKey hostKey =
createHostKey(configuration.getWebsocketURI());
- VertxWebsocketHost vertxWebsocketHost =
vertxHostRegistry.remove(hostKey);
+ VertxWebsocketHost vertxWebsocketHost = vertxHostRegistry.get(hostKey);
if (vertxWebsocketHost != null) {
vertxWebsocketHost.disconnect(configuration.getWebsocketURI().getPath());
+
+ // every consumer on this host and port shares the one host, which
stops its server as its last
+ // consumer goes. Forgetting it any earlier would leave the
consumers still on it unable to disconnect.
+ // computeIfPresent decides that against the same key
connectConsumer computes on, so a consumer
+ // connecting at this moment either keeps the host or gets a fresh
one, never a forgotten one
+ vertxHostRegistry.computeIfPresent(hostKey, (key, host) ->
host.isServingConsumers() ? host : null);
}
}
diff --git
a/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketHost.java
b/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketHost.java
index a3222fa81256..630d85f1ad20 100644
---
a/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketHost.java
+++
b/components/camel-vertx/camel-vertx-websocket/src/main/java/org/apache/camel/component/vertx/websocket/VertxWebsocketHost.java
@@ -51,10 +51,11 @@ public class VertxWebsocketHost {
private final VertxWebsocketHostConfiguration hostConfiguration;
private final VertxWebsocketHostKey hostKey;
+ // guarded by this host's monitor, together with the server it decides the
lifecycle of
private final Map<String, Route> routeRegistry = new HashMap<>();
private final List<VertxWebsocketPeer> connectedPeers = new
CopyOnWriteArrayList<>(); // thread-safe
private final CamelContext camelContext;
- private HttpServer server;
+ private volatile HttpServer server;
private int port = VertxWebsocketConstants.DEFAULT_VERTX_SERVER_PORT;
public VertxWebsocketHost(CamelContext camelContext,
VertxWebsocketHostConfiguration websocketHostConfiguration,
@@ -67,7 +68,7 @@ public class VertxWebsocketHost {
/**
* Sets up a Vert.x route and handler for the WebSocket path specified by
the consumer configuration
*/
- public void connect(VertxWebsocketConsumer consumer) {
+ public synchronized void connect(VertxWebsocketConsumer consumer) {
VertxWebsocketEndpoint endpoint = consumer.getEndpoint();
VertxWebsocketConfiguration configuration =
endpoint.getConfiguration();
@@ -159,10 +160,12 @@ public class VertxWebsocketHost {
/**
* Removes the Vert.x route and handler for the WebSocket path specified
by the consumer configuration
*/
- public void disconnect(String path) {
+ public synchronized void disconnect(String path) {
LOG.info("Disconnected consumer for path {}", path);
Route route = routeRegistry.remove(path);
- route.remove();
+ if (route != null) {
+ route.remove();
+ }
if (routeRegistry.isEmpty()) {
try {
stop();
@@ -175,7 +178,7 @@ public class VertxWebsocketHost {
/**
* Starts a Vert.x HTTP server to host the WebSocket router
*/
- public void start() throws Exception {
+ public synchronized void start() throws Exception {
if (server == null) {
Vertx vertx = hostConfiguration.getVertx();
Router router = hostConfiguration.getRouter();
@@ -213,7 +216,7 @@ public class VertxWebsocketHost {
/**
* Stops a previously started Vert.x HTTP server
*/
- public void stop() throws ExecutionException, InterruptedException {
+ public synchronized void stop() throws ExecutionException,
InterruptedException {
if (server != null) {
LOG.info("Stopping server");
try {
@@ -236,6 +239,14 @@ public class VertxWebsocketHost {
port = VertxWebsocketConstants.DEFAULT_VERTX_SERVER_PORT;
}
+ /**
+ * Whether this host still serves any consumer. Every consumer bound to
the same host and port shares one instance,
+ * so the host outlives the first consumer that stops, and only once the
last one goes is its server stopped.
+ */
+ public synchronized boolean isServingConsumers() {
+ return !routeRegistry.isEmpty();
+ }
+
/**
* Gets all WebSocket peers connected to the Vert.x HTTP sever together
with their associated connection key
*/
diff --git
a/components/camel-vertx/camel-vertx-websocket/src/test/java/org/apache/camel/component/vertx/websocket/VertxWebsocketMultiConsumerLifecycleTest.java
b/components/camel-vertx/camel-vertx-websocket/src/test/java/org/apache/camel/component/vertx/websocket/VertxWebsocketMultiConsumerLifecycleTest.java
new file mode 100644
index 000000000000..224305f70908
--- /dev/null
+++
b/components/camel-vertx/camel-vertx-websocket/src/test/java/org/apache/camel/component/vertx/websocket/VertxWebsocketMultiConsumerLifecycleTest.java
@@ -0,0 +1,132 @@
+/*
+ * 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.vertx.websocket;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Consumers bound to the same host and port share one {@link
VertxWebsocketHost}, so stopping one of them must leave
+ * the others able to serve and, later, to disconnect.
+ */
+public class VertxWebsocketMultiConsumerLifecycleTest extends
VertxWebSocketTestSupport {
+
+ private Map<VertxWebsocketHostKey, VertxWebsocketHost> hostRegistry() {
+ return context.getComponent("vertx-websocket",
VertxWebsocketComponent.class).getVertxHostRegistry();
+ }
+
+ @Test
+ void theHostOutlivesTheFirstConsumerToStop() throws Exception {
+ assertEquals(1, hostRegistry().size());
+
+ context.getRouteController().stopRoute("a");
+
+ // the host still serves route b, so it must still be reachable - it
used to be dropped from the
+ // registry here, which left b unable to ever disconnect and its
server running for good
+ assertEquals(1, hostRegistry().size());
+
+ MockEndpoint mockEndpoint = getMockEndpoint("mock:result");
+ mockEndpoint.expectedBodiesReceived("Hello b");
+
+ template.sendBody("vertx-websocket:localhost:" + port + "/test/b",
"b");
+
+ mockEndpoint.assertIsSatisfied();
+ }
+
+ @Test
+ void theHostGoesWithTheLastConsumerToStop() throws Exception {
+ assertEquals(1, hostRegistry().size());
+
+ context.getRouteController().stopRoute("a");
+ context.getRouteController().stopRoute("b");
+
+ assertTrue(hostRegistry().isEmpty());
+ }
+
+ @Test
+ void stoppingAConsumerTwiceIsHarmless() throws Exception {
+ context.getRouteController().stopRoute("a");
+ context.getRouteController().startRoute("a");
+ context.getRouteController().stopRoute("a");
+
+ assertEquals(1, hostRegistry().size());
+ }
+
+ @Test
+ void consumersSharingAHostCanBeStoppedConcurrently() throws Exception {
+ assertEquals(1, hostRegistry().size());
+
+ List<Throwable> failures = new CopyOnWriteArrayList<>();
+ CountDownLatch startLine = new CountDownLatch(1);
+ CountDownLatch finished = new CountDownLatch(2);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+
+ try {
+ for (String routeId : List.of("a", "b")) {
+ executor.submit(() -> {
+ try {
+ startLine.await();
+ context.getRouteController().stopRoute(routeId);
+ } catch (Throwable t) {
+ failures.add(t);
+ } finally {
+ finished.countDown();
+ }
+ });
+ }
+
+ startLine.countDown();
+ assertTrue(finished.await(30, TimeUnit.SECONDS), "the consumers
did not stop in time");
+ } finally {
+ executor.shutdownNow();
+ }
+
+ assertTrue(failures.isEmpty(), "stopping the consumers concurrently
failed with " + failures);
+ assertTrue(hostRegistry().isEmpty(), "the host outlived both of its
consumers");
+ }
+
+ @Override
+ protected RoutesBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ fromF("vertx-websocket:localhost:%d/test/a", port.getPort())
+ .routeId("a")
+ .setBody(simple("Hello ${body}"))
+ .to("mock:result");
+
+ fromF("vertx-websocket:localhost:%d/test/b", port.getPort())
+ .routeId("b")
+ .setBody(simple("Hello ${body}"))
+ .to("mock:result");
+ }
+ };
+ }
+}