jamesnetherton commented on code in PR #9264:
URL: https://github.com/apache/camel-quarkus/pull/9264#discussion_r4156811566


##########
extensions-jvm/cli-connector/deployment/src/main/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorProcessor.java:
##########
@@ -51,6 +59,29 @@ CamelBeanBuildItem 
camelBeanBuildItem(CamelCliConnectorRecorder recorder) {
                 recorder.createCliConnectorFactory(Version.getVersion()));
     }
 
+    /**
+     * The WebSocket transport uses the WebSockets Next client when the 
application has it, the JDK client otherwise.
+     */
+    @BuildStep
+    void webSocketClient(
+            Capabilities capabilities,
+            BuildProducer<AdditionalBeanBuildItem> additionalBeans,
+            BuildProducer<RunTimeConfigurationDefaultBuildItem> 
configDefaults) {
+        if (capabilities.isPresent(Capability.WEBSOCKETS_NEXT)) {
+            // the back-pressure of the WebSockets Next client (since Quarkus 
3.40) fetches one frame per message
+            // received, so every message received in several frames (over 64 
KB) leaves fewer frames to fetch, until
+            // the connection stops reading: no back-pressure by default, as 
before Quarkus 3.40
+            // TODO: Remove when 
https://github.com/quarkusio/quarkus/issues/57079 is fixed
+            configDefaults.produce(new RunTimeConfigurationDefaultBuildItem(
+                    "quarkus.websockets-next.client.max-pending-messages", 
"0"));

Review Comment:
   This turns off back-pressure for every WebSockets Next client in the 
application whenever both extensions are present, including when the connector 
is on the file transport or `camel.cli.websocket.client=jdk`, and in 
production. I see it is documented and tracked in quarkusio/quarkus#57079, but 
until that is fixed an application's own clients silently lose the default 
limit of 256 pending messages.
   
   Would it be safer to leave the Quarkus default alone, and instead log a 
warning at startup when the WebSocket transport is in use with the WebSockets 
Next client and `max-pending-messages` is not `0`?



##########
extensions-jvm/cli-connector/runtime/src/main/doc/usage.adoc:
##########
@@ -0,0 +1,38 @@
+=== WebSocket transport
+
+With `camel.cli.transport=websocket`, the application dials out to a developer 
tool over a WebSocket instead of
+exchanging files with the Camel CLI (see the
+xref:{cq-camel-components}:others:cli-connector.adoc[CLI Connector] 
documentation for the protocol and the options).
+
+[source,properties]
+----
+camel.cli.transport = websocket
+camel.cli.websocket.url = ws://127.0.0.1:8000/connect?executionId=run-1
+----
+
+The tool gets full control of the application: the transport refuses to start 
with the `prod` profile
+(`camel.main.profile=prod`).

Review Comment:
   This guard checks the Camel profile, and Camel Quarkus only sets that in dev 
mode (`CamelMainRecorder.customizeDevModeCamelMain`). A packaged application 
running under the Quarkus `prod` profile has no Camel profile unless 
`camel.main.profile` is set explicitly, so the transport still starts if 
`camel.cli.transport=websocket` is left unscoped in `application.properties`.
   
   Could `QuarkusLocalCliConnector` refuse the WebSocket transport in 
`LaunchMode.NORMAL`? Otherwise the docs should say the guard only applies when 
`camel.main.profile=prod` is set explicitly, since readers will assume the 
Quarkus `prod` profile is covered.



##########
extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusCliWebSocketClient.java:
##########
@@ -0,0 +1,156 @@
+/*
+ * 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.quarkus.component.cli.connector;
+
+import java.net.URI;
+import java.util.Map;
+import java.util.OptionalInt;
+import java.util.concurrent.CompletionStage;
+
+import io.quarkus.arc.Arc;
+import io.quarkus.arc.InstanceHandle;
+import io.quarkus.websockets.next.BasicWebSocketConnector;
+import io.quarkus.websockets.next.BasicWebSocketConnector.ExecutionModel;
+import io.quarkus.websockets.next.CloseReason;
+import io.quarkus.websockets.next.WebSocketClientConnection;
+import io.vertx.core.buffer.Buffer;
+import io.vertx.core.http.UpgradeRejectedException;
+import jakarta.annotation.PostConstruct;
+import jakarta.inject.Inject;
+import org.apache.camel.cli.connector.CliWebSocketClient;
+import org.apache.camel.cli.connector.CliWebSocketHandshakeException;
+import org.eclipse.microprofile.config.ConfigProvider;
+import org.jboss.logging.Logger;
+
+/**
+ * The {@link CliWebSocketClient} of the Camel CLI connector WebSocket 
transport, with the Quarkus WebSockets Next
+ * client: registered as a bean only when the application has {@code 
quarkus-websockets-next}.
+ * <p/>
+ * The callbacks run on the Vert.x event loop: they only hand over to the 
transport, which parses and sends on its own
+ * threads.
+ */
+public class QuarkusCliWebSocketClient implements CliWebSocketClient {
+
+    static final String NAME = "quarkus-websockets-next";
+    private static final Logger LOG = 
Logger.getLogger(QuarkusCliWebSocketClient.class);
+    private static final String MAX_MESSAGE_SIZE_KEY = 
"quarkus.websockets-next.client.max-message-size";
+    private static final int GOING_AWAY = 1001;
+
+    @Inject
+    CamelCliConnectorRunTimeConfig config;
+
+    @PostConstruct
+    void checkMaxMessageSize() {
+        // set globally, it replaces the size of every connector (see 
customizeOptions)
+        OptionalInt size = 
ConfigProvider.getConfig().getOptionalValue(MAX_MESSAGE_SIZE_KEY, Integer.class)
+                .map(OptionalInt::of).orElse(OptionalInt.empty());
+        if (size.isPresent() && size.getAsInt() < MAX_MESSAGE_SIZE) {
+            LOG.warnf("%s=%d: the Camel CLI connector cannot receive actions 
larger than that (up to %d supported)",
+                    MAX_MESSAGE_SIZE_KEY, size.getAsInt(), MAX_MESSAGE_SIZE);
+        }
+    }
+
+    @Override
+    public String getName() {
+        return NAME;
+    }
+
+    @Override
+    public CompletionStage<Channel> connect(URI url, Map<String, String> 
headers, Listener listener) {
+        // a new connector for each connection, released when it is closed
+        InstanceHandle<BasicWebSocketConnector> handle = 
Arc.container().instance(BasicWebSocketConnector.class);
+        String path = url.getRawPath() != null ? url.getRawPath() : "";
+        int last = path.lastIndexOf('/') + 1;
+        BasicWebSocketConnector connector = handle.get()
+                .baseUri(baseUri(url, path.substring(0, last)))
+                // the last segment as the path: WebSockets Next adds a '/' 
after the base uri path otherwise
+                // TODO: Remove when 
https://github.com/quarkusio/quarkus/issues/57081 is fixed
+                .path(path.substring(last))
+                .executionModel(ExecutionModel.NON_BLOCKING)
+                // not the frame size: Vert.x also splits the messages it 
sends into frames of that size, and servers
+                // commonly refuse frames over 64 KB

Review Comment:
   Understood that raising the frame size also changes how outgoing messages 
are split. The consequence, though, is that a tool sending an unfragmented 
message over 64 KB gets its connection dropped with this client, while the same 
tool works with the JDK client. So adding `quarkus-websockets-next` to an 
application can break a working setup. I'm not sure that most WebSocket 
libraries fragment by default; several send one frame per message. The tests 
can't catch this, because the tool server is Vert.x, which fragments at 64 KB.
   
   Is that limitation acceptable for the tools this targets?
   
   Minor, related: this limit is counted in bytes while 
`CliWebSocketClient.MAX_MESSAGE_SIZE` is defined in chars, so a non-ASCII 
message under the limit can still be rejected.



##########
extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusCliWebSocketClient.java:
##########
@@ -0,0 +1,156 @@
+/*
+ * 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.quarkus.component.cli.connector;
+
+import java.net.URI;
+import java.util.Map;
+import java.util.OptionalInt;
+import java.util.concurrent.CompletionStage;
+
+import io.quarkus.arc.Arc;
+import io.quarkus.arc.InstanceHandle;
+import io.quarkus.websockets.next.BasicWebSocketConnector;
+import io.quarkus.websockets.next.BasicWebSocketConnector.ExecutionModel;
+import io.quarkus.websockets.next.CloseReason;
+import io.quarkus.websockets.next.WebSocketClientConnection;
+import io.vertx.core.buffer.Buffer;
+import io.vertx.core.http.UpgradeRejectedException;
+import jakarta.annotation.PostConstruct;
+import jakarta.inject.Inject;
+import org.apache.camel.cli.connector.CliWebSocketClient;
+import org.apache.camel.cli.connector.CliWebSocketHandshakeException;
+import org.eclipse.microprofile.config.ConfigProvider;
+import org.jboss.logging.Logger;
+
+/**
+ * The {@link CliWebSocketClient} of the Camel CLI connector WebSocket 
transport, with the Quarkus WebSockets Next
+ * client: registered as a bean only when the application has {@code 
quarkus-websockets-next}.
+ * <p/>
+ * The callbacks run on the Vert.x event loop: they only hand over to the 
transport, which parses and sends on its own
+ * threads.
+ */
+public class QuarkusCliWebSocketClient implements CliWebSocketClient {

Review Comment:
   An open question, not a change request: did you consider a plain Vert.x 
`WebSocketClient` from the injected `Vertx` instead of WebSockets Next?
   
   This client has to work around WebSockets Next in several places: the global 
back-pressure default, no `onClose`, the path split and `%25` re-escaping in 
`baseUri()`, and an `abort()` that cannot drop the socket. Three of those now 
have upstream issues, so it may be fine to wait for them. The Vert.x client 
would give per-client options, a close handler with the status code and the raw 
request URI, with TLS still coming from the TLS registry. Vert.x is already on 
the classpath through `camel-quarkus-console` (`quarkus-vertx-http`), so it 
would also remove the optional dependency, the capability check and the by-name 
bean.



##########
extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusCliWebSocketClient.java:
##########
@@ -0,0 +1,156 @@
+/*
+ * 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.quarkus.component.cli.connector;
+
+import java.net.URI;
+import java.util.Map;
+import java.util.OptionalInt;
+import java.util.concurrent.CompletionStage;
+
+import io.quarkus.arc.Arc;
+import io.quarkus.arc.InstanceHandle;
+import io.quarkus.websockets.next.BasicWebSocketConnector;
+import io.quarkus.websockets.next.BasicWebSocketConnector.ExecutionModel;
+import io.quarkus.websockets.next.CloseReason;
+import io.quarkus.websockets.next.WebSocketClientConnection;
+import io.vertx.core.buffer.Buffer;
+import io.vertx.core.http.UpgradeRejectedException;
+import jakarta.annotation.PostConstruct;
+import jakarta.inject.Inject;
+import org.apache.camel.cli.connector.CliWebSocketClient;
+import org.apache.camel.cli.connector.CliWebSocketHandshakeException;
+import org.eclipse.microprofile.config.ConfigProvider;
+import org.jboss.logging.Logger;
+
+/**
+ * The {@link CliWebSocketClient} of the Camel CLI connector WebSocket 
transport, with the Quarkus WebSockets Next
+ * client: registered as a bean only when the application has {@code 
quarkus-websockets-next}.
+ * <p/>
+ * The callbacks run on the Vert.x event loop: they only hand over to the 
transport, which parses and sends on its own
+ * threads.
+ */
+public class QuarkusCliWebSocketClient implements CliWebSocketClient {
+
+    static final String NAME = "quarkus-websockets-next";
+    private static final Logger LOG = 
Logger.getLogger(QuarkusCliWebSocketClient.class);
+    private static final String MAX_MESSAGE_SIZE_KEY = 
"quarkus.websockets-next.client.max-message-size";
+    private static final int GOING_AWAY = 1001;
+
+    @Inject
+    CamelCliConnectorRunTimeConfig config;
+
+    @PostConstruct
+    void checkMaxMessageSize() {
+        // set globally, it replaces the size of every connector (see 
customizeOptions)
+        OptionalInt size = 
ConfigProvider.getConfig().getOptionalValue(MAX_MESSAGE_SIZE_KEY, Integer.class)
+                .map(OptionalInt::of).orElse(OptionalInt.empty());
+        if (size.isPresent() && size.getAsInt() < MAX_MESSAGE_SIZE) {
+            LOG.warnf("%s=%d: the Camel CLI connector cannot receive actions 
larger than that (up to %d supported)",
+                    MAX_MESSAGE_SIZE_KEY, size.getAsInt(), MAX_MESSAGE_SIZE);
+        }
+    }
+
+    @Override
+    public String getName() {
+        return NAME;
+    }
+
+    @Override
+    public CompletionStage<Channel> connect(URI url, Map<String, String> 
headers, Listener listener) {
+        // a new connector for each connection, released when it is closed
+        InstanceHandle<BasicWebSocketConnector> handle = 
Arc.container().instance(BasicWebSocketConnector.class);
+        String path = url.getRawPath() != null ? url.getRawPath() : "";
+        int last = path.lastIndexOf('/') + 1;
+        BasicWebSocketConnector connector = handle.get()
+                .baseUri(baseUri(url, path.substring(0, last)))
+                // the last segment as the path: WebSockets Next adds a '/' 
after the base uri path otherwise
+                // TODO: Remove when 
https://github.com/quarkusio/quarkus/issues/57081 is fixed
+                .path(path.substring(last))
+                .executionModel(ExecutionModel.NON_BLOCKING)
+                // not the frame size: Vert.x also splits the messages it 
sends into frames of that size, and servers
+                // commonly refuse frames over 64 KB
+                .customizeOptions((connect, client) -> 
client.setMaxMessageSize(MAX_MESSAGE_SIZE))

Review Comment:
   No connect or handshake timeout is set here, whereas the JDK client uses 10 
s for both. If the tool accepts the TCP connection but never answers the 
upgrade (paused in a debugger, or a port-forward with nothing behind it), I 
don't think `connect()` ever completes. As far as I can tell the transport only 
reconnects when that stage completes, so the connector would stay stuck until 
the application restarts.
   
   Suggest setting `connect.setTimeout(...)` and 
`client.setConnectTimeout(...)` here to match the JDK client.



##########
extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusCliWebSocketClient.java:
##########
@@ -0,0 +1,156 @@
+/*
+ * 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.quarkus.component.cli.connector;
+
+import java.net.URI;
+import java.util.Map;
+import java.util.OptionalInt;
+import java.util.concurrent.CompletionStage;
+
+import io.quarkus.arc.Arc;
+import io.quarkus.arc.InstanceHandle;
+import io.quarkus.websockets.next.BasicWebSocketConnector;
+import io.quarkus.websockets.next.BasicWebSocketConnector.ExecutionModel;
+import io.quarkus.websockets.next.CloseReason;
+import io.quarkus.websockets.next.WebSocketClientConnection;
+import io.vertx.core.buffer.Buffer;
+import io.vertx.core.http.UpgradeRejectedException;
+import jakarta.annotation.PostConstruct;
+import jakarta.inject.Inject;
+import org.apache.camel.cli.connector.CliWebSocketClient;
+import org.apache.camel.cli.connector.CliWebSocketHandshakeException;
+import org.eclipse.microprofile.config.ConfigProvider;
+import org.jboss.logging.Logger;
+
+/**
+ * The {@link CliWebSocketClient} of the Camel CLI connector WebSocket 
transport, with the Quarkus WebSockets Next
+ * client: registered as a bean only when the application has {@code 
quarkus-websockets-next}.
+ * <p/>
+ * The callbacks run on the Vert.x event loop: they only hand over to the 
transport, which parses and sends on its own
+ * threads.
+ */
+public class QuarkusCliWebSocketClient implements CliWebSocketClient {
+
+    static final String NAME = "quarkus-websockets-next";
+    private static final Logger LOG = 
Logger.getLogger(QuarkusCliWebSocketClient.class);
+    private static final String MAX_MESSAGE_SIZE_KEY = 
"quarkus.websockets-next.client.max-message-size";
+    private static final int GOING_AWAY = 1001;
+
+    @Inject
+    CamelCliConnectorRunTimeConfig config;
+
+    @PostConstruct
+    void checkMaxMessageSize() {
+        // set globally, it replaces the size of every connector (see 
customizeOptions)
+        OptionalInt size = 
ConfigProvider.getConfig().getOptionalValue(MAX_MESSAGE_SIZE_KEY, Integer.class)
+                .map(OptionalInt::of).orElse(OptionalInt.empty());

Review Comment:
   Nit: the `OptionalInt` conversion isn't needed. 
`getOptionalValue(MAX_MESSAGE_SIZE_KEY, Integer.class).filter(s -> s < 
MAX_MESSAGE_SIZE).ifPresent(s -> LOG.warnf(...))` does the same.



##########
integration-tests-jvm/cli-connector-websockets-next/pom.xml:
##########
@@ -0,0 +1,184 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+    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.
+
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0"; 
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"; 
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
http://maven.apache.org/xsd/maven-4.0.0.xsd";>
+    <modelVersion>4.0.0</modelVersion>
+    <parent>
+        <groupId>org.apache.camel.quarkus</groupId>
+        <artifactId>camel-quarkus-build-parent-it</artifactId>
+        <version>4.0.0-SNAPSHOT</version>
+        <relativePath>../../poms/build-parent-it/pom.xml</relativePath>
+    </parent>
+
+    
<artifactId>camel-quarkus-integration-test-cli-connector-websockets-next</artifactId>
+    <name>Camel Quarkus :: Integration Tests :: CLI Connector :: WebSockets 
Next</name>
+    <description>Integration tests for Camel Quarkus CLI Connector 
extension</description>
+
+    <dependencies>
+        <dependency>
+            <groupId>org.apache.camel.quarkus</groupId>
+            <artifactId>camel-quarkus-cli-connector</artifactId>
+        </dependency>
+        <dependency>
+            <groupId>io.quarkus</groupId>
+            <artifactId>quarkus-websockets-next</artifactId>
+        </dependency>
+        <dependency>
+            <groupId>org.apache.camel.quarkus</groupId>
+            <artifactId>camel-quarkus-management</artifactId>
+        </dependency>
+        <dependency>
+            <groupId>org.apache.camel.quarkus</groupId>
+            <artifactId>camel-quarkus-direct</artifactId>
+        </dependency>
+        <!-- test dependencies -->
+        <dependency>
+            <groupId>io.quarkus</groupId>
+            <artifactId>quarkus-junit</artifactId>
+            <scope>test</scope>
+        </dependency>
+        <dependency>
+            <groupId>io.rest-assured</groupId>
+            <artifactId>rest-assured</artifactId>
+            <scope>test</scope>
+        </dependency>

Review Comment:
   Nit: `rest-assured` looks unused now that `CliConnectorResource` and the 
RestAssured test are gone. Same in 
`integration-tests-jvm/cli-connector/pom.xml`.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to