This is an automated email from the ASF dual-hosted git repository.
jamesnetherton pushed a commit to branch camel-main
in repository https://gitbox.apache.org/repos/asf/camel-quarkus.git
The following commit(s) were added to refs/heads/camel-main by this push:
new 0478ca48a1 CAMEL-25197: cli-connector - use the Vert.x WebSocket
client for the WebSocket transport
0478ca48a1 is described below
commit 0478ca48a1118c3147dc33b065bdaf311a9b0fa5
Author: Federico Mariani <[email protected]>
AuthorDate: Thu Oct 1 19:44:21 2026 +0200
CAMEL-25197: cli-connector - use the Vert.x WebSocket client for the
WebSocket transport
* CAMEL-25197: cli-connector - use the Quarkus WebSockets Next client for
the WebSocket transport
When the application has quarkus-websockets-next, the Camel CLI connector
WebSocket transport
(camel.cli.transport=websocket) uses its client through the new
CliWebSocketClient SPI of
camel-cli-connector, otherwise the JDK client. No dependency is added to
the application for this:
quarkus-websockets-next is an optional dependency, and the client bean is
only registered when the
capability is present.
- QuarkusCliWebSocketClient: BasicWebSocketConnector, non-blocking
callbacks on the event loop, 16 MB
messages, 64 KB frames, path and query sent as given, handshake
rejections mapped to their HTTP status
- quarkus.camel.cli.websocket.tls-configuration-name (runtime) selects a
TLS registry configuration
- workarounds for WebSockets Next 3.40: back-pressure disabled by default
(quarkus.websockets-next.client.max-pending-messages=0, it stalls
connections receiving fragmented
messages), and no onClose callback (a 1006 close fails in WebSockets Next
and skips the cleanup)
- integration tests for both clients: cli-connector (JDK) and
cli-connector-websockets-next, sharing the
test application and a WebSocket tool server
Co-Authored-By: Claude Opus 5.5 <[email protected]>
* CAMEL-25197: cli-connector - use the Vert.x WebSocket client, refuse the
WebSocket transport in prod
Review follow-up:
- the Vert.x WebSocket client of the application replaces the WebSockets
Next client, and is the default
on Camel Quarkus (camel.cli.websocket.client=jdk uses the JDK client): no
optional dependency, no
application-wide back-pressure setting, close status and raw request uri
as given, a real abort
- connect and handshake timeouts of 10 s, as the JDK client
- whole messages accepted in a single frame (sizes in bytes), messages sent
in frames of 64 KB at most
- the WebSocket transport is refused with the Quarkus prod profile: Camel
Quarkus only sets the Camel
profile in dev mode, so the Camel prod profile check did not apply to
packaged applications
- tests: a 1 MB action in a single frame, a tool that does not answer the
handshake, the prod profile
(QuarkusProdModeTest); unused rest-assured removed
Co-Authored-By: Claude Opus 5.5 <[email protected]>
* CAMEL-25197: cli-connector - match the transport name in any case, apply
the whole TLS configuration
- the Quarkus prod profile guard matches camel.cli.transport in any case,
as Camel does
- the TLS registry configuration is applied with TlsConfigUtils (cipher
suites, protocols, CRLs, ...)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---------
Co-authored-by: Claude Opus 5.5 <[email protected]>
---
.../pages/reference/extensions/cli-connector.adoc | 37 ++++
extensions-jvm/cli-connector/deployment/pom.xml | 4 +
.../deployment/CliConnectorProcessor.java | 15 ++
.../CliConnectorWebSocketProdProfileTest.java | 48 +++++
extensions-jvm/cli-connector/runtime/pom.xml | 4 +
.../cli-connector/runtime/src/main/doc/usage.adoc | 24 +++
.../connector/CamelCliConnectorRunTimeConfig.java | 45 +++++
.../cli/connector/QuarkusLocalCliConnector.java | 15 ++
.../cli/connector/VertxCliWebSocketClient.java | 198 +++++++++++++++++++++
integration-tests-jvm/cli-connector/pom.xml | 28 ++-
.../cli/connector/it/CliConnectorResource.java | 47 -----
.../cli/connector/it/CliConnectorRoutes.java} | 19 +-
.../src/main/resources/application.properties | 25 +++
.../cli/connector/it/CliConnectorTest.java | 139 ++++++++++++++-
.../component/cli/connector/it/CliToolServer.java | 164 +++++++++++++++++
15 files changed, 741 insertions(+), 71 deletions(-)
diff --git a/docs/modules/ROOT/pages/reference/extensions/cli-connector.adoc
b/docs/modules/ROOT/pages/reference/extensions/cli-connector.adoc
index 55c74a4c75..c5152000f6 100644
--- a/docs/modules/ROOT/pages/reference/extensions/cli-connector.adoc
+++ b/docs/modules/ROOT/pages/reference/extensions/cli-connector.adoc
@@ -40,6 +40,36 @@ ifeval::[{doc-show-user-guide-link} == true]
Check the xref:user-guide/index.adoc[User guide] for more information about
writing Camel Quarkus applications.
endif::[]
+[id="extensions-cli-connector-usage"]
+== Usage
+[id="extensions-cli-connector-usage-websocket-transport"]
+=== 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]
+----
+%dev.camel.cli.transport = websocket
+%dev.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 Quarkus `prod` profile (the
+default profile of a packaged application), as with the Camel `prod` profile
(`camel.main.profile=prod`). To use it
+from a packaged application, for example on a remote development cluster, run
it with another profile.
+
+[id="extensions-cli-connector-usage-websocket-client"]
+==== WebSocket client
+
+The connector uses the Vert.x WebSocket client of the application. The client
in use is logged at startup
+(`Camel CLI connector uses the vertx WebSocket client`) and reported to the
tool. Set
+`camel.cli.websocket.client=jdk` to use the JDK WebSocket client instead.
+
+For a `wss://` tool, `quarkus.camel.cli.websocket.tls-configuration-name`
selects a TLS configuration of the
+https://quarkus.io/guides/tls-registry-reference[Quarkus TLS registry];
otherwise the JVM default trust store is used.
+
+
[id="extensions-cli-connector-additional-camel-quarkus-configuration"]
== Additional Camel Quarkus configuration
@@ -53,6 +83,13 @@ a|icon:lock[title=Fixed at build time]
[[quarkus-camel-cli-enabled]]`link:#quark
Sets whether to enable Camel CLI Connector support.
| `boolean`
| `true`
+
+a|
[[quarkus-camel-cli-websocket-tls-configuration-name]]`link:#quarkus-camel-cli-websocket-tls-configuration-name[quarkus.camel.cli.websocket.tls-configuration-name]`
+
+The name of the TLS configuration (from the Quarkus TLS registry) used to
connect to a `wss://` tool. When
+not set, the JVM default trust store is used.
+| `string`
+|
|===
[.configuration-legend]
diff --git a/extensions-jvm/cli-connector/deployment/pom.xml
b/extensions-jvm/cli-connector/deployment/pom.xml
index 4ecc63428c..47cbf556e0 100644
--- a/extensions-jvm/cli-connector/deployment/pom.xml
+++ b/extensions-jvm/cli-connector/deployment/pom.xml
@@ -46,6 +46,10 @@
<groupId>org.apache.camel.quarkus</groupId>
<artifactId>camel-quarkus-cli-connector</artifactId>
</dependency>
+ <dependency>
+ <groupId>io.quarkus</groupId>
+ <artifactId>quarkus-vertx-http-deployment</artifactId>
+ </dependency>
<dependency>
<groupId>io.quarkus</groupId>
diff --git
a/extensions-jvm/cli-connector/deployment/src/main/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorProcessor.java
b/extensions-jvm/cli-connector/deployment/src/main/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorProcessor.java
index 8461d504f3..1564b1ec62 100644
---
a/extensions-jvm/cli-connector/deployment/src/main/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorProcessor.java
+++
b/extensions-jvm/cli-connector/deployment/src/main/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorProcessor.java
@@ -18,6 +18,8 @@ package
org.apache.camel.quarkus.component.cli.connector.deployment;
import java.util.function.BooleanSupplier;
+import io.quarkus.arc.deployment.AdditionalBeanBuildItem;
+import io.quarkus.arc.processor.DotNames;
import io.quarkus.builder.Version;
import io.quarkus.deployment.annotations.BuildStep;
import io.quarkus.deployment.annotations.BuildSteps;
@@ -27,6 +29,7 @@ import io.quarkus.deployment.builditem.FeatureBuildItem;
import io.quarkus.deployment.pkg.steps.NativeOrNativeSourcesBuild;
import
org.apache.camel.quarkus.component.cli.connector.CamelCliConnectorConfig;
import
org.apache.camel.quarkus.component.cli.connector.CamelCliConnectorRecorder;
+import
org.apache.camel.quarkus.component.cli.connector.VertxCliWebSocketClient;
import org.apache.camel.quarkus.core.JvmOnlyRecorder;
import org.apache.camel.quarkus.core.deployment.spi.CamelBeanBuildItem;
import org.apache.camel.spi.CliConnectorFactory;
@@ -51,6 +54,18 @@ class CliConnectorProcessor {
recorder.createCliConnectorFactory(Version.getVersion()));
}
+ /**
+ * The WebSocket transport uses the Vert.x client, unless
camel.cli.websocket.client=jdk.
+ */
+ @BuildStep
+ AdditionalBeanBuildItem webSocketClient() {
+ return AdditionalBeanBuildItem.builder()
+ .addBeanClasses(VertxCliWebSocketClient.class)
+ .setDefaultScope(DotNames.SINGLETON)
+ .setUnremovable()
+ .build();
+ }
+
/**
* Remove this once this extension starts supporting the native mode.
*/
diff --git
a/extensions-jvm/cli-connector/deployment/src/test/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorWebSocketProdProfileTest.java
b/extensions-jvm/cli-connector/deployment/src/test/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorWebSocketProdProfileTest.java
new file mode 100644
index 0000000000..f6b871d2ea
--- /dev/null
+++
b/extensions-jvm/cli-connector/deployment/src/test/java/org/apache/camel/quarkus/component/cli/connector/deployment/CliConnectorWebSocketProdProfileTest.java
@@ -0,0 +1,48 @@
+/*
+ * 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.deployment;
+
+import java.util.Map;
+
+import io.quarkus.test.QuarkusProdModeTest;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class CliConnectorWebSocketProdProfileTest {
+ // a packaged application: the Quarkus prod profile
+ @RegisterExtension
+ static final QuarkusProdModeTest CONFIG = new QuarkusProdModeTest()
+ .withEmptyApplication()
+ .setApplicationName("cli-connector-prod")
+ .setApplicationVersion("1.0")
+ .setRun(true)
+ .setExpectExit(true)
+ .setRuntimeProperties(Map.of(
+ // any case, as Camel matches it
+ "camel.cli.transport", "WebSocket",
+ "camel.cli.websocket.url", "ws://127.0.0.1:9/connect"));
+
+ @Test
+ void refusesTheWebSocketTransport() {
+ String output = CONFIG.getStartupConsoleOutput();
+ assertTrue(output.contains("cannot be used with the Quarkus prod
profile"), output);
+ assertNotEquals(0, CONFIG.getExitCode());
+ }
+}
diff --git a/extensions-jvm/cli-connector/runtime/pom.xml
b/extensions-jvm/cli-connector/runtime/pom.xml
index 3e2a535579..084176c02e 100644
--- a/extensions-jvm/cli-connector/runtime/pom.xml
+++ b/extensions-jvm/cli-connector/runtime/pom.xml
@@ -51,6 +51,10 @@
<groupId>org.apache.camel</groupId>
<artifactId>camel-cli-connector</artifactId>
</dependency>
+ <dependency>
+ <groupId>io.quarkus</groupId>
+ <artifactId>quarkus-vertx-http</artifactId>
+ </dependency>
</dependencies>
<build>
diff --git a/extensions-jvm/cli-connector/runtime/src/main/doc/usage.adoc
b/extensions-jvm/cli-connector/runtime/src/main/doc/usage.adoc
new file mode 100644
index 0000000000..2a3f07e2d6
--- /dev/null
+++ b/extensions-jvm/cli-connector/runtime/src/main/doc/usage.adoc
@@ -0,0 +1,24 @@
+=== 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]
+----
+%dev.camel.cli.transport = websocket
+%dev.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 Quarkus `prod` profile (the
+default profile of a packaged application), as with the Camel `prod` profile
(`camel.main.profile=prod`). To use it
+from a packaged application, for example on a remote development cluster, run
it with another profile.
+
+==== WebSocket client
+
+The connector uses the Vert.x WebSocket client of the application. The client
in use is logged at startup
+(`Camel CLI connector uses the vertx WebSocket client`) and reported to the
tool. Set
+`camel.cli.websocket.client=jdk` to use the JDK WebSocket client instead.
+
+For a `wss://` tool, `quarkus.camel.cli.websocket.tls-configuration-name`
selects a TLS configuration of the
+https://quarkus.io/guides/tls-registry-reference[Quarkus TLS registry];
otherwise the JVM default trust store is used.
diff --git
a/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/CamelCliConnectorRunTimeConfig.java
b/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/CamelCliConnectorRunTimeConfig.java
new file mode 100644
index 0000000000..0de5cef305
--- /dev/null
+++
b/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/CamelCliConnectorRunTimeConfig.java
@@ -0,0 +1,45 @@
+/*
+ * 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.util.Optional;
+
+import io.quarkus.runtime.annotations.ConfigPhase;
+import io.quarkus.runtime.annotations.ConfigRoot;
+import io.smallrye.config.ConfigMapping;
+
+@ConfigRoot(phase = ConfigPhase.RUN_TIME)
+@ConfigMapping(prefix = "quarkus.camel.cli")
+public interface CamelCliConnectorRunTimeConfig {
+ /**
+ * Options of the WebSocket transport (`camel.cli.transport=websocket`)
when the Vert.x client is used (the
+ * default).
+ *
+ * @asciidoclet
+ */
+ WebSocketConfig websocket();
+
+ interface WebSocketConfig {
+ /**
+ * The name of the TLS configuration (from the Quarkus TLS registry)
used to connect to a `wss://` tool. When
+ * not set, the JVM default trust store is used.
+ *
+ * @asciidoclet
+ */
+ Optional<String> tlsConfigurationName();
+ }
+}
diff --git
a/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusLocalCliConnector.java
b/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusLocalCliConnector.java
index d8fb9f4e22..f3db169b22 100644
---
a/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusLocalCliConnector.java
+++
b/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/QuarkusLocalCliConnector.java
@@ -18,6 +18,8 @@ package org.apache.camel.quarkus.component.cli.connector;
import io.quarkus.runtime.LaunchMode;
import io.quarkus.runtime.Quarkus;
+import io.quarkus.runtime.configuration.ConfigUtils;
+import org.apache.camel.cli.connector.CliConnectorTransport;
import org.apache.camel.cli.connector.LocalCliConnector;
import org.apache.camel.spi.CliConnectorFactory;
@@ -26,6 +28,19 @@ public class QuarkusLocalCliConnector extends
LocalCliConnector {
super(cliConnectorFactory);
}
+ @Override
+ protected CliConnectorTransport createTransport(String name) {
+ // Camel refuses it with the Camel prod profile, which Camel Quarkus
does not set outside dev mode
+ // matched as Camel does (any case)
+ if ("websocket".equalsIgnoreCase(name) &&
ConfigUtils.isProfileActive("prod")) {
+ throw new IllegalStateException(
+ "The Camel CLI connector websocket transport gives the
connected tool full control of this"
+ + " application and cannot be used with the
Quarkus prod profile."
+ + " Remove camel.cli.transport=websocket, or use
another profile.");
+ }
+ return super.createTransport(name);
+ }
+
@Override
public void sigterm() {
if (LaunchMode.current().equals(LaunchMode.DEVELOPMENT)) {
diff --git
a/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/VertxCliWebSocketClient.java
b/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/VertxCliWebSocketClient.java
new file mode 100644
index 0000000000..be63336cec
--- /dev/null
+++
b/extensions-jvm/cli-connector/runtime/src/main/java/org/apache/camel/quarkus/component/cli/connector/VertxCliWebSocketClient.java
@@ -0,0 +1,198 @@
+/*
+ * 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.Locale;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionStage;
+import java.util.concurrent.TimeoutException;
+
+import io.quarkus.tls.TlsConfiguration;
+import io.quarkus.tls.TlsConfigurationRegistry;
+import io.quarkus.tls.runtime.config.TlsConfigUtils;
+import io.vertx.core.Future;
+import io.vertx.core.Vertx;
+import io.vertx.core.buffer.Buffer;
+import io.vertx.core.http.UpgradeRejectedException;
+import io.vertx.core.http.WebSocket;
+import io.vertx.core.http.WebSocketClient;
+import io.vertx.core.http.WebSocketClientOptions;
+import io.vertx.core.http.WebSocketConnectOptions;
+import io.vertx.core.http.WebSocketFrame;
+import jakarta.annotation.PostConstruct;
+import jakarta.inject.Inject;
+import org.apache.camel.cli.connector.CliWebSocketClient;
+import org.apache.camel.cli.connector.CliWebSocketHandshakeException;
+import org.jboss.logging.Logger;
+
+/**
+ * The {@link CliWebSocketClient} of the Camel CLI connector WebSocket
transport on Camel Quarkus, with the Vert.x
+ * WebSocket client of the application. {@code camel.cli.websocket.client=jdk}
uses the JDK client instead.
+ * <p/>
+ * The handlers run on the Vert.x event loop: they only hand over to the
transport, which parses and sends on its own
+ * threads.
+ */
+public class VertxCliWebSocketClient implements CliWebSocketClient {
+
+ static final String NAME = "vertx";
+ private static final Logger LOG =
Logger.getLogger(VertxCliWebSocketClient.class);
+ private static final int TIMEOUT = 10000;
+ // the frames sent: servers commonly refuse frames over 64 KB (Vert.x,
Quarkus), and a char takes up to 3 bytes
+ private static final int FRAME_CHARS = 16 * 1024;
+ // the messages received, in bytes: up to 3 per char
+ private static final int MAX_MESSAGE_BYTES = 3 * MAX_MESSAGE_SIZE;
+
+ @Inject
+ Vertx vertx;
+
+ @Inject
+ TlsConfigurationRegistry tlsRegistry;
+
+ @Inject
+ CamelCliConnectorRunTimeConfig config;
+
+ private TlsConfiguration tls;
+
+ @PostConstruct
+ void resolveTlsConfiguration() {
+ Optional<String> name = config.websocket().tlsConfigurationName();
+ if (name.isPresent()) {
+ tls = tlsRegistry.get(name.get()).orElseThrow(() -> new
IllegalStateException(
+ "No TLS configuration named " + name.get() + "
(quarkus.camel.cli.websocket.tls-configuration-name)"));
+ }
+ }
+
+ @Override
+ public String getName() {
+ return NAME;
+ }
+
+ @Override
+ public CompletionStage<Channel> connect(URI url, Map<String, String>
headers, Listener listener) {
+ boolean ssl = "wss".equals(url.getScheme().toLowerCase(Locale.ROOT));
+ WebSocketClientOptions options = new WebSocketClientOptions()
+ .setConnectTimeout(TIMEOUT)
+ // how long the tool has to close the connection once the
close frame is sent (seconds): abort() closes
+ // the client, which closes the connection like this
+ .setClosingTimeout(1)
+ // a tool can send a whole message in a single frame
+ .setMaxFrameSize(MAX_MESSAGE_BYTES)
+ .setMaxMessageSize(MAX_MESSAGE_BYTES);
+ if (ssl && tls != null) {
+ TlsConfigUtils.configure(options, tls);
+ }
+ String path = url.getRawPath() == null || url.getRawPath().isEmpty() ?
"/" : url.getRawPath();
+ WebSocketConnectOptions connect = new WebSocketConnectOptions()
+ .setHost(url.getHost())
+ .setPort(url.getPort() != -1 ? url.getPort() : ssl ? 443 : 80)
+ .setSsl(ssl)
+ // the path and query as given, encoded characters included
+ .setURI(url.getRawQuery() != null ? path + "?" +
url.getRawQuery() : path);
+ headers.forEach(connect::addHeader);
+
+ // a client for each connection, closed with it
+ WebSocketClient client = vertx.createWebSocketClient(options);
+ CompletableFuture<Channel> answer = new CompletableFuture<>();
+ // the handshake: not WebSocketConnectOptions.setTimeout, which is
also an idle timeout that can close the
+ // connection once open
+ long timer = vertx.setTimer(TIMEOUT, id -> {
+ if (answer.completeExceptionally(new TimeoutException("WebSocket
handshake timed out after " + TIMEOUT + " ms"))) {
+ client.close();
+ }
+ });
+ client.connect(connect).onComplete(ar -> {
+ vertx.cancelTimer(timer);
+ if (answer.isDone()) {
+ // timed out
+ if (ar.succeeded()) {
+ ar.result().close();
+ }
+ return;
+ }
+ if (ar.failed()) {
+ client.close();
+ answer.completeExceptionally(translate(ar.cause()));
+ return;
+ }
+ WebSocket ws = ar.result();
+ ws.textMessageHandler(listener::onText);
+ ws.pongHandler(data -> listener.onPong());
+ ws.exceptionHandler(e -> {
+ LOG.debugf(e, "Camel CLI connector WebSocket error");
+ listener.onError(e);
+ });
+ ws.closeHandler(v -> {
+ client.close();
+ // no status when the connection is lost without a close frame
+ Short code = ws.closeStatusCode();
+ LOG.debugf("Camel CLI connector WebSocket closed: %s %s",
code, ws.closeReason());
+ listener.onClose(code != null ? code : 1006, ws.closeReason());
+ });
+ answer.complete(new VertxChannel(ws, client));
+ });
+ return answer;
+ }
+
+ private static Throwable translate(Throwable e) {
+ if (e instanceof UpgradeRejectedException rejected) {
+ return new CliWebSocketHandshakeException(rejected.getStatus(),
rejected);
+ }
+ return e;
+ }
+
+ private record VertxChannel(WebSocket ws, WebSocketClient client)
implements Channel {
+
+ @Override
+ public CompletionStage<?> sendText(String text) {
+ // in frames of FRAME_CHARS, never splitting a surrogate pair
+ Future<Void> last = null;
+ int start = 0;
+ do {
+ int end = Math.min(text.length(), start + FRAME_CHARS);
+ if (end < text.length() &&
Character.isHighSurrogate(text.charAt(end - 1))) {
+ end--;
+ }
+ String part = text.substring(start, end);
+ boolean fin = end == text.length();
+ last = ws.writeFrame(start == 0
+ ? WebSocketFrame.textFrame(part, fin)
+ :
WebSocketFrame.continuationFrame(Buffer.buffer(part), fin));
+ start = end;
+ } while (start < text.length());
+ return last.toCompletionStage();
+ }
+
+ @Override
+ public CompletionStage<?> sendPing() {
+ return ws.writePing(Buffer.buffer()).toCompletionStage();
+ }
+
+ @Override
+ public CompletionStage<?> close(int code, String reason) {
+ return ws.close((short) code, reason).toCompletionStage();
+ }
+
+ @Override
+ public void abort() {
+ // closes the connection right away, without a close frame
+ client.close();
+ }
+ }
+}
diff --git a/integration-tests-jvm/cli-connector/pom.xml
b/integration-tests-jvm/cli-connector/pom.xml
index b9e9b1e314..4f01f6ca89 100644
--- a/integration-tests-jvm/cli-connector/pom.xml
+++ b/integration-tests-jvm/cli-connector/pom.xml
@@ -36,12 +36,12 @@
<artifactId>camel-quarkus-cli-connector</artifactId>
</dependency>
<dependency>
- <groupId>io.quarkus</groupId>
- <artifactId>quarkus-resteasy</artifactId>
+ <groupId>org.apache.camel.quarkus</groupId>
+ <artifactId>camel-quarkus-management</artifactId>
</dependency>
<dependency>
<groupId>org.apache.camel.quarkus</groupId>
- <artifactId>camel-quarkus-management</artifactId>
+ <artifactId>camel-quarkus-direct</artifactId>
</dependency>
<!-- test dependencies -->
<dependency>
@@ -50,8 +50,13 @@
<scope>test</scope>
</dependency>
<dependency>
- <groupId>io.rest-assured</groupId>
- <artifactId>rest-assured</artifactId>
+ <groupId>org.awaitility</groupId>
+ <artifactId>awaitility</artifactId>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.assertj</groupId>
+ <artifactId>assertj-core</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
@@ -79,6 +84,19 @@
</exclusion>
</exclusions>
</dependency>
+ <dependency>
+ <groupId>org.apache.camel.quarkus</groupId>
+ <artifactId>camel-quarkus-direct-deployment</artifactId>
+ <version>${project.version}</version>
+ <type>pom</type>
+ <scope>test</scope>
+ <exclusions>
+ <exclusion>
+ <groupId>*</groupId>
+ <artifactId>*</artifactId>
+ </exclusion>
+ </exclusions>
+ </dependency>
<dependency>
<groupId>org.apache.camel.quarkus</groupId>
<artifactId>camel-quarkus-management-deployment</artifactId>
diff --git
a/integration-tests-jvm/cli-connector/src/main/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorResource.java
b/integration-tests-jvm/cli-connector/src/main/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorResource.java
deleted file mode 100644
index c6cdb58d7c..0000000000
---
a/integration-tests-jvm/cli-connector/src/main/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorResource.java
+++ /dev/null
@@ -1,47 +0,0 @@
-/*
- * 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.it;
-
-import jakarta.enterprise.context.ApplicationScoped;
-import jakarta.inject.Inject;
-import jakarta.ws.rs.GET;
-import jakarta.ws.rs.Path;
-import jakarta.ws.rs.Produces;
-import jakarta.ws.rs.core.MediaType;
-import jakarta.ws.rs.core.Response;
-import org.apache.camel.CamelContext;
-import org.jboss.logging.Logger;
-
-@Path("/cli-connector")
-@ApplicationScoped
-public class CliConnectorResource {
-
- private static final Logger LOG =
Logger.getLogger(CliConnectorResource.class);
-
- private static final String OTHER_CLI_CONNECTOR = "cli-connector";
- @Inject
- CamelContext context;
-
- @Path("/load/other/cli-connector")
- @GET
- @Produces(MediaType.TEXT_PLAIN)
- public Response loadOtherCliConnector() throws Exception {
- /* This is an autogenerated test */
- /* No way to test a Camel artifact of kind "other" */
- return Response.ok().build();
- }
-}
diff --git
a/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorTest.java
b/integration-tests-jvm/cli-connector/src/main/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorRoutes.java
similarity index 70%
copy from
integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorTest.java
copy to
integration-tests-jvm/cli-connector/src/main/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorRoutes.java
index b7b3ac7fe3..aff5d26c0f 100644
---
a/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorTest.java
+++
b/integration-tests-jvm/cli-connector/src/main/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorRoutes.java
@@ -16,19 +16,12 @@
*/
package org.apache.camel.quarkus.component.cli.connector.it;
-import io.quarkus.test.junit.QuarkusTest;
-import io.restassured.RestAssured;
-import org.junit.jupiter.api.Test;
+import org.apache.camel.builder.RouteBuilder;
-@QuarkusTest
-class CliConnectorTest {
-
- @Test
- public void loadOtherCliConnector() {
- /* A simple autogenerated test */
- RestAssured.get("/cli-connector/load/other/cli-connector")
- .then()
- .statusCode(200);
+public class CliConnectorRoutes extends RouteBuilder {
+ @Override
+ public void configure() {
+ from("direct:echo").routeId("echo").setBody(simple("Hello ${body}"));
+
from("direct:length").routeId("length").setBody(simple("length=${body.length()}"));
}
-
}
diff --git
a/integration-tests-jvm/cli-connector/src/main/resources/application.properties
b/integration-tests-jvm/cli-connector/src/main/resources/application.properties
new file mode 100644
index 0000000000..df15566cdd
--- /dev/null
+++
b/integration-tests-jvm/cli-connector/src/main/resources/application.properties
@@ -0,0 +1,25 @@
+## ---------------------------------------------------------------------------
+## 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.
+## ---------------------------------------------------------------------------
+# camel.cli.websocket.url and camel.cli.websocket.token are set by
CliToolServer
+camel.cli.transport = websocket
+camel.cli.websocket.reconnect-delay = 100
+camel.cli.websocket.reconnect-max-delay = 500
+# trace snapshots, sent in messages of up to 128 KB
+camel.trace.enabled = true
+# the tests check what the connector logs
+quarkus.log.file.enabled = true
+quarkus.log.file.path = target/quarkus.log
diff --git
a/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorTest.java
b/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorTest.java
index b7b3ac7fe3..1e056bc990 100644
---
a/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorTest.java
+++
b/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliConnectorTest.java
@@ -16,19 +16,146 @@
*/
package org.apache.camel.quarkus.component.cli.connector.it;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.util.concurrent.TimeUnit;
+
+import io.quarkus.test.common.QuarkusTestResource;
import io.quarkus.test.junit.QuarkusTest;
-import io.restassured.RestAssured;
+import io.vertx.core.http.ServerWebSocket;
+import org.apache.camel.util.json.JsonObject;
import org.junit.jupiter.api.Test;
+import static
org.apache.camel.quarkus.component.cli.connector.it.CliToolServer.action;
+import static
org.apache.camel.quarkus.component.cli.connector.it.CliToolServer.awaitResult;
+import static
org.apache.camel.quarkus.component.cli.connector.it.CliToolServer.reconnect;
+import static
org.apache.camel.quarkus.component.cli.connector.it.CliToolServer.send;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * The Vert.x client, the default on Camel Quarkus.
+ */
@QuarkusTest
+@QuarkusTestResource(CliToolServer.class)
class CliConnectorTest {
+ private static final Path LOG = Paths.get("target/quarkus.log");
+
+ @Test
+ void connectsWithTheExpectedClient() throws Exception {
+ JsonObject hello = reconnect();
+
+ assertThat(hello.getString("transport")).isEqualTo("vertx");
+ assertThat(hello.getString("camelVersion")).isNotBlank();
+ // the url is used as given, encoded characters included
+ assertThat(CliToolServer.REQUEST_URIS).isNotEmpty().allSatisfy(uri ->
assertThat(uri)
+ .isEqualTo(CliToolServer.REQUEST_URI));
+ }
+
+ @Test
+ void executesActions() throws Exception {
+ reconnect();
+
+ send(action("r1", "send", "endpoint", "direct:echo", "body", "World",
"exchangePattern", "InOut"));
+ JsonObject result = awaitResult("r1");
+ assertThat(result.getBoolean("ok")).isTrue();
+ assertThat(resultJson(result)).contains("Hello World");
+
+ send(action("r2", "does-not-exist"));
+ result = awaitResult("r2");
+ assertThat(result.getBoolean("ok")).isFalse();
+ assertThat(result.getString("error")).isEqualTo("Unknown action:
does-not-exist");
+ }
+
+ @Test
+ void exchangesLargeMessages() throws Exception {
+ reconnect();
+ // well over the 256 KB messages and the 64 KB frames that Vert.x
accepts by default
+ String body = "x".repeat(4 * 1024 * 1024);
+
+ // each of them in 64 frames: more frames in total than a client
fetching a fixed number of frames would read
+ for (int i = 0; i < 6; i++) {
+ send(action("r" + i, "send", "endpoint", "direct:length", "body",
body, "exchangePattern", "InOut"));
+ JsonObject result = awaitResult("r" + i);
+ assertThat(result.getBoolean("ok")).isTrue();
+ assertThat(resultJson(result)).contains("length=" + body.length());
+ }
+
+ // and back: trace snapshots are sent in messages of up to 128 KB,
over the 64 KB frames the tool accepts
+ String traced = "y".repeat(6000);
+ for (int i = 0; i < 40; i++) {
+ send(action("t" + i, "send", "endpoint", "direct:echo", "body",
traced, "exchangePattern", "InOut"));
+ assertThat(awaitResult("t" + i).getBoolean("ok")).isTrue();
+ }
+ await().atMost(20, TimeUnit.SECONDS)
+ .untilAsserted(() ->
assertThat(CliToolServer.LARGEST_TRACE_SNAPSHOT).hasValueGreaterThan(64 *
1024));
+ }
+
+ @Test
+ void reconnectsOnceTheTokenIsAccepted() throws Exception {
+ reconnect();
+ int rejected = CliToolServer.REJECTED.get();
+ try {
+ CliToolServer.reject = true;
+ CliToolServer.SOCKETS.forEach(ServerWebSocket::close);
+
+ await().atMost(10, TimeUnit.SECONDS)
+ .untilAsserted(() ->
assertThat(CliToolServer.REJECTED).hasValueGreaterThan(rejected + 1));
+ // the client reports the HTTP status of the rejected handshake
+ await().atMost(10, TimeUnit.SECONDS).untilAsserted(() ->
assertThat(Files.readString(LOG))
+ .contains("Camel CLI connector was rejected by
ws://127.0.0.1:")
+ .contains("(HTTP 401): check camel.cli.websocket.token"));
+ } finally {
+ CliToolServer.reject = false;
+ }
+ // connected again once the token is accepted
+ assertThat(CliToolServer.awaitFrame(f ->
"hello".equals(f.getString("type")))).isNotNull();
+ }
+
+ @Test
+ void receivesLargeActionsInASingleFrame() throws Exception {
+ reconnect();
+ // over the 64 KB frames that Vert.x accepts by default
+ String body = "x".repeat(1024 * 1024);
+
+ CliToolServer.sendInASingleFrame(
+ action("r1", "send", "endpoint", "direct:length", "body",
body, "exchangePattern", "InOut"));
+ JsonObject result = awaitResult("r1");
+ assertThat(result.getBoolean("ok")).isTrue();
+ assertThat(resultJson(result)).contains("length=" + body.length());
+ }
+
+ @Test
+ void staysConnected() throws Exception {
+ reconnect();
+ int connections = CliToolServer.CONNECTIONS.get();
+
+ // longer than the connect and handshake timeouts (10 s), which must
not apply once connected
+ await().during(15, TimeUnit.SECONDS).atMost(20, TimeUnit.SECONDS)
+ .untilAsserted(() ->
assertThat(CliToolServer.CONNECTIONS).hasValue(connections));
+ }
+
@Test
- public void loadOtherCliConnector() {
- /* A simple autogenerated test */
- RestAssured.get("/cli-connector/load/other/cli-connector")
- .then()
- .statusCode(200);
+ void reconnectsWhenTheToolDoesNotAnswerTheHandshake() throws Exception {
+ reconnect();
+ int unanswered = CliToolServer.UNANSWERED.get();
+ try {
+ CliToolServer.silent = true;
+ CliToolServer.SOCKETS.forEach(ServerWebSocket::close);
+
+ // the handshake times out after 10 s, and the connector tries
again
+ await().atMost(40, TimeUnit.SECONDS)
+ .untilAsserted(() ->
assertThat(CliToolServer.UNANSWERED).hasValueGreaterThan(unanswered + 1));
+ } finally {
+ CliToolServer.silent = false;
+ }
+ assertThat(CliToolServer.awaitFrame(f ->
"hello".equals(f.getString("type")))).isNotNull();
}
+ private static String resultJson(JsonObject result) {
+ JsonObject json = result.getMap("result");
+ return json.toJson();
+ }
}
diff --git
a/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliToolServer.java
b/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliToolServer.java
new file mode 100644
index 0000000000..80d2705272
--- /dev/null
+++
b/integration-tests-jvm/cli-connector/src/test/java/org/apache/camel/quarkus/component/cli/connector/it/CliToolServer.java
@@ -0,0 +1,164 @@
+/*
+ * 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.it;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Predicate;
+
+import io.quarkus.test.common.QuarkusTestResourceLifecycleManager;
+import io.vertx.core.Vertx;
+import io.vertx.core.http.HttpServer;
+import io.vertx.core.http.HttpServerOptions;
+import io.vertx.core.http.ServerWebSocket;
+import io.vertx.core.http.WebSocketFrame;
+import org.apache.camel.util.json.JsonObject;
+import org.apache.camel.util.json.Jsoner;
+
+import static org.awaitility.Awaitility.await;
+
+/**
+ * Plays the developer tool: a WebSocket server the Camel CLI connector of the
application dials out to.
+ */
+public class CliToolServer implements QuarkusTestResourceLifecycleManager {
+
+ static final String TOKEN = "it-s3cret";
+ // a query parameter holding an encoded '&': it must reach the tool as
given
+ static final String REQUEST_URI = "/v1/connect?executionId=it%261";
+ private static final int MAX_MESSAGE_SIZE = 32 * 1024 * 1024;
+
+ static final BlockingQueue<JsonObject> FRAMES = new
LinkedBlockingQueue<>();
+ static final List<ServerWebSocket> SOCKETS = new CopyOnWriteArrayList<>();
+ static final List<String> REQUEST_URIS = new CopyOnWriteArrayList<>();
+ static final AtomicInteger REJECTED = new AtomicInteger();
+ static final AtomicInteger CONNECTIONS = new AtomicInteger();
+ static final AtomicInteger UNANSWERED = new AtomicInteger();
+ static volatile boolean silent;
+ static final AtomicInteger LARGEST_TRACE_SNAPSHOT = new AtomicInteger();
+ static volatile boolean reject;
+
+ private Vertx vertx;
+
+ @Override
+ public Map<String, String> start() {
+ vertx = Vertx.vertx();
+ try {
+ // frames of 64 KB at most (the Vert.x default), in both
directions, as Vert.x and Quarkus servers
+ HttpServer server = vertx.createHttpServer(new HttpServerOptions()
+ .setMaxWebSocketMessageSize(MAX_MESSAGE_SIZE))
+ .webSocketHandshakeHandler(handshake -> {
+ REQUEST_URIS.add(handshake.uri());
+ if (silent) {
+ // never answers the upgrade
+ UNANSWERED.incrementAndGet();
+ return;
+ }
+ if (reject || !("Bearer " +
TOKEN).equals(handshake.headers().get("Authorization"))) {
+ REJECTED.incrementAndGet();
+ handshake.reject(401);
+ } else {
+ handshake.accept();
+ }
+ })
+ .webSocketHandler(ws -> {
+ SOCKETS.add(ws);
+ CONNECTIONS.incrementAndGet();
+ ws.textMessageHandler(text -> {
+ try {
+ JsonObject frame = (JsonObject)
Jsoner.deserialize(text);
+ if ("trace".equals(frame.getString("kind"))) {
+
LARGEST_TRACE_SNAPSHOT.accumulateAndGet(text.length(), Math::max);
+ }
+ FRAMES.add(frame);
+ } catch (Exception e) {
+ throw new IllegalStateException(e);
+ }
+ });
+ ws.closeHandler(v -> SOCKETS.remove(ws));
+ })
+ .listen(0,
"127.0.0.1").toCompletionStage().toCompletableFuture().get(10,
TimeUnit.SECONDS);
+ return Map.of(
+ "camel.cli.websocket.url", "ws://127.0.0.1:" +
server.actualPort() + REQUEST_URI,
+ "camel.cli.websocket.token", TOKEN);
+ } catch (Exception e) {
+ throw new IllegalStateException(e);
+ }
+ }
+
+ @Override
+ public void stop() {
+ if (vertx != null) {
+ vertx.close();
+ }
+ }
+
+ /**
+ * Closes the connection, and waits for the application to connect again
and say hello. The old sockets are not
+ * waited for: they are closed once the application answers the close, or
after the Vert.x closing timeout.
+ */
+ static JsonObject reconnect() throws InterruptedException {
+ int connections = CONNECTIONS.get();
+ FRAMES.clear();
+ SOCKETS.forEach(ServerWebSocket::close);
+ await().atMost(20, TimeUnit.SECONDS).until(() -> CONNECTIONS.get() >
connections);
+ return awaitFrame(f -> "hello".equals(f.getString("type")));
+ }
+
+ static void send(JsonObject frame) {
+ await().atMost(10, TimeUnit.SECONDS).until(() -> !SOCKETS.isEmpty());
+ // in frames of 64 KB at most, as Vert.x and Quarkus servers
+ SOCKETS.get(SOCKETS.size() - 1).writeTextMessage(frame.toJson());
+ }
+
+ static void sendInASingleFrame(JsonObject frame) {
+ await().atMost(10, TimeUnit.SECONDS).until(() -> !SOCKETS.isEmpty());
+ // as tools that never split messages, whatever the size
+ SOCKETS.get(SOCKETS.size() -
1).writeFrame(WebSocketFrame.textFrame(frame.toJson(), true));
+ }
+
+ static JsonObject action(String requestId, String name, String...
keyValues) {
+ JsonObject action = new JsonObject();
+ action.put("action", name);
+ for (int i = 0; i < keyValues.length; i += 2) {
+ action.put(keyValues[i], keyValues[i + 1]);
+ }
+ return new JsonObject(Map.of("v", 1, "type", "action", "requestId",
requestId, "action", action));
+ }
+
+ /**
+ * Waits for a frame matching the predicate, skipping the others (mostly
snapshots).
+ */
+ static JsonObject awaitFrame(Predicate<JsonObject> predicate) throws
InterruptedException {
+ long deadline = System.currentTimeMillis() + 20000;
+ while (System.currentTimeMillis() < deadline) {
+ JsonObject frame = FRAMES.poll(100, TimeUnit.MILLISECONDS);
+ if (frame != null && predicate.test(frame)) {
+ return frame;
+ }
+ }
+ throw new AssertionError("No matching frame received within 20
seconds");
+ }
+
+ static JsonObject awaitResult(String requestId) throws
InterruptedException {
+ return awaitFrame(f -> "result".equals(f.getString("type")) &&
requestId.equals(f.getString("requestId")));
+ }
+}