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

Croway 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 186bee616ee0 CAMEL-25197: camel-cli-connector - pluggable WebSocket 
client for the websocket transport (#27210)
186bee616ee0 is described below

commit 186bee616ee031c0058924d2412c2af7a55d72fa
Author: Federico Mariani <[email protected]>
AuthorDate: Thu Oct 1 16:46:01 2026 +0200

    CAMEL-25197: camel-cli-connector - pluggable WebSocket client for the 
websocket transport (#27210)
    
    * CAMEL-25197: camel-cli-connector - pluggable WebSocket client for the 
websocket transport
    
    The websocket transport now leaves the socket I/O to a CliWebSocketClient, 
so that runtimes can plug in
    their own WebSocket client (for example the Quarkus WebSockets Next client, 
or Spring's), while the
    protocol stays in the transport: envelope, hello, actions, snapshots and 
their splitting, single writer,
    send timeout, lone surrogates, heartbeat, reconnect backoff, inbound size 
cap and security checks.
    
    - CliWebSocketClient (connect, Channel, Listener), 
CliWebSocketHandshakeException (HTTP status of a
      rejected handshake) and JdkCliWebSocketClient (the previous JDK code, 
moved)
    - the transport uses the single CliWebSocketClient in the registry, 
otherwise the JDK client;
      camel.cli.websocket.client=jdk always uses the JDK client
    - the client in use is logged and reported in hello.transport
    - incoming frames are parsed on the transport thread, not on the client 
callback thread (which can be
      an event loop)
    - a close or error reported before the client completes the connect is 
handled (reconnects)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
    
    * CAMEL-25197: camel-cli-connector - keep the frames received before the 
connection is reported open
    
    A tool that does not wait for the hello can send an action before the 
client completes connect(): the
    action ran but its result was dropped. Such frames are now kept (up to 64) 
and handled once the
    connection is reported open, after the hello. Frames of a connection that 
is gone are no longer run,
    as their result cannot be sent.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
    
    ---------
    
    Co-authored-by: Claude Opus 5.5 <[email protected]>
---
 .../apache/camel/catalog/docs/cli-connector.adoc   |  21 ++-
 .../src/main/docs/cli-connector.adoc               |  21 ++-
 .../camel/cli/connector/CliWebSocketClient.java    | 110 ++++++++++++++
 .../connector/CliWebSocketHandshakeException.java  |  38 +++++
 .../camel/cli/connector/JdkCliWebSocketClient.java | 160 +++++++++++++++++++++
 .../connector/WebSocketCliConnectorTransport.java  | 158 +++++++++++---------
 .../WebSocketCliConnectorTransportTest.java        | 133 ++++++++++++++++-
 7 files changed, 571 insertions(+), 70 deletions(-)

diff --git 
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/cli-connector.adoc
 
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/cli-connector.adoc
index d4c769f2a9bb..6ec4916195b0 100644
--- 
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/cli-connector.adoc
+++ 
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/cli-connector.adoc
@@ -44,7 +44,8 @@ connector. Choose it with the `camel.cli.transport` property:
 
 The integration connects to the tool, so the tool does not need access to the 
filesystem of the integration
 (containers, remote development clusters), and snapshots are pushed as they 
are collected. The integration does not
-open any port. The JDK WebSocket client is used, so no extra dependency is 
needed.
+open any port. The JDK WebSocket client is used by default, so no extra 
dependency is needed (see
+<<WebSocket client>>).
 
 [source,bash]
 ----
@@ -69,6 +70,8 @@ java -Dcamel.cli.transport=websocket \
 | `camel.cli.websocket.reconnect-max-delay` | `30000` | Maximum delay before 
reconnecting, in millis.
 | `camel.cli.websocket.allow-insecure` | `false` | Allow `ws://` (unencrypted) 
to a host that is not a loopback
   address. Prefer `wss://` or a tunnel.
+| `camel.cli.websocket.client` | `auto` | The WebSocket client: `auto` uses 
the `CliWebSocketClient` from the
+  registry if there is exactly one, otherwise the JDK client; `jdk` always 
uses the JDK client.
 |===
 
 As with Camel Main options, the camelCase form of the keys is accepted too 
(for example
@@ -99,6 +102,22 @@ same JSON the file transport reads and writes:
   accepts messages of that size.
 * The extra `stop` action stops the integration.
 
+=== WebSocket client
+
+The transport does the protocol (frames, snapshots, heartbeat, reconnecting, 
the security checks below) and leaves
+the socket I/O to a `org.apache.camel.cli.connector.CliWebSocketClient`. A 
runtime can provide its own client, for
+example to reuse the WebSocket client and TLS configuration of the 
application, by putting a single bean of that type
+in the registry; the client in use is logged at startup and sent in the 
`transport` field of the `hello` frame.
+
+A client must:
+
+* pass whole text messages to the listener (reassemble messages sent in 
several frames), and accept messages of up
+  to 16 MB (`CliWebSocketClient.MAX_MESSAGE_SIZE`);
+* fail the connection with a `CliWebSocketHandshakeException` holding the HTTP 
status code when the tool rejects the
+  handshake;
+* never block in the listener callbacks, which may run on any thread (for 
example an event loop): the transport
+  parses the frames and sends on its own threads, one message at a time, and 
waits for the returned stages there.
+
 === Security
 
 The WebSocket transport is a development tool: *the tool it connects to gets 
full control of the integration*,
diff --git a/dsl/camel-cli-connector/src/main/docs/cli-connector.adoc 
b/dsl/camel-cli-connector/src/main/docs/cli-connector.adoc
index d4c769f2a9bb..6ec4916195b0 100644
--- a/dsl/camel-cli-connector/src/main/docs/cli-connector.adoc
+++ b/dsl/camel-cli-connector/src/main/docs/cli-connector.adoc
@@ -44,7 +44,8 @@ connector. Choose it with the `camel.cli.transport` property:
 
 The integration connects to the tool, so the tool does not need access to the 
filesystem of the integration
 (containers, remote development clusters), and snapshots are pushed as they 
are collected. The integration does not
-open any port. The JDK WebSocket client is used, so no extra dependency is 
needed.
+open any port. The JDK WebSocket client is used by default, so no extra 
dependency is needed (see
+<<WebSocket client>>).
 
 [source,bash]
 ----
@@ -69,6 +70,8 @@ java -Dcamel.cli.transport=websocket \
 | `camel.cli.websocket.reconnect-max-delay` | `30000` | Maximum delay before 
reconnecting, in millis.
 | `camel.cli.websocket.allow-insecure` | `false` | Allow `ws://` (unencrypted) 
to a host that is not a loopback
   address. Prefer `wss://` or a tunnel.
+| `camel.cli.websocket.client` | `auto` | The WebSocket client: `auto` uses 
the `CliWebSocketClient` from the
+  registry if there is exactly one, otherwise the JDK client; `jdk` always 
uses the JDK client.
 |===
 
 As with Camel Main options, the camelCase form of the keys is accepted too 
(for example
@@ -99,6 +102,22 @@ same JSON the file transport reads and writes:
   accepts messages of that size.
 * The extra `stop` action stops the integration.
 
+=== WebSocket client
+
+The transport does the protocol (frames, snapshots, heartbeat, reconnecting, 
the security checks below) and leaves
+the socket I/O to a `org.apache.camel.cli.connector.CliWebSocketClient`. A 
runtime can provide its own client, for
+example to reuse the WebSocket client and TLS configuration of the 
application, by putting a single bean of that type
+in the registry; the client in use is logged at startup and sent in the 
`transport` field of the `hello` frame.
+
+A client must:
+
+* pass whole text messages to the listener (reassemble messages sent in 
several frames), and accept messages of up
+  to 16 MB (`CliWebSocketClient.MAX_MESSAGE_SIZE`);
+* fail the connection with a `CliWebSocketHandshakeException` holding the HTTP 
status code when the tool rejects the
+  handshake;
+* never block in the listener callbacks, which may run on any thread (for 
example an event loop): the transport
+  parses the frames and sends on its own threads, one message at a time, and 
waits for the returned stages there.
+
 === Security
 
 The WebSocket transport is a development tool: *the tool it connects to gets 
full control of the integration*,
diff --git 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliWebSocketClient.java
 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliWebSocketClient.java
new file mode 100644
index 000000000000..e486216a695a
--- /dev/null
+++ 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliWebSocketClient.java
@@ -0,0 +1,110 @@
+/*
+ * 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.cli.connector;
+
+import java.net.URI;
+import java.util.Map;
+import java.util.concurrent.CompletionStage;
+
+/**
+ * The WebSocket client of the websocket transport 
(<tt>camel.cli.transport=websocket</tt>): only the socket I/O, the
+ * protocol (frames, snapshots, heartbeat, reconnecting, security checks) is 
done by the transport.
+ * <p/>
+ * The transport uses the single bean of this type in the registry, if any, 
otherwise the JDK client
+ * ({@link JdkCliWebSocketClient}). Set 
<tt>camel.cli.websocket.client=jdk</tt> to always use the JDK client.
+ * <p/>
+ * Threads: the transport never has two messages in flight on a {@link 
Channel}, and may wait for the returned stages on
+ * its own threads. The {@link Listener} callbacks may be called on any 
thread, for example an event loop: they do not
+ * block.
+ *
+ * @since 4.23
+ */
+public interface CliWebSocketClient {
+
+    /**
+     * The largest message the client must accept from the tool, in chars (16 
MB).
+     */
+    int MAX_MESSAGE_SIZE = 16 * 1024 * 1024;
+
+    /**
+     * The name of the client, sent to the tool in the <tt>hello</tt> frame 
(<tt>transport</tt>) and logged.
+     */
+    String getName();
+
+    /**
+     * Opens a connection.
+     *
+     * @param  url      the <tt>ws://</tt> or <tt>wss://</tt> url of the tool, 
query string included as given
+     * @param  headers  headers to send with the handshake (such as 
<tt>Authorization</tt>)
+     * @param  listener receives what the tool sends on this connection
+     * @return          completes with the connection once open, or fails with 
a {@link CliWebSocketHandshakeException}
+     *                  when the tool rejects the handshake
+     */
+    CompletionStage<Channel> connect(URI url, Map<String, String> headers, 
Listener listener);
+
+    /**
+     * An open connection.
+     */
+    interface Channel {
+
+        /**
+         * Sends a whole text message.
+         */
+        CompletionStage<?> sendText(String text);
+
+        /**
+         * Sends a ping, the tool answers with a pong.
+         */
+        CompletionStage<?> sendPing();
+
+        /**
+         * Sends a close frame. The transport calls {@link #abort()} if it 
does not complete in time.
+         */
+        CompletionStage<?> close(int code, String reason);
+
+        /**
+         * Closes the connection right away, without waiting for anything.
+         */
+        void abort();
+    }
+
+    /**
+     * Receives what the tool sends. The methods must not block.
+     */
+    interface Listener {
+
+        /**
+         * A whole text message (a message sent in several frames is passed 
once complete).
+         */
+        void onText(String text);
+
+        /**
+         * A pong, answering a {@link Channel#sendPing()}.
+         */
+        void onPong();
+
+        /**
+         * The connection was closed (by the tool, or by the network).
+         */
+        void onClose(int code, String reason);
+
+        /**
+         * The connection failed, it cannot be used anymore.
+         */
+        void onError(Throwable error);
+    }
+}
diff --git 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliWebSocketHandshakeException.java
 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliWebSocketHandshakeException.java
new file mode 100644
index 000000000000..9e5c25da3909
--- /dev/null
+++ 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliWebSocketHandshakeException.java
@@ -0,0 +1,38 @@
+/*
+ * 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.cli.connector;
+
+import java.io.IOException;
+
+/**
+ * The tool rejected the WebSocket handshake with this HTTP status code (for 
example 401 for a wrong token).
+ *
+ * @since 4.23
+ */
+public class CliWebSocketHandshakeException extends IOException {
+
+    private final int statusCode;
+
+    public CliWebSocketHandshakeException(int statusCode, Throwable cause) {
+        super("WebSocket handshake rejected with HTTP " + statusCode, cause);
+        this.statusCode = statusCode;
+    }
+
+    public int getStatusCode() {
+        return statusCode;
+    }
+}
diff --git 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/JdkCliWebSocketClient.java
 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/JdkCliWebSocketClient.java
new file mode 100644
index 000000000000..dfdea7dcaf36
--- /dev/null
+++ 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/JdkCliWebSocketClient.java
@@ -0,0 +1,160 @@
+/*
+ * 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.cli.connector;
+
+import java.io.IOException;
+import java.net.URI;
+import java.net.http.HttpClient;
+import java.net.http.WebSocket;
+import java.net.http.WebSocketHandshakeException;
+import java.nio.ByteBuffer;
+import java.time.Duration;
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
+import java.util.concurrent.CompletionStage;
+
+/**
+ * The default {@link CliWebSocketClient}, with the JDK client 
(<tt>java.net.http</tt>), so no extra dependency is
+ * needed.
+ *
+ * @since 4.23
+ */
+public class JdkCliWebSocketClient implements CliWebSocketClient {
+
+    private static final Duration CONNECT_TIMEOUT = Duration.ofSeconds(10);
+
+    private final HttpClient client = 
HttpClient.newBuilder().connectTimeout(CONNECT_TIMEOUT).build();
+
+    @Override
+    public String getName() {
+        return "jdk";
+    }
+
+    @Override
+    public CompletionStage<Channel> connect(URI url, Map<String, String> 
headers, Listener listener) {
+        WebSocket.Builder builder = 
client.newWebSocketBuilder().connectTimeout(CONNECT_TIMEOUT);
+        headers.forEach(builder::header);
+        CompletableFuture<Channel> answer = new CompletableFuture<>();
+        builder.buildAsync(url, new Adapter(listener)).whenComplete((ws, e) -> 
{
+            if (e != null) {
+                answer.completeExceptionally(translate(e));
+            } else {
+                answer.complete(new JdkChannel(ws));
+            }
+        });
+        return answer;
+    }
+
+    private static Throwable translate(Throwable e) {
+        Throwable cause = e instanceof CompletionException && e.getCause() != 
null ? e.getCause() : e;
+        if (cause instanceof WebSocketHandshakeException he) {
+            return new 
CliWebSocketHandshakeException(he.getResponse().statusCode(), he);
+        }
+        return cause;
+    }
+
+    private record JdkChannel(WebSocket ws) implements Channel {
+
+        @Override
+        public CompletionStage<?> sendText(String text) {
+            return ws.sendText(text, true);
+        }
+
+        @Override
+        public CompletionStage<?> sendPing() {
+            return ws.sendPing(ByteBuffer.allocate(0));
+        }
+
+        @Override
+        public CompletionStage<?> close(int code, String reason) {
+            return ws.sendClose(code, reason);
+        }
+
+        @Override
+        public void abort() {
+            ws.abort();
+        }
+    }
+
+    /**
+     * Reassembles the text messages sent in several frames, one frame at a 
time.
+     */
+    private static final class Adapter implements WebSocket.Listener {
+
+        private final Listener listener;
+        private final StringBuilder partial = new StringBuilder();
+
+        Adapter(Listener listener) {
+            this.listener = listener;
+        }
+
+        @Override
+        public void onOpen(WebSocket webSocket) {
+            webSocket.request(1);
+        }
+
+        @Override
+        public CompletionStage<?> onText(WebSocket webSocket, CharSequence 
data, boolean last) {
+            partial.append(data);
+            if (partial.length() > MAX_MESSAGE_SIZE) {
+                partial.setLength(0);
+                webSocket.abort();
+                listener.onError(new IOException("Message larger than " + 
MAX_MESSAGE_SIZE + " chars"));
+                return null;
+            }
+            if (last) {
+                String text = partial.toString();
+                partial.setLength(0);
+                listener.onText(text);
+            }
+            webSocket.request(1);
+            return null;
+        }
+
+        @Override
+        public CompletionStage<?> onBinary(WebSocket webSocket, ByteBuffer 
data, boolean last) {
+            // not part of the protocol
+            webSocket.request(1);
+            return null;
+        }
+
+        @Override
+        public CompletionStage<?> onPing(WebSocket webSocket, ByteBuffer 
message) {
+            // the JDK replies with a pong
+            return WebSocket.Listener.super.onPing(webSocket, message);
+        }
+
+        @Override
+        public CompletionStage<?> onPong(WebSocket webSocket, ByteBuffer 
message) {
+            listener.onPong();
+            webSocket.request(1);
+            return null;
+        }
+
+        @Override
+        public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, 
String reason) {
+            listener.onClose(statusCode, reason);
+            return null;
+        }
+
+        @Override
+        public void onError(WebSocket webSocket, Throwable error) {
+            listener.onError(error);
+        }
+    }
+}
diff --git 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
index 2eed2a483428..7ce0700d54d7 100644
--- 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
+++ 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
@@ -17,18 +17,14 @@
 package org.apache.camel.cli.connector;
 
 import java.net.URI;
-import java.net.http.HttpClient;
-import java.net.http.WebSocket;
-import java.net.http.WebSocketHandshakeException;
-import java.nio.ByteBuffer;
 import java.nio.charset.StandardCharsets;
-import java.time.Duration;
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Locale;
 import java.util.Map;
-import java.util.concurrent.CompletionStage;
+import java.util.Set;
+import java.util.concurrent.CompletionException;
 import java.util.concurrent.Executors;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.RejectedExecutionException;
@@ -51,8 +47,8 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 /**
- * Transport that dials out to a developer tool over a WebSocket (JDK client), 
for tools that do not share the
- * filesystem of the integration or need events pushed instead of polled.
+ * Transport that dials out to a developer tool over a WebSocket, for tools 
that do not share the filesystem of the
+ * integration or need events pushed instead of polled.
  * <p/>
  * Every frame is a JSON envelope <tt>{"v":1,"type":...}</tt>. The tool sends 
<tt>action</tt> frames holding the same
  * action JSON the Camel CLI writes to its action file, and gets 
<tt>result</tt> frames back correlated by
@@ -61,9 +57,12 @@ import org.slf4j.LoggerFactory;
  * <p/>
  * The connection is kept alive with pings, and re-established with an 
exponential backoff when it is lost.
  * <p/>
+ * The socket I/O is done by a {@link CliWebSocketClient}: the single one in 
the registry, if any, otherwise the JDK
+ * client (<tt>camel.cli.websocket.client=jdk</tt> always uses the JDK client).
+ * <p/>
  * Threads: actions run one at a time on their own thread (the dispatcher is 
not thread-safe, and an action can block
- * for a long time); connecting, snapshots, heartbeats and every write to the 
socket run on a second thread, so the JDK
- * WebSocket never has two sends in flight.
+ * for a long time); connecting, parsing the incoming frames, snapshots, 
heartbeats and every write to the socket run on
+ * a second thread, so the client never has two sends in flight, and its 
callbacks never wait.
  */
 public class WebSocketCliConnectorTransport extends ServiceSupport implements 
CliConnectorTransport {
 
@@ -72,7 +71,8 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
     private static final Logger LOG = 
LoggerFactory.getLogger(WebSocketCliConnectorTransport.class);
     private static final long SEND_TIMEOUT = 10000;
     private static final long STABLE_CONNECTION = 10000;
-    private static final int MAX_FRAME_SIZE = 16 * 1024 * 1024;
+    private static final int MAX_FRAME_SIZE = 
CliWebSocketClient.MAX_MESSAGE_SIZE;
+    private static final int NORMAL_CLOSURE = 1000;
     private static final int MAX_PENDING_ACTIONS = 64;
     // WebSocket servers commonly refuse messages over 256 KB (Vert.x, 
Quarkus): trace and receive snapshots can be
     // much larger (after a reconnect they hold every retained message), so 
they are split into frames of this size
@@ -91,7 +91,7 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
     private long heartbeatInterval;
 
     private String threadNamePattern;
-    private HttpClient client;
+    private CliWebSocketClient client;
     private ThreadPoolExecutor actions;
     private ScheduledExecutorService scheduler;
     // fields below are only used from the scheduler thread, except the 
volatile ones
@@ -156,7 +156,8 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
             LOG.warn("Camel CLI connector uses an unencrypted connection to a 
remote host; use wss:// or a tunnel");
         }
 
-        client = 
HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(10)).build();
+        client = resolveClient();
+        LOG.info("Camel CLI connector uses the {} WebSocket client", 
client.getName());
         // Camel's thread factory (naming, virtual threads when enabled) but 
not Camel's thread pools: these threads must
         // keep running while Camel is stopping, to report it and to send the 
close frame
         threadNamePattern = 
camelContext.getExecutorServiceManager().getThreadNamePattern();
@@ -192,6 +193,24 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
         scheduler.execute(this::connect);
     }
 
+    private CliWebSocketClient resolveClient() {
+        String name = property("camel.cli.websocket.client", "auto");
+        if ("jdk".equalsIgnoreCase(name)) {
+            return new JdkCliWebSocketClient();
+        }
+        if (!"auto".equalsIgnoreCase(name)) {
+            throw new IllegalArgumentException("camel.cli.websocket.client 
must be auto or jdk: " + name);
+        }
+        Set<CliWebSocketClient> found = 
camelContext.getRegistry().findByType(CliWebSocketClient.class);
+        if (found.size() == 1) {
+            return found.iterator().next();
+        }
+        if (found.size() > 1) {
+            LOG.warn("Camel CLI connector found {} WebSocket clients in the 
registry, using the JDK client", found.size());
+        }
+        return new JdkCliWebSocketClient();
+    }
+
     @Override
     public void updateDelay(int delay) {
         // camel-cli-debug makes it faster so breakpoints show up quickly: 
only the debug snapshot needs that
@@ -232,7 +251,7 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
                     Connection c = connection;
                     if (c != null) {
                         connection = null;
-                        c.close(WebSocket.NORMAL_CLOSURE, "stopping");
+                        c.close(NORMAL_CLOSURE, "stopping");
                     }
                 }).get(SEND_TIMEOUT, TimeUnit.MILLISECONDS);
             } catch (Exception e) {
@@ -255,13 +274,16 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
             return;
         }
         try {
-            WebSocket.Builder builder = 
client.newWebSocketBuilder().connectTimeout(Duration.ofSeconds(10));
+            Map<String, String> headers = new HashMap<>();
             if (token != null && !token.isBlank()) {
-                builder.header("Authorization", "Bearer " + token);
+                headers.put("Authorization", "Bearer " + token);
             }
-            builder.buildAsync(url, new Connection()).whenComplete((ws, e) -> {
+            Connection c = new Connection();
+            client.connect(url, headers, c).whenComplete((channel, e) -> {
                 if (e != null) {
                     execute(() -> connectFailed(e));
+                } else {
+                    execute(() -> opened(c, channel));
                 }
             });
         } catch (Exception e) {
@@ -271,9 +293,9 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
 
     private void connectFailed(Throwable e) {
         failures++;
-        Throwable cause = e.getCause() != null ? e.getCause() : e;
-        if (cause instanceof WebSocketHandshakeException he) {
-            int code = he.getResponse().statusCode();
+        Throwable cause = e instanceof CompletionException && e.getCause() != 
null ? e.getCause() : e;
+        if (cause instanceof CliWebSocketHandshakeException he) {
+            int code = he.getStatusCode();
             if (code == 401 || code == 403) {
                 LOG.warn("Camel CLI connector was rejected by {} (HTTP {}): 
check camel.cli.websocket.token", where(), code);
             } else {
@@ -287,23 +309,38 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
         reconnect();
     }
 
-    private void opened(Connection c) {
+    private void opened(Connection c, CliWebSocketClient.Channel channel) {
+        c.channel = channel;
         if (stopping) {
-            c.ws.abort();
+            channel.abort();
+            return;
+        }
+        if (c.lost) {
+            // closed or failed before the client reported it open
+            channel.abort();
+            failures++;
+            reconnect();
             return;
         }
+        c.openedAt = System.currentTimeMillis();
+        c.lastSeen = c.openedAt;
         connection = c;
         LOG.info("Camel CLI connector connected to {}", where());
         sayHello(c);
+        // what the tool sent before the client reported the connection open, 
in order
+        List<String> early = new ArrayList<>(c.early);
+        c.early.clear();
+        early.forEach(text -> onFrame(c, text));
     }
 
     private void closed(Connection c, String reason) {
         if (connection != c) {
-            // an old connection, or already handled
+            // an old connection, already handled, or not open yet (see opened)
+            c.lost = true;
             return;
         }
         connection = null;
-        c.ws.abort();
+        c.channel.abort();
         // a connection that drops right away counts as a failure, so a tool 
that keeps closing us is not hammered
         if (System.currentTimeMillis() - c.openedAt < STABLE_CONNECTION) {
             failures++;
@@ -339,7 +376,7 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
             // the tool went away without closing the connection (sleep, 
network change, crash)
             closed(c, "no heartbeat");
         } else {
-            c.ws.sendPing(ByteBuffer.allocate(0));
+            c.channel.sendPing();
         }
     }
 
@@ -352,7 +389,7 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
         JsonObject frame = envelope("hello");
         frame.put("camelVersion", camelContext.getVersion());
         frame.put("name", camelContext.getName());
-        frame.put("transport", "jdk");
+        frame.put("transport", client.getName());
         try {
             JsonObject status = snapshots.status();
             frame.put("runtime", status.get("runtime"));
@@ -463,7 +500,7 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
             return;
         }
         try {
-            c.ws.sendText(wellFormed(frame.toJson()), true).get(SEND_TIMEOUT, 
TimeUnit.MILLISECONDS);
+            
c.channel.sendText(wellFormed(frame.toJson())).toCompletableFuture().get(SEND_TIMEOUT,
 TimeUnit.MILLISECONDS);
         } catch (Exception e) {
             // a send that fails or hangs leaves the socket unusable
             closed(c, "send failed: " + describe(e));
@@ -473,6 +510,14 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
     // ---- incoming frames ----
 
     private void onFrame(Connection c, String text) {
+        if (c != connection) {
+            if (c.channel == null && !c.lost && c.early.size() < 
MAX_PENDING_ACTIONS) {
+                // the client has not reported the connection open yet: kept 
until it does (see opened)
+                c.early.add(text);
+            }
+            // otherwise the connection is gone, and the result could not be 
sent
+            return;
+        }
         String requestId = null;
         try {
             JsonObject frame = (JsonObject) Jsoner.deserialize(text);
@@ -661,75 +706,54 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
     }
 
     /**
-     * One WebSocket connection, and what has been sent on it.
+     * One WebSocket connection, and what has been sent on it. The callbacks 
come from the client, on any thread.
      */
-    private final class Connection implements WebSocket.Listener {
+    private final class Connection implements CliWebSocketClient.Listener {
 
-        private WebSocket ws;
-        private final StringBuilder partial = new StringBuilder();
-        private final long openedAt = System.currentTimeMillis();
+        // set on the scheduler thread, before the connection is used
+        private CliWebSocketClient.Channel channel;
+        private boolean lost;
+        private long openedAt = System.currentTimeMillis();
         private volatile long lastSeen = openedAt;
         private boolean helloSent;
         private long lastTraceUid;
         private long lastReceiveUid;
         private final Map<String, String> lastSent = new HashMap<>();
+        // frames received before the connection is reported open, only used 
on the scheduler thread
+        private final List<String> early = new ArrayList<>();
 
         @Override
-        public void onOpen(WebSocket webSocket) {
-            this.ws = webSocket;
-            webSocket.request(1);
-            execute(() -> opened(this));
-        }
-
-        @Override
-        public CompletionStage<?> onText(WebSocket webSocket, CharSequence 
data, boolean last) {
+        public void onText(String text) {
             lastSeen = System.currentTimeMillis();
-            partial.append(data);
-            if (partial.length() > MAX_FRAME_SIZE) {
-                partial.setLength(0);
+            if (text.length() > MAX_FRAME_SIZE) {
                 execute(() -> closed(this, "frame larger than " + 
MAX_FRAME_SIZE + " chars"));
-                return null;
-            }
-            if (last) {
-                String text = partial.toString();
-                partial.setLength(0);
-                onFrame(this, text);
+                return;
             }
-            webSocket.request(1);
-            return null;
-        }
-
-        @Override
-        public CompletionStage<?> onPing(WebSocket webSocket, ByteBuffer 
message) {
-            lastSeen = System.currentTimeMillis();
-            // the JDK replies with a pong
-            return WebSocket.Listener.super.onPing(webSocket, message);
+            execute(() -> onFrame(this, text));
         }
 
         @Override
-        public CompletionStage<?> onPong(WebSocket webSocket, ByteBuffer 
message) {
+        public void onPong() {
             lastSeen = System.currentTimeMillis();
-            webSocket.request(1);
-            return null;
         }
 
         @Override
-        public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, 
String reason) {
-            execute(() -> closed(this, "closed by the tool: " + statusCode + 
(reason.isEmpty() ? "" : " " + reason)));
-            return null;
+        public void onClose(int statusCode, String reason) {
+            execute(() -> closed(this, "closed by the tool: " + statusCode
+                                       + (reason == null || reason.isEmpty() ? 
"" : " " + reason)));
         }
 
         @Override
-        public void onError(WebSocket webSocket, Throwable error) {
+        public void onError(Throwable error) {
             LOG.debug("Camel CLI connector websocket error: {}", 
error.getMessage(), error);
             execute(() -> closed(this, "error: " + describe(error)));
         }
 
         void close(int code, String reason) {
             try {
-                ws.sendClose(code, reason).get(SEND_TIMEOUT, 
TimeUnit.MILLISECONDS);
+                channel.close(code, 
reason).toCompletableFuture().get(SEND_TIMEOUT, TimeUnit.MILLISECONDS);
             } catch (Exception e) {
-                ws.abort();
+                channel.abort();
             }
         }
 
diff --git 
a/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
 
b/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
index 80933ff5e7f2..f285369027a0 100644
--- 
a/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
+++ 
b/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
@@ -16,9 +16,12 @@
  */
 package org.apache.camel.cli.connector;
 
+import java.net.URI;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionStage;
 import java.util.concurrent.CopyOnWriteArrayList;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.TimeUnit;
@@ -304,6 +307,95 @@ class WebSocketCliConnectorTransportTest extends 
CamelTestSupport {
                 "camel.cli.websocket.token must be set when 
camel.cli.websocket.url is not a loopback address");
     }
 
+    @Test
+    void usesTheJdkClientByDefault() throws Exception {
+        startConnector();
+
+        assertThat(tool.awaitFrame(f -> 
"hello".equals(f.getString("type"))).getString("transport")).isEqualTo("jdk");
+    }
+
+    @Test
+    void usesTheClientFromTheRegistry() throws Exception {
+        RecordingClient client = new RecordingClient();
+        context.getRegistry().bind("myClient", client);
+        startConnector();
+
+        assertThat(tool.awaitFrame(f -> 
"hello".equals(f.getString("type"))).getString("transport")).isEqualTo("test");
+        assertThat(client.connects).hasValue(1);
+
+        tool.send(action("r1", "send", "endpoint", "direct:hello", "body", 
"World", "exchangePattern", "InOut"));
+        JsonObject result = tool.awaitResult("r1");
+        assertThat(result.getBoolean("ok")).isTrue();
+        assertThat(map(result, "result").toJson()).contains("Hello World");
+    }
+
+    @Test
+    void usesTheJdkClientWhenAsked() throws Exception {
+        RecordingClient client = new RecordingClient();
+        context.getRegistry().bind("myClient", client);
+        property("camel.cli.websocket.client", "jdk");
+        startConnector();
+
+        assertThat(tool.awaitFrame(f -> 
"hello".equals(f.getString("type"))).getString("transport")).isEqualTo("jdk");
+        assertThat(client.connects).hasValue(0);
+    }
+
+    @Test
+    void refusesAnUnknownClient() {
+        property("camel.cli.websocket.client", "netty");
+
+        
assertThatThrownBy(this::startConnector).hasMessage("camel.cli.websocket.client 
must be auto or jdk: netty");
+    }
+
+    @Test
+    void sendsACloseFrameWhenStopping() throws Exception {
+        startConnector();
+        tool.awaitFrame(f -> "hello".equals(f.getString("type")));
+
+        connector.stop();
+        connector = null;
+
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertThat(tool.closeCodes).containsExactly((short) 1000));
+    }
+
+    @Test
+    void staysConnectedWhileTheToolAnswersPings() throws Exception {
+        // the tool never sends anything: only its pongs show it is alive
+        property("camel.cli.websocket.heartbeatInterval", "200");
+        startConnector();
+        tool.awaitFrame(f -> "hello".equals(f.getString("type")));
+
+        await().during(1500, TimeUnit.MILLISECONDS).atMost(5, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertThat(tool.handshakes).hasValue(1));
+    }
+
+    @Test
+    void reconnectsWhenClosedBeforeTheClientReportsItOpen() throws Exception {
+        // a client calling the listener on its own threads can report the 
close before the connection itself
+        RecordingClient client = new RecordingClient();
+        client.closeFirstEarly = true;
+        context.getRegistry().bind("myClient", client);
+        startConnector();
+
+        tool.awaitFrame(f -> "hello".equals(f.getString("type")));
+        assertThat(client.connects).hasValue(2);
+    }
+
+    @Test
+    void runsActionsReceivedBeforeTheClientReportsTheConnectionOpen() throws 
Exception {
+        // a tool that does not wait for the hello: its first action can 
arrive before the client completes connect()
+        RecordingClient client = new RecordingClient();
+        client.textFirstEarly = action("r1", "send", "endpoint", 
"direct:hello", "body", "Early", "exchangePattern", "InOut")
+                .toJson();
+        context.getRegistry().bind("myClient", client);
+        startConnector();
+
+        tool.awaitFrame(f -> "hello".equals(f.getString("type")));
+        JsonObject result = tool.awaitResult("r1");
+        assertThat(result.getBoolean("ok")).isTrue();
+        assertThat(map(result, "result").toJson()).contains("Hello Early");
+    }
+
     private void startConnector() {
         connector = new LocalCliConnector(new DefaultCliConnectorFactory()) {
             @Override
@@ -337,6 +429,39 @@ class WebSocketCliConnectorTransportTest extends 
CamelTestSupport {
         return new JsonObject(Map.of("v", 1, "type", "action", "requestId", 
requestId, "action", action));
     }
 
+    /**
+     * A client from the registry: the JDK client under another name.
+     */
+    private static class RecordingClient implements CliWebSocketClient {
+
+        final JdkCliWebSocketClient delegate = new JdkCliWebSocketClient();
+        final AtomicInteger connects = new AtomicInteger();
+        volatile boolean closeFirstEarly;
+        volatile String textFirstEarly;
+
+        @Override
+        public String getName() {
+            return "test";
+        }
+
+        @Override
+        public CompletionStage<Channel> connect(URI url, Map<String, String> 
headers, Listener listener) {
+            boolean first = connects.incrementAndGet() == 1;
+            return delegate.connect(url, headers, 
listener).thenCompose(channel -> {
+                if (first && closeFirstEarly) {
+                    listener.onClose(1001, "gone early");
+                }
+                if (first && textFirstEarly != null) {
+                    listener.onText(textFirstEarly);
+                    // and reports the connection open well after it
+                    return CompletableFuture.supplyAsync(() -> channel,
+                            CompletableFuture.delayedExecutor(500, 
TimeUnit.MILLISECONDS));
+                }
+                return CompletableFuture.completedFuture(channel);
+            });
+        }
+    }
+
     /**
      * Plays the tool: a WebSocket server the connector dials out to.
      */
@@ -345,6 +470,7 @@ class WebSocketCliConnectorTransportTest extends 
CamelTestSupport {
         final BlockingQueue<JsonObject> frames = new LinkedBlockingQueue<>();
         final List<ServerWebSocket> sockets = new CopyOnWriteArrayList<>();
         final AtomicInteger handshakes = new AtomicInteger();
+        final List<Short> closeCodes = new CopyOnWriteArrayList<>();
         volatile String requiredToken;
         int port;
         private Vertx vertx;
@@ -371,7 +497,12 @@ class WebSocketCliConnectorTransportTest extends 
CamelTestSupport {
                                 throw new IllegalStateException(e);
                             }
                         });
-                        ws.closeHandler(v -> sockets.remove(ws));
+                        ws.closeHandler(v -> {
+                            sockets.remove(ws);
+                            if (ws.closeStatusCode() != null) {
+                                closeCodes.add(ws.closeStatusCode());
+                            }
+                        });
                     })
                     .listen(port, 
"127.0.0.1").toCompletionStage().toCompletableFuture().get(10, 
TimeUnit.SECONDS);
             this.port = server.actualPort();

Reply via email to