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-spring-boot.git
commit 98e0ceed23f2c6a254ec596a26464b02e1527012 Author: croway <[email protected]> AuthorDate: Thu Oct 1 15:08:58 2026 +0200 CAMEL-25197: camel-cli-connector-starter - use the Spring WebSocket client for the websocket transport When the application has spring-websocket and a Jakarta WebSocket implementation (spring-boot-starter-websocket), the starter registers a CliWebSocketClient backed by Spring's StandardWebSocketClient, which the websocket transport uses instead of the JDK client. Without them nothing is registered and Camel uses the JDK client; spring-websocket is an optional dependency, nothing is added to the application. - sends through the asynchronous Jakarta remote endpoint, so the transport's send timeout can abort a hanging send - reassembles partial messages itself (16 MB cap, the connection is aborted above it) instead of raising the Jakarta text buffer, which Tomcat allocates up front for every connection - reports the HTTP status of a rejected handshake (Tomcat puts it in the message) - camel.cli.websocket.ssl-bundle selects an SSL bundle for wss:// - camel.cli.transport and camel.cli.websocket.* in the configuration metadata, and a starter documentation page The stop action (and camel stop) now closes the Spring application context instead of only stopping it: a stopped web application kept the JVM running. Co-Authored-By: Claude Opus 5.5 <[email protected]> --- docs/spring-boot/modules/ROOT/nav.adoc | 1 + .../modules/ROOT/pages/starters/cli-connector.adoc | 71 ++++++ dsl-starter/camel-cli-connector-starter/pom.xml | 48 ++++ .../src/main/doc/intro.adoc | 3 + .../src/main/doc/usage.adoc | 30 +++ .../src/main/docs/cli-connector.json | 78 +++++- .../connector/CliConnectorAutoConfiguration.java | 30 +++ .../cli/connector/CliConnectorConfiguration.java | 147 ++++++++++++ .../cli/connector/SpringCliWebSocketClient.java | 263 +++++++++++++++++++++ .../cli/connector/SpringLocalCliConnector.java | 4 +- .../CliConnectorAutoConfigurationTest.java | 84 +++++++ .../cli/connector/CliConnectorStopActionTest.java | 68 ++++++ .../CliConnectorWebSocketJdkClientTest.java | 45 ++++ .../CliConnectorWebSocketSpringClientTest.java | 44 ++++ .../CliConnectorWebSocketTestSupport.java | 91 +++++++ .../connector/SpringCliWebSocketClientTest.java | 148 ++++++++++++ .../camel/springboot/cli/connector/ToolServer.java | 180 ++++++++++++++ 17 files changed, 1333 insertions(+), 2 deletions(-) diff --git a/docs/spring-boot/modules/ROOT/nav.adoc b/docs/spring-boot/modules/ROOT/nav.adoc index a54292f14bd..38ad861c5a7 100644 --- a/docs/spring-boot/modules/ROOT/nav.adoc +++ b/docs/spring-boot/modules/ROOT/nav.adoc @@ -87,6 +87,7 @@ ** xref:starters/cbor.adoc[CBOR] ** xref:starters/chatscript.adoc[ChatScript] ** xref:starters/chunk.adoc[Chunk] +** xref:starters/cli-connector.adoc[Cli Connector] ** xref:starters/clickhouse.adoc[ClickHouse] ** xref:starters/clickup.adoc[ClickUp] ** xref:starters/cloudevents.adoc[Cloudevents] diff --git a/docs/spring-boot/modules/ROOT/pages/starters/cli-connector.adoc b/docs/spring-boot/modules/ROOT/pages/starters/cli-connector.adoc new file mode 100644 index 00000000000..8214a543820 --- /dev/null +++ b/docs/spring-boot/modules/ROOT/pages/starters/cli-connector.adoc @@ -0,0 +1,71 @@ +// Do not edit directly! +// This file was generated by camel-spring-boot-generator-maven-plugin += Cli Connector +:artifactid: camel-cli-connector-starter + +Spring Boot auto-configuration for the Camel CLI connector. + +The connector lets developer tools manage a running Spring Boot application: the Camel CLI (`camel get`, `camel trace`, `camel stop`, ...) through files in `~/.camel`, or a tool such as an IDE over a WebSocket the application dials out to. + +== Maven coordinates + +[source,xml] +---- +<dependency> + <groupId>org.apache.camel.springboot</groupId> + <artifactId>camel-cli-connector-starter</artifactId> +</dependency> +---- + +== Usage + +Adding this starter to the classpath enables the connector with the file transport, used by the Camel CLI. It can be turned off with: + +[source,properties] +---- +camel.cli.enabled = false +---- + +The `stop` action of the tooling (for example `camel stop`) closes the Spring application context, so the application exits. + +=== WebSocket transport + +With `camel.cli.transport = websocket`, the application connects to the tool at `camel.cli.websocket.url`. See the CLI Connector documentation of Camel for the protocol and the security rules: *the tool gets full control of the application*, the transport is for development only and is refused with the `prod` profile. + +[source,properties] +---- +camel.cli.transport = websocket +camel.cli.websocket.url = ws://127.0.0.1:8000/connect +---- + +The WebSocket client of Spring is used when the application has `spring-websocket` and a Jakarta WebSocket client, for example with `spring-boot-starter-websocket`. Otherwise, the JDK client is used: the starter does not add any dependency for it. The log, and the `transport` field of the `hello` frame sent to the tool, say which client is used (`spring` or `jdk`). Set `camel.cli.websocket.client = jdk` to always use the JDK client. + +For a `wss://` URL, the Spring client can trust the tool with an SSL bundle: + +[source,properties] +---- +spring.ssl.bundle.pem.tool.truststore.certificate = classpath:tool.crt +camel.cli.websocket.url = wss://tool.example.com/connect +camel.cli.websocket.token = ${TOOL_TOKEN} +camel.cli.websocket.ssl-bundle = tool +---- + +== Spring Boot Auto-Configuration + +The starter supports 11 options, which are listed below. + +[width="100%",cols="2,5,^1,2",options="header"] +|=== +| Name | Description | Default | Type +| camel.cli.enabled | Whether CLI connector is enabled. | true | Boolean +| camel.cli.transport | How the connector talks to the tooling: file (exchanges files in ~/.camel with the Camel CLI) or websocket (dials out to a developer tool over a WebSocket, see camel.cli.websocket). | file | String +| camel.cli.websocket.allow-insecure | Allow ws:// (unencrypted) to a host that is not a loopback address. Prefer wss:// or a tunnel. | false | Boolean +| camel.cli.websocket.client | The WebSocket client: auto uses the Spring WebSocket client when the application has spring-websocket and a Jakarta WebSocket implementation (such as with spring-boot-starter-websocket), otherwise the JDK client; jdk always uses the JDK client. | auto | String +| camel.cli.websocket.heartbeat-interval | How often the connection is checked with a ping, in millis. The connection is re-established when the tool does not answer for three intervals. | 10000 | Long +| camel.cli.websocket.reconnect-delay | Delay before reconnecting, in millis. It doubles after each failed attempt (with some random jitter). | 1000 | Long +| camel.cli.websocket.reconnect-max-delay | Maximum delay before reconnecting, in millis. | 30000 | Long +| camel.cli.websocket.snapshot-interval | How often snapshots are sent, in millis. | 1000 | Long +| camel.cli.websocket.ssl-bundle | The name of an SSL bundle (spring.ssl.bundle) to trust the tool with wss:// when using the Spring WebSocket client. | | String +| camel.cli.websocket.token | Sent as Authorization: Bearer token when connecting. Required unless the URL is a loopback address. Prefer the CAMEL_CLI_WEBSOCKET_TOKEN environment variable to keep it out of the process list. | | String +| camel.cli.websocket.url | The ws:// or wss:// URL of the tool. Required when camel.cli.transport=websocket. | | String +|=== diff --git a/dsl-starter/camel-cli-connector-starter/pom.xml b/dsl-starter/camel-cli-connector-starter/pom.xml index bcce6c763f5..cc021eebf07 100644 --- a/dsl-starter/camel-cli-connector-starter/pom.xml +++ b/dsl-starter/camel-cli-connector-starter/pom.xml @@ -44,5 +44,53 @@ <artifactId>camel-cli-connector</artifactId> <version>${camel-version}</version> </dependency> + <!-- generates META-INF/spring-configuration-metadata.json, so IDEs complete the camel.cli options --> + <dependency> + <groupId>org.springframework.boot</groupId> + <artifactId>spring-boot-configuration-processor</artifactId> + <version>${spring-boot-version}</version> + <scope>provided</scope> + </dependency> + <!-- the Spring WebSocket client is used when the application has it, otherwise the JDK one --> + <dependency> + <groupId>org.springframework</groupId> + <artifactId>spring-websocket</artifactId> + <version>${spring-version}</version> + <optional>true</optional> + </dependency> + <dependency> + <groupId>jakarta.websocket</groupId> + <artifactId>jakarta.websocket-client-api</artifactId> + <optional>true</optional> + </dependency> + <dependency> + <groupId>org.apache.camel.springboot</groupId> + <artifactId>camel-spring-boot-starter</artifactId> + <version>${project.version}</version> + <scope>test</scope> + </dependency> + <dependency> + <groupId>org.apache.camel</groupId> + <artifactId>camel-test-spring-junit6</artifactId> + <version>${camel-version}</version> + <scope>test</scope> + </dependency> + <dependency> + <groupId>org.springframework.boot</groupId> + <artifactId>spring-boot-starter-test</artifactId> + <version>${spring-boot-version}</version> + <scope>test</scope> + </dependency> + <dependency> + <groupId>org.springframework.boot</groupId> + <artifactId>spring-boot-starter-websocket</artifactId> + <version>${spring-boot-version}</version> + <scope>test</scope> + </dependency> + <dependency> + <groupId>org.awaitility</groupId> + <artifactId>awaitility</artifactId> + <scope>test</scope> + </dependency> </dependencies> </project> diff --git a/dsl-starter/camel-cli-connector-starter/src/main/doc/intro.adoc b/dsl-starter/camel-cli-connector-starter/src/main/doc/intro.adoc new file mode 100644 index 00000000000..758a3390e22 --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/main/doc/intro.adoc @@ -0,0 +1,3 @@ +Spring Boot auto-configuration for the Camel CLI connector. + +The connector lets developer tools manage a running Spring Boot application: the Camel CLI (`camel get`, `camel trace`, `camel stop`, ...) through files in `~/.camel`, or a tool such as an IDE over a WebSocket the application dials out to. diff --git a/dsl-starter/camel-cli-connector-starter/src/main/doc/usage.adoc b/dsl-starter/camel-cli-connector-starter/src/main/doc/usage.adoc new file mode 100644 index 00000000000..149fc023ed1 --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/main/doc/usage.adoc @@ -0,0 +1,30 @@ +Adding this starter to the classpath enables the connector with the file transport, used by the Camel CLI. It can be turned off with: + +[source,properties] +---- +camel.cli.enabled = false +---- + +The `stop` action of the tooling (for example `camel stop`) closes the Spring application context, so the application exits. + +=== WebSocket transport + +With `camel.cli.transport = websocket`, the application connects to the tool at `camel.cli.websocket.url`. See the CLI Connector documentation of Camel for the protocol and the security rules: *the tool gets full control of the application*, the transport is for development only and is refused with the `prod` profile. + +[source,properties] +---- +camel.cli.transport = websocket +camel.cli.websocket.url = ws://127.0.0.1:8000/connect +---- + +The WebSocket client of Spring is used when the application has `spring-websocket` and a Jakarta WebSocket client, for example with `spring-boot-starter-websocket`. Otherwise, the JDK client is used: the starter does not add any dependency for it. The log, and the `transport` field of the `hello` frame sent to the tool, say which client is used (`spring` or `jdk`). Set `camel.cli.websocket.client = jdk` to always use the JDK client. + +For a `wss://` URL, the Spring client can trust the tool with an SSL bundle: + +[source,properties] +---- +spring.ssl.bundle.pem.tool.truststore.certificate = classpath:tool.crt +camel.cli.websocket.url = wss://tool.example.com/connect +camel.cli.websocket.token = ${TOOL_TOKEN} +camel.cli.websocket.ssl-bundle = tool +---- diff --git a/dsl-starter/camel-cli-connector-starter/src/main/docs/cli-connector.json b/dsl-starter/camel-cli-connector-starter/src/main/docs/cli-connector.json index fe81755c7df..542ebb608bc 100644 --- a/dsl-starter/camel-cli-connector-starter/src/main/docs/cli-connector.json +++ b/dsl-starter/camel-cli-connector-starter/src/main/docs/cli-connector.json @@ -4,6 +4,12 @@ "name": "camel.cli", "type": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration", "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration" + }, + { + "name": "camel.cli.websocket", + "type": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration", + "sourceMethod": "getWebsocket()" } ], "properties": [ @@ -13,7 +19,77 @@ "description": "Whether CLI connector is enabled.", "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration", "defaultValue": true + }, + { + "name": "camel.cli.transport", + "type": "java.lang.String", + "description": "How the connector talks to the tooling: file (exchanges files in ~\/.camel with the Camel CLI) or websocket (dials out to a developer tool over a WebSocket, see camel.cli.websocket).", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration", + "defaultValue": "file" + }, + { + "name": "camel.cli.websocket.allow-insecure", + "type": "java.lang.Boolean", + "description": "Allow ws:\/\/ (unencrypted) to a host that is not a loopback address. Prefer wss:\/\/ or a tunnel.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket", + "defaultValue": false + }, + { + "name": "camel.cli.websocket.client", + "type": "java.lang.String", + "description": "The WebSocket client: auto uses the Spring WebSocket client when the application has spring-websocket and a Jakarta WebSocket implementation (such as with spring-boot-starter-websocket), otherwise the JDK client; jdk always uses the JDK client.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket", + "defaultValue": "auto" + }, + { + "name": "camel.cli.websocket.heartbeat-interval", + "type": "java.lang.Long", + "description": "How often the connection is checked with a ping, in millis. The connection is re-established when the tool does not answer for three intervals.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket", + "defaultValue": 10000 + }, + { + "name": "camel.cli.websocket.reconnect-delay", + "type": "java.lang.Long", + "description": "Delay before reconnecting, in millis. It doubles after each failed attempt (with some random jitter).", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket", + "defaultValue": 1000 + }, + { + "name": "camel.cli.websocket.reconnect-max-delay", + "type": "java.lang.Long", + "description": "Maximum delay before reconnecting, in millis.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket", + "defaultValue": 30000 + }, + { + "name": "camel.cli.websocket.snapshot-interval", + "type": "java.lang.Long", + "description": "How often snapshots are sent, in millis.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket", + "defaultValue": 1000 + }, + { + "name": "camel.cli.websocket.ssl-bundle", + "type": "java.lang.String", + "description": "The name of an SSL bundle (spring.ssl.bundle) to trust the tool with wss:\/\/ when using the Spring WebSocket client.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket" + }, + { + "name": "camel.cli.websocket.token", + "type": "java.lang.String", + "description": "Sent as Authorization: Bearer token when connecting. Required unless the URL is a loopback address. Prefer the CAMEL_CLI_WEBSOCKET_TOKEN environment variable to keep it out of the process list.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket" + }, + { + "name": "camel.cli.websocket.url", + "type": "java.lang.String", + "description": "The ws:\/\/ or wss:\/\/ URL of the tool. Required when camel.cli.transport=websocket.", + "sourceType": "org.apache.camel.springboot.cli.connector.CliConnectorConfiguration$Websocket" } ], - "hints": [] + "hints": [], + "ignored": { + "properties": [] + } } \ No newline at end of file diff --git a/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorAutoConfiguration.java b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorAutoConfiguration.java index 74c506466b2..bd3886d1950 100644 --- a/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorAutoConfiguration.java +++ b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorAutoConfiguration.java @@ -21,12 +21,16 @@ import java.net.URL; import java.util.Enumeration; import java.util.jar.Manifest; +import org.apache.camel.cli.connector.CliWebSocketClient; import org.apache.camel.spi.CliConnectorFactory; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.SpringBootVersion; import org.springframework.boot.autoconfigure.AutoConfigureBefore; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.boot.ssl.SslBundles; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.support.AbstractApplicationContext; @@ -69,4 +73,30 @@ public class CliConnectorAutoConfiguration { return answer; } + /** + * The Spring WebSocket client for the websocket transport, when the application has spring-websocket and a Jakarta + * WebSocket client (such as Tomcat with spring-boot-starter-websocket). Otherwise Camel uses the JDK client. + */ + @Configuration(proxyBeanMethods = false) + @ConditionalOnClass(name = { + "org.springframework.web.socket.client.standard.StandardWebSocketClient", + "jakarta.websocket.ContainerProvider" }) + static class SpringWebSocketClientConfiguration { + + @Bean + @ConditionalOnMissingBean(CliWebSocketClient.class) + public CliWebSocketClient cliWebSocketClient( + CliConnectorConfiguration config, ObjectProvider<SslBundles> sslBundles) { + String bundle = config.getWebsocket().getSslBundle(); + if (bundle == null || bundle.isBlank()) { + return new SpringCliWebSocketClient(); + } + SslBundles bundles = sslBundles.getIfAvailable(); + if (bundles == null) { + throw new IllegalStateException( + "camel.cli.websocket.ssl-bundle=" + bundle + " but there are no SSL bundles (spring.ssl.bundle)"); + } + return new SpringCliWebSocketClient(bundles.getBundle(bundle).createSslContext()); + } + } } diff --git a/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorConfiguration.java b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorConfiguration.java index d3ad2aaa314..5f3bccf4720 100644 --- a/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorConfiguration.java +++ b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/CliConnectorConfiguration.java @@ -28,6 +28,16 @@ public class CliConnectorConfiguration { * Whether CLI connector is enabled. */ private Boolean enabled = true; + /** + * How the connector talks to the tooling: file (exchanges files in ~/.camel with the Camel CLI) or websocket (dials + * out to a developer tool over a WebSocket, see camel.cli.websocket). + */ + private String transport = "file"; + /** + * The WebSocket transport (camel.cli.transport=websocket). It gives the connected tool full control of the + * application, it is for development only and refused with the prod profile. + */ + private Websocket websocket = new Websocket(); public Boolean getEnabled() { return enabled; @@ -36,4 +46,141 @@ public class CliConnectorConfiguration { public void setEnabled(Boolean enabled) { this.enabled = enabled; } + + public String getTransport() { + return transport; + } + + public void setTransport(String transport) { + this.transport = transport; + } + + public Websocket getWebsocket() { + return websocket; + } + + public void setWebsocket(Websocket websocket) { + this.websocket = websocket; + } + + /** + * Options of the WebSocket transport. Camel reads them from the Spring environment: they are declared here so they + * are documented and completed by IDEs. + */ + public static class Websocket { + + /** + * The ws:// or wss:// URL of the tool. Required when camel.cli.transport=websocket. + */ + private String url; + /** + * Sent as Authorization: Bearer token when connecting. Required unless the URL is a loopback address. Prefer + * the CAMEL_CLI_WEBSOCKET_TOKEN environment variable to keep it out of the process list. + */ + private String token; + /** + * How often snapshots are sent, in millis. + */ + private Long snapshotInterval = 1000L; + /** + * How often the connection is checked with a ping, in millis. The connection is re-established when the tool + * does not answer for three intervals. + */ + private Long heartbeatInterval = 10000L; + /** + * Delay before reconnecting, in millis. It doubles after each failed attempt (with some random jitter). + */ + private Long reconnectDelay = 1000L; + /** + * Maximum delay before reconnecting, in millis. + */ + private Long reconnectMaxDelay = 30000L; + /** + * Allow ws:// (unencrypted) to a host that is not a loopback address. Prefer wss:// or a tunnel. + */ + private Boolean allowInsecure = false; + /** + * The WebSocket client: auto uses the Spring WebSocket client when the application has spring-websocket and a + * Jakarta WebSocket implementation (such as with spring-boot-starter-websocket), otherwise the JDK client; jdk + * always uses the JDK client. + */ + private String client = "auto"; + /** + * The name of an SSL bundle (spring.ssl.bundle) to trust the tool with wss:// when using the Spring WebSocket + * client. + */ + private String sslBundle; + + public String getUrl() { + return url; + } + + public void setUrl(String url) { + this.url = url; + } + + public String getToken() { + return token; + } + + public void setToken(String token) { + this.token = token; + } + + public Long getSnapshotInterval() { + return snapshotInterval; + } + + public void setSnapshotInterval(Long snapshotInterval) { + this.snapshotInterval = snapshotInterval; + } + + public Long getHeartbeatInterval() { + return heartbeatInterval; + } + + public void setHeartbeatInterval(Long heartbeatInterval) { + this.heartbeatInterval = heartbeatInterval; + } + + public Long getReconnectDelay() { + return reconnectDelay; + } + + public void setReconnectDelay(Long reconnectDelay) { + this.reconnectDelay = reconnectDelay; + } + + public Long getReconnectMaxDelay() { + return reconnectMaxDelay; + } + + public void setReconnectMaxDelay(Long reconnectMaxDelay) { + this.reconnectMaxDelay = reconnectMaxDelay; + } + + public Boolean getAllowInsecure() { + return allowInsecure; + } + + public void setAllowInsecure(Boolean allowInsecure) { + this.allowInsecure = allowInsecure; + } + + public String getClient() { + return client; + } + + public void setClient(String client) { + this.client = client; + } + + public String getSslBundle() { + return sslBundle; + } + + public void setSslBundle(String sslBundle) { + this.sslBundle = sslBundle; + } + } } diff --git a/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/SpringCliWebSocketClient.java b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/SpringCliWebSocketClient.java new file mode 100644 index 00000000000..781a61aca9b --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/SpringCliWebSocketClient.java @@ -0,0 +1,263 @@ +/* + * 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.springboot.cli.connector; + +import java.io.IOException; +import java.net.URI; +import java.nio.ByteBuffer; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CompletionStage; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import javax.net.ssl.SSLContext; + +import jakarta.websocket.CloseReason; +import jakarta.websocket.ContainerProvider; +import jakarta.websocket.DeploymentException; +import jakarta.websocket.Session; + +import org.apache.camel.cli.connector.CliWebSocketClient; +import org.apache.camel.cli.connector.CliWebSocketHandshakeException; +import org.springframework.core.task.SimpleAsyncTaskExecutor; +import org.springframework.web.socket.CloseStatus; +import org.springframework.web.socket.adapter.NativeWebSocketSession; +import org.springframework.web.socket.PongMessage; +import org.springframework.web.socket.TextMessage; +import org.springframework.web.socket.WebSocketHttpHeaders; +import org.springframework.web.socket.WebSocketSession; +import org.springframework.web.socket.client.standard.StandardWebSocketClient; +import org.springframework.web.socket.handler.AbstractWebSocketHandler; + +/** + * The Camel CLI connector WebSocket client of Spring ({@link StandardWebSocketClient}, Jakarta WebSocket, for example + * Tomcat's <tt>tomcat-embed-websocket</tt>), used when the application has it instead of the JDK client. + * <p/> + * Sends go through the asynchronous remote endpoint of the Jakarta session, so a send that hangs does not block a + * thread and can be aborted by the transport. Partial messages are reassembled here: raising the Jakarta text buffer + * to the largest accepted message would allocate that buffer for every connection. + */ +public class SpringCliWebSocketClient implements CliWebSocketClient { + + // Tomcat client options (user properties of the endpoint config), ignored by other Jakarta WebSocket clients + static final String IO_TIMEOUT_PROPERTY = "org.apache.tomcat.websocket.IO_TIMEOUT_MS"; + static final String BLOCKING_SEND_TIMEOUT_PROPERTY = "org.apache.tomcat.websocket.BLOCKING_SEND_TIMEOUT"; + private static final long TIMEOUT = 10000; + + // Jakarta WebSocket has no status code for a rejected upgrade: Tomcat puts it in the message, as [401] + private static final Pattern HTTP_STATUS = Pattern.compile("\\[([1-5][0-9]{2})]"); + + private final SSLContext sslContext; + private volatile StandardWebSocketClient client; + + public SpringCliWebSocketClient() { + this(null); + } + + /** + * @param sslContext for wss:// urls, or null for the default one + */ + public SpringCliWebSocketClient(SSLContext sslContext) { + this.sslContext = sslContext; + } + + @Override + public String getName() { + return "spring"; + } + + SSLContext getSslContext() { + return sslContext; + } + + @Override + public CompletionStage<Channel> connect(URI url, Map<String, String> headers, Listener listener) { + WebSocketHttpHeaders handshakeHeaders = new WebSocketHttpHeaders(); + headers.forEach(handshakeHeaders::add); + Handler handler = new Handler(listener); + try { + return client().execute(handler, handshakeHeaders, url) + .handle((session, e) -> { + if (e != null) { + throw new CompletionException(translate(e)); + } + return handler.channel; + }); + } catch (RuntimeException e) { + // no Jakarta WebSocket implementation on the classpath + return CompletableFuture.failedFuture(e); + } + } + + /** + * Created on first use, so the application starts even without a Jakarta WebSocket implementation when this client + * is not used. + */ + private StandardWebSocketClient client() { + StandardWebSocketClient answer = client; + if (answer == null) { + synchronized (this) { + answer = client; + if (answer == null) { + answer = new StandardWebSocketClient(ContainerProvider.getWebSocketContainer()); + answer.setSslContext(sslContext); + Map<String, Object> properties = new HashMap<>(); + // the JDK client times out a connect after 10 seconds too + properties.put(IO_TIMEOUT_PROPERTY, Long.toString(TIMEOUT)); + // pings and close frames are blocking sends + properties.put(BLOCKING_SEND_TIMEOUT_PROPERTY, TIMEOUT); + answer.setUserProperties(properties); + answer.setTaskExecutor(new SimpleAsyncTaskExecutor("CliConnectorWebSocketConnect-")); + client = answer; + } + } + } + return answer; + } + + static Throwable translate(Throwable e) { + Throwable cause = e instanceof CompletionException && e.getCause() != null ? e.getCause() : e; + for (Throwable t = cause; t != null; t = t.getCause() != t ? t.getCause() : null) { + if (t instanceof DeploymentException && t.getMessage() != null) { + Matcher m = HTTP_STATUS.matcher(t.getMessage()); + if (m.find()) { + return new CliWebSocketHandshakeException(Integer.parseInt(m.group(1)), cause); + } + } + } + return cause; + } + + private static final class SpringChannel implements Channel { + + private final Session session; + + SpringChannel(WebSocketSession session) { + this.session = ((NativeWebSocketSession) session).getNativeSession(Session.class); + } + + @Override + public CompletionStage<?> sendText(String text) { + CompletableFuture<Void> answer = new CompletableFuture<>(); + try { + session.getAsyncRemote().sendText(text, result -> { + if (result.isOK()) { + answer.complete(null); + } else { + answer.completeExceptionally(result.getException()); + } + }); + } catch (RuntimeException e) { + answer.completeExceptionally(e); + } + return answer; + } + + @Override + public CompletionStage<?> sendPing() { + try { + session.getAsyncRemote().sendPing(ByteBuffer.allocate(0)); + return CompletableFuture.completedFuture(null); + } catch (IOException | RuntimeException e) { + return CompletableFuture.failedFuture(e); + } + } + + @Override + public CompletionStage<?> close(int code, String reason) { + try { + session.close(new CloseReason(CloseReason.CloseCodes.getCloseCode(code), reason)); + return CompletableFuture.completedFuture(null); + } catch (IOException | RuntimeException e) { + return CompletableFuture.failedFuture(e); + } + } + + @Override + public void abort() { + abort(session); + } + + static void abort(WebSocketSession session) { + abort(((NativeWebSocketSession) session).getNativeSession(Session.class)); + } + + static void abort(Session session) { + try { + // not sent as such: Tomcat tries a close frame for a short time only, then closes the socket + session.close(new CloseReason(CloseReason.CloseCodes.CLOSED_ABNORMALLY, "aborted")); + } catch (IOException | RuntimeException e) { + // already closed + } + } + } + + private static final class Handler extends AbstractWebSocketHandler { + + private final Listener listener; + private final StringBuilder partial = new StringBuilder(); + private volatile Channel channel; + + Handler(Listener listener) { + this.listener = listener; + } + + @Override + public boolean supportsPartialMessages() { + return true; + } + + @Override + public void afterConnectionEstablished(WebSocketSession session) { + // before connect() completes, and before any message + channel = new SpringChannel(session); + } + + @Override + protected void handleTextMessage(WebSocketSession session, TextMessage message) { + partial.append(message.getPayload()); + if (partial.length() > MAX_MESSAGE_SIZE) { + partial.setLength(0); + // as the JDK client: stop reading right away, the transport reconnects + SpringChannel.abort(session); + listener.onError(new IOException("Message larger than " + MAX_MESSAGE_SIZE + " chars")); + } else if (message.isLast()) { + String text = partial.toString(); + partial.setLength(0); + listener.onText(text); + } + } + + @Override + protected void handlePongMessage(WebSocketSession session, PongMessage message) { + listener.onPong(); + } + + @Override + public void handleTransportError(WebSocketSession session, Throwable exception) { + listener.onError(exception); + } + + @Override + public void afterConnectionClosed(WebSocketSession session, CloseStatus status) { + listener.onClose(status.getCode(), status.getReason()); + } + } +} diff --git a/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/SpringLocalCliConnector.java b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/SpringLocalCliConnector.java index 35e61384edd..23c63db2843 100644 --- a/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/SpringLocalCliConnector.java +++ b/dsl-starter/camel-cli-connector-starter/src/main/java/org/apache/camel/springboot/cli/connector/SpringLocalCliConnector.java @@ -35,7 +35,9 @@ public class SpringLocalCliConnector extends LocalCliConnector { try { super.sigterm(); } finally { - applicationContext.stop(); + // close, not only stop: a stopped context keeps the JVM running (for example the embedded web server), + // while the Camel CLI (camel stop) and the websocket tools expect the application to exit + applicationContext.close(); } } } diff --git a/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorAutoConfigurationTest.java b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorAutoConfigurationTest.java new file mode 100644 index 00000000000..844353943b2 --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorAutoConfigurationTest.java @@ -0,0 +1,84 @@ +/* + * 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.springboot.cli.connector; + +import javax.net.ssl.SSLContext; + +import org.apache.camel.cli.connector.CliWebSocketClient; +import org.junit.jupiter.api.Test; +import org.springframework.boot.autoconfigure.AutoConfigurations; +import org.springframework.boot.ssl.DefaultSslBundleRegistry; +import org.springframework.boot.ssl.SslBundle; +import org.springframework.boot.ssl.SslBundles; +import org.springframework.boot.ssl.SslStoreBundle; +import org.springframework.boot.test.context.FilteredClassLoader; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.web.socket.client.standard.StandardWebSocketClient; + +import static org.assertj.core.api.Assertions.assertThat; + +class CliConnectorAutoConfigurationTest { + + private final ApplicationContextRunner runner = new ApplicationContextRunner() + .withConfiguration(AutoConfigurations.of(CliConnectorAutoConfiguration.class)); + + @Test + void usesTheSpringClientWhenTheApplicationHasSpringWebSocket() { + runner.run(context -> assertThat(context).getBean(CliWebSocketClient.class) + .isInstanceOf(SpringCliWebSocketClient.class) + .extracting(CliWebSocketClient::getName).isEqualTo("spring")); + } + + @Test + void leavesTheJdkClientToCamelWithoutSpringWebSocket() { + runner.withClassLoader(new FilteredClassLoader(StandardWebSocketClient.class)) + .run(context -> assertThat(context).hasNotFailed().doesNotHaveBean(CliWebSocketClient.class)); + } + + @Test + void leavesTheJdkClientToCamelWithoutJakartaWebSocket() { + runner.withClassLoader(new FilteredClassLoader("jakarta.websocket.")) + .run(context -> assertThat(context).hasNotFailed().doesNotHaveBean(CliWebSocketClient.class)); + } + + @Test + void keepsTheClientOfTheApplication() { + CliWebSocketClient custom = new SpringCliWebSocketClient(); + runner.withBean(CliWebSocketClient.class, () -> custom) + .run(context -> assertThat(context).getBean(CliWebSocketClient.class).isSameAs(custom)); + } + + @Test + void usesTheSslBundle() throws Exception { + SslBundle bundle = SslBundle.of(SslStoreBundle.NONE); + runner.withBean(SslBundles.class, () -> new DefaultSslBundleRegistry("tool", bundle)) + .withPropertyValues("camel.cli.websocket.ssl-bundle=tool") + .run(context -> { + SSLContext ssl = ((SpringCliWebSocketClient) context.getBean(CliWebSocketClient.class)).getSslContext(); + assertThat(ssl).isNotNull(); + assertThat(ssl.getProtocol()).isEqualTo(bundle.getProtocol()); + }); + } + + @Test + void failsOnAnSslBundleWithoutSslBundles() { + runner.withPropertyValues("camel.cli.websocket.ssl-bundle=tool") + .run(context -> assertThat(context).hasFailed().getFailure() + .hasRootCauseMessage( + "camel.cli.websocket.ssl-bundle=tool but there are no SSL bundles (spring.ssl.bundle)")); + } +} diff --git a/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorStopActionTest.java b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorStopActionTest.java new file mode 100644 index 00000000000..a35aa132bb7 --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorStopActionTest.java @@ -0,0 +1,68 @@ +/* + * 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.springboot.cli.connector; + +import java.util.Map; +import java.util.concurrent.TimeUnit; + +import org.apache.camel.spring.boot.CamelAutoConfiguration; +import org.apache.camel.util.json.JsonObject; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.context.ConfigurableApplicationContext; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; + +/** + * The stop action closes the Spring application context, so the application exits. + */ +class CliConnectorStopActionTest { + + private final ToolServer tool = new ToolServer(); + private ConfigurableApplicationContext context; + + @AfterEach + void stopAll() throws Exception { + if (context != null) { + context.close(); + } + tool.close(); + } + + @Test + void stopActionClosesTheApplicationContext() throws Exception { + tool.start(); + context = new SpringApplicationBuilder( + CamelAutoConfiguration.class, CliConnectorAutoConfiguration.class, + CliConnectorWebSocketTestSupport.Routes.class) + .web(WebApplicationType.NONE) + .properties( + "camel.cli.transport=websocket", + "camel.cli.websocket.url=" + tool.url()) + .run(); + tool.awaitFrame(f -> "hello".equals(f.getString("type"))); + + JsonObject action = new JsonObject(Map.of("action", "stop")); + tool.send(new JsonObject(Map.of("v", 1, "type", "action", "requestId", "r1", "action", action))); + + assertThat(tool.awaitResult("r1").getBoolean("ok")).isTrue(); + await().atMost(20, TimeUnit.SECONDS).untilAsserted(() -> assertThat(context.isActive()).isFalse()); + } +} diff --git a/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketJdkClientTest.java b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketJdkClientTest.java new file mode 100644 index 00000000000..7735864adf4 --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketJdkClientTest.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.springboot.cli.connector; + +import org.apache.camel.spring.boot.CamelAutoConfiguration; +import org.apache.camel.test.spring.junit6.CamelSpringBootTest; +import org.apache.camel.test.spring.junit6.DisableJmx; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.annotation.DirtiesContext; + +/** + * The connector uses the JDK client when asked to, even with the Spring WebSocket client available. + */ +@DirtiesContext +@CamelSpringBootTest +// the status sent to the tool is built from the JMX management layer +@DisableJmx(false) +@SpringBootTest(classes = { + CamelAutoConfiguration.class, CliConnectorAutoConfiguration.class, + CliConnectorWebSocketTestSupport.Routes.class }, + properties = { + "camel.cli.transport=websocket", + "camel.cli.websocket.snapshot-interval=200", + "camel.cli.websocket.client=jdk" }) +class CliConnectorWebSocketJdkClientTest extends CliConnectorWebSocketTestSupport { + + @Override + String expectedClient() { + return "jdk"; + } +} diff --git a/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketSpringClientTest.java b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketSpringClientTest.java new file mode 100644 index 00000000000..ae50ef557ad --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketSpringClientTest.java @@ -0,0 +1,44 @@ +/* + * 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.springboot.cli.connector; + +import org.apache.camel.spring.boot.CamelAutoConfiguration; +import org.apache.camel.test.spring.junit6.CamelSpringBootTest; +import org.apache.camel.test.spring.junit6.DisableJmx; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.annotation.DirtiesContext; + +/** + * The connector uses the Spring WebSocket client, as the application has spring-websocket and Tomcat. + */ +@DirtiesContext +@CamelSpringBootTest +// the status sent to the tool is built from the JMX management layer +@DisableJmx(false) +@SpringBootTest(classes = { + CamelAutoConfiguration.class, CliConnectorAutoConfiguration.class, + CliConnectorWebSocketTestSupport.Routes.class }, + properties = { + "camel.cli.transport=websocket", + "camel.cli.websocket.snapshot-interval=200" }) +class CliConnectorWebSocketSpringClientTest extends CliConnectorWebSocketTestSupport { + + @Override + String expectedClient() { + return "spring"; + } +} diff --git a/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketTestSupport.java b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketTestSupport.java new file mode 100644 index 00000000000..dc033cc75cb --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/CliConnectorWebSocketTestSupport.java @@ -0,0 +1,91 @@ +/* + * 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.springboot.cli.connector; + +import java.util.Map; + +import org.apache.camel.CamelContext; +import org.apache.camel.ServiceStatus; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.util.json.JsonObject; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * A Spring Boot application with camel.cli.transport=websocket connects to the tool, says hello with the name of its + * WebSocket client, and executes an action. + */ +abstract class CliConnectorWebSocketTestSupport { + + static ToolServer tool; + + @Autowired + CamelContext context; + + @DynamicPropertySource + static void tool(DynamicPropertyRegistry registry) throws Exception { + tool = new ToolServer().start(); + tool.requiredToken = "t0k3n-4-t3st"; + // read by Camel from the Spring environment + registry.add("camel.cli.websocket.url", tool::url); + registry.add("camel.cli.websocket.token", () -> tool.requiredToken); + } + + @AfterAll + static void stopTool() throws Exception { + tool.close(); + } + + abstract String expectedClient(); + + @Test + void connectsAndExecutesActions() throws Exception { + JsonObject hello = tool.awaitFrame(f -> "hello".equals(f.getString("type"))); + assertThat(hello.getString("transport")).isEqualTo(expectedClient()); + assertThat(hello.getString("name")).isEqualTo(context.getName()); + + JsonObject action = new JsonObject(Map.of("action", "route", "command", "stop", "id", "hello")); + tool.send(new JsonObject(Map.of("v", 1, "type", "action", "requestId", "r1", "action", action))); + + assertThat(tool.awaitResult("r1").getBoolean("ok")).isTrue(); + assertThat(context.getRouteController().getRouteStatus("hello")).isEqualTo(ServiceStatus.Stopped); + JsonObject status = tool.awaitFrame(f -> "snapshot".equals(f.getString("type")) + && "status".equals(f.getString("kind"))).getMap("data"); + assertThat(status.toJson()).contains("\"routeId\":\"hello\""); + } + + @Configuration + static class Routes { + + @Bean + RouteBuilder routes() { + return new RouteBuilder() { + @Override + public void configure() { + from("direct:hello").routeId("hello").setBody(simple("Hello ${body}")); + } + }; + } + } +} diff --git a/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/SpringCliWebSocketClientTest.java b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/SpringCliWebSocketClientTest.java new file mode 100644 index 00000000000..b12937ce123 --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/SpringCliWebSocketClientTest.java @@ -0,0 +1,148 @@ +/* + * 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.springboot.cli.connector; + +import java.net.URI; +import java.util.Map; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.CompletionException; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; + +import org.apache.camel.cli.connector.CliWebSocketClient; +import org.apache.camel.cli.connector.CliWebSocketHandshakeException; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.awaitility.Awaitility.await; + +class SpringCliWebSocketClientTest { + + private final ToolServer tool = new ToolServer(); + private final SpringCliWebSocketClient client = new SpringCliWebSocketClient(); + private final Recorder listener = new Recorder(); + + @BeforeEach + void startTool() throws Exception { + tool.start(); + } + + @AfterEach + void stopTool() throws Exception { + tool.close(); + } + + @Test + void sendsTheHeadersAndReassemblesLargeMessages() throws Exception { + tool.requiredToken = "t0k3n"; + // much larger than the Jakarta WebSocket text buffer (8 KB): Tomcat passes it in parts + String large = "x".repeat(2 * 1024 * 1024); + tool.greeting = "{\"big\":\"" + large + "\"}"; + + CliWebSocketClient.Channel channel = client.connect(URI.create(tool.url()), + Map.of("Authorization", "Bearer t0k3n"), listener).toCompletableFuture().get(10, TimeUnit.SECONDS); + + String text = listener.texts.poll(10, TimeUnit.SECONDS); + assertThat(text).isEqualTo(tool.greeting); + + channel.sendText("{\"v\":1}").toCompletableFuture().get(10, TimeUnit.SECONDS); + assertThat(tool.frames.poll(10, TimeUnit.SECONDS).getInteger("v")).isEqualTo(1); + + channel.sendPing().toCompletableFuture().get(10, TimeUnit.SECONDS); + await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> assertThat(listener.pongs).isPositive()); + + channel.close(1000, "bye").toCompletableFuture().get(10, TimeUnit.SECONDS); + await().atMost(10, TimeUnit.SECONDS).until(tool.sessions::isEmpty); + } + + @Test + void failsTheConnectionOnAMessageLargerThanTheLimit() throws Exception { + tool.greeting = "x".repeat(CliWebSocketClient.MAX_MESSAGE_SIZE + 1); + + client.connect(URI.create(tool.url()), Map.of(), listener).toCompletableFuture().get(10, TimeUnit.SECONDS); + + await().atMost(20, TimeUnit.SECONDS).untilAsserted(() -> assertThat(listener.error) + .hasMessage("Message larger than " + CliWebSocketClient.MAX_MESSAGE_SIZE + " chars")); + await().atMost(10, TimeUnit.SECONDS).until(tool.sessions::isEmpty); + assertThat(listener.texts).isEmpty(); + } + + @Test + void reportsTheHttpStatusOfARejectedHandshake() { + tool.requiredToken = "t0k3n"; + + assertThatThrownBy(() -> client.connect(URI.create(tool.url()), Map.of("Authorization", "Bearer wrong"), listener) + .toCompletableFuture().join()) + .isInstanceOf(CompletionException.class) + .cause().isInstanceOfSatisfying(CliWebSocketHandshakeException.class, + e -> assertThat(e.getStatusCode()).isEqualTo(401)); + } + + @Test + void reportsWhenTheToolClosesTheConnection() throws Exception { + client.connect(URI.create(tool.url()), Map.of(), listener).toCompletableFuture().get(10, TimeUnit.SECONDS); + await().atMost(10, TimeUnit.SECONDS).until(() -> !tool.sessions.isEmpty()); + + tool.sessions.get(0).close(); + + await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> assertThat(listener.closeCode).isEqualTo(1000)); + } + + @Test + void abortClosesTheConnection() throws Exception { + CliWebSocketClient.Channel channel = client.connect(URI.create(tool.url()), Map.of(), listener) + .toCompletableFuture().get(10, TimeUnit.SECONDS); + await().atMost(10, TimeUnit.SECONDS).until(() -> !tool.sessions.isEmpty()); + + channel.abort(); + + await().atMost(10, TimeUnit.SECONDS).until(tool.sessions::isEmpty); + assertThatThrownBy(() -> channel.sendText("{}").toCompletableFuture().get(10, TimeUnit.SECONDS)) + .hasRootCauseInstanceOf(IllegalStateException.class); + } + + private static class Recorder implements CliWebSocketClient.Listener { + + final BlockingQueue<String> texts = new LinkedBlockingQueue<>(); + volatile int pongs; + volatile int closeCode; + volatile Throwable error; + + @Override + public void onText(String text) { + texts.add(text); + } + + @Override + public void onPong() { + pongs++; + } + + @Override + public void onClose(int code, String reason) { + closeCode = code; + } + + @Override + public void onError(Throwable error) { + this.error = error; + } + } +} diff --git a/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/ToolServer.java b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/ToolServer.java new file mode 100644 index 00000000000..ec715b6a4b3 --- /dev/null +++ b/dsl-starter/camel-cli-connector-starter/src/test/java/org/apache/camel/springboot/cli/connector/ToolServer.java @@ -0,0 +1,180 @@ +/* + * 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.springboot.cli.connector; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.function.Predicate; + +import jakarta.servlet.http.HttpServlet; +import jakarta.servlet.http.HttpServletRequest; +import jakarta.servlet.http.HttpServletResponse; +import jakarta.websocket.CloseReason; +import jakarta.websocket.Endpoint; +import jakarta.websocket.EndpointConfig; +import jakarta.websocket.MessageHandler; +import jakarta.websocket.Session; +import jakarta.websocket.server.ServerContainer; +import jakarta.websocket.server.ServerEndpointConfig; + +import org.apache.catalina.Context; +import org.apache.catalina.startup.Tomcat; +import org.apache.camel.util.json.JsonObject; +import org.apache.camel.util.json.Jsoner; +import org.apache.tomcat.util.descriptor.web.FilterDef; +import org.apache.tomcat.util.descriptor.web.FilterMap; +import org.apache.tomcat.websocket.server.WsSci; +import org.springframework.util.FileSystemUtils; + +/** + * Plays the tool: a WebSocket server (embedded Tomcat) the connector dials out to. + */ +class ToolServer implements AutoCloseable { + + final BlockingQueue<JsonObject> frames = new LinkedBlockingQueue<>(); + final List<Session> sessions = new CopyOnWriteArrayList<>(); + /** + * When set, a handshake without Authorization: Bearer requiredToken is rejected with HTTP 401. + */ + volatile String requiredToken; + /** + * When set, sent to every connection once open. + */ + volatile String greeting; + int port; + private Path baseDir; + private Tomcat tomcat; + + ToolServer start() throws Exception { + tomcat = new Tomcat(); + baseDir = Files.createTempDirectory("tool-server"); + tomcat.setBaseDir(baseDir.toString()); + tomcat.setPort(0); + tomcat.getConnector(); + Context ctx = tomcat.addContext("", null); + Tomcat.addServlet(ctx, "default", new HttpServlet() { + }); + ctx.addServletMappingDecoded("/*", "default"); + + FilterDef auth = new FilterDef(); + auth.setFilterName("auth"); + auth.setFilter((request, response, chain) -> { + String header = ((HttpServletRequest) request).getHeader("Authorization"); + if (requiredToken != null && !("Bearer " + requiredToken).equals(header)) { + ((HttpServletResponse) response).sendError(401); + } else { + chain.doFilter(request, response); + } + }); + ctx.addFilterDef(auth); + FilterMap authMap = new FilterMap(); + authMap.setFilterName("auth"); + authMap.addURLPattern("/*"); + ctx.addFilterMap(authMap); + + ctx.addServletContainerInitializer(new WsSci(), null); + ctx.addServletContainerInitializer((classes, servletContext) -> { + ServerContainer container = (ServerContainer) servletContext.getAttribute(ServerContainer.class.getName()); + // the status snapshot is larger than the default 8 KB + container.setDefaultMaxTextMessageBufferSize(1024 * 1024); + try { + container.addEndpoint(ServerEndpointConfig.Builder.create(ToolEndpoint.class, "/connect") + .configurator(new ServerEndpointConfig.Configurator() { + @Override + @SuppressWarnings("unchecked") + public <T> T getEndpointInstance(Class<T> endpointClass) { + return (T) new ToolEndpoint(); + } + }).build()); + } catch (Exception e) { + throw new IllegalStateException(e); + } + }, null); + tomcat.start(); + port = tomcat.getConnector().getLocalPort(); + return this; + } + + String url() { + return "ws://127.0.0.1:" + port + "/connect?executionId=it-1"; + } + + void send(JsonObject frame) throws IOException { + sessions.get(sessions.size() - 1).getBasicRemote().sendText(frame.toJson()); + } + + /** + * Waits for a frame matching the predicate, skipping the others (mostly snapshots). + */ + 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"); + } + + JsonObject awaitResult(String requestId) throws InterruptedException { + return awaitFrame(f -> "result".equals(f.getString("type")) && requestId.equals(f.getString("requestId"))); + } + + @Override + public void close() throws Exception { + if (tomcat != null) { + tomcat.stop(); + tomcat.destroy(); + tomcat = null; + } + if (baseDir != null) { + FileSystemUtils.deleteRecursively(baseDir); + baseDir = null; + } + } + + class ToolEndpoint extends Endpoint { + + @Override + public void onOpen(Session session, EndpointConfig config) { + sessions.add(session); + session.addMessageHandler(String.class, (MessageHandler.Whole<String>) text -> { + try { + frames.add((JsonObject) Jsoner.deserialize(text)); + } catch (Exception e) { + throw new IllegalStateException(e); + } + }); + String text = greeting; + if (text != null) { + session.getAsyncRemote().sendText(text); + } + } + + @Override + public void onClose(Session session, CloseReason closeReason) { + sessions.remove(session); + } + } +}
