davsclaus commented on code in PR #27145: URL: https://github.com/apache/camel/pull/27145#discussion_r4154637863
########## dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java: ########## @@ -0,0 +1,757 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.camel.cli.connector; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.WebSocket; +import java.net.http.WebSocketHandshakeException; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import org.apache.camel.CamelContext; +import org.apache.camel.ExtendedStartupListener; +import org.apache.camel.spi.PropertiesComponent; +import org.apache.camel.support.service.ServiceSupport; +import org.apache.camel.util.StringHelper; +import org.apache.camel.util.concurrent.CamelThreadFactory; +import org.apache.camel.util.json.JsonArray; +import org.apache.camel.util.json.JsonObject; +import org.apache.camel.util.json.Jsoner; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Transport that dials out to a developer tool over a WebSocket (JDK client), for tools that do not share the + * filesystem of the integration or need events pushed instead of polled. + * <p/> + * Every frame is a JSON envelope <tt>{"v":1,"type":...}</tt>. The tool sends <tt>action</tt> frames holding the same + * action JSON the Camel CLI writes to its action file, and gets <tt>result</tt> frames back correlated by + * <tt>requestId</tt>. The transport sends a <tt>hello</tt> frame once the integration is started, and <tt>snapshot</tt> + * frames (the content of the status, trace, ... files) while connected. + * <p/> + * The connection is kept alive with pings, and re-established with an exponential backoff when it is lost. + * <p/> + * Threads: actions run one at a time on their own thread (the dispatcher is not thread-safe, and an action can block + * for a long time); connecting, snapshots, heartbeats and every write to the socket run on a second thread, so the JDK + * WebSocket never has two sends in flight. + */ +public class WebSocketCliConnectorTransport extends ServiceSupport implements CliConnectorTransport { + + static final int VERSION = 1; + + private static final Logger LOG = LoggerFactory.getLogger(WebSocketCliConnectorTransport.class); + private static final long SEND_TIMEOUT = 10000; + private static final long STABLE_CONNECTION = 10000; + private static final int MAX_FRAME_SIZE = 16 * 1024 * 1024; + private static final int MAX_PENDING_ACTIONS = 64; + // WebSocket servers commonly refuse messages over 256 KB (Vert.x, Quarkus): trace and receive snapshots can be + // much larger (after a reconnect they hold every retained message), so they are split into frames of this size + private static final int MAX_SNAPSHOT_SIZE = 128 * 1024; + + private CamelContext camelContext; + private CliActionDispatcher dispatcher; + private CliSnapshotProducer snapshots; + private Runnable shutdown; + + private URI url; + private String token; + private long reconnectDelay; + private long reconnectMaxDelay; + private long snapshotInterval; + private long heartbeatInterval; + + private String threadNamePattern; + private HttpClient client; + private ThreadPoolExecutor actions; + private ScheduledExecutorService scheduler; + // fields below are only used from the scheduler thread, except the volatile ones + private volatile Connection connection; + private volatile boolean ready; + private volatile boolean stopping; + private boolean listenerAdded; + // the snapshot task, only (re)scheduled from the scheduler thread or before it runs anything + private ScheduledFuture<?> snapshotFuture; + private long debugInterval; + private ScheduledFuture<?> debugFuture; + private int failures; + private long ticks; + + @Override + public void configure( + CamelContext camelContext, CliActionDispatcher dispatcher, CliSnapshotProducer snapshots, Runnable shutdown) { + this.camelContext = camelContext; + this.dispatcher = dispatcher; + this.snapshots = snapshots; + this.shutdown = shutdown; + } + + @Override + protected void doStart() throws Exception { + if ("prod".equals(camelContext.getCamelContextExtension().getProfile())) { + throw new IllegalStateException( + "The Camel CLI connector websocket transport gives the connected tool full control of this" + + " application and cannot be used with the prod profile." + + " Remove camel.cli.transport=websocket."); + } + String value = property("camel.cli.websocket.url", null); + if (value == null || value.isBlank()) { + throw new IllegalArgumentException("camel.cli.websocket.url must be set when camel.cli.transport=websocket"); + } + url = URI.create(value); + String scheme = url.getScheme() != null ? url.getScheme().toLowerCase(Locale.ROOT) : ""; + if (!"ws".equals(scheme) && !"wss".equals(scheme)) { + throw new IllegalArgumentException("camel.cli.websocket.url must be a ws:// or wss:// url: " + value); + } + boolean loopback = isLoopback(url.getHost()); + token = property("camel.cli.websocket.token", null); + if ((token == null || token.isBlank()) && !loopback) { + throw new IllegalArgumentException( + "camel.cli.websocket.token must be set when camel.cli.websocket.url is not a loopback address"); + } + reconnectDelay = Long.parseLong(property("camel.cli.websocket.reconnect-delay", "1000")); + reconnectMaxDelay = Long.parseLong(property("camel.cli.websocket.reconnect-max-delay", "30000")); + snapshotInterval = Long.parseLong(property("camel.cli.websocket.snapshot-interval", "1000")); + // when debugging, the debug snapshot (breakpoints) is sent faster than the others, as the file transport does + debugInterval = camelContext.isDebugging() ? 100 : 0; + heartbeatInterval = Long.parseLong(property("camel.cli.websocket.heartbeat-interval", "10000")); + + LOG.warn("Camel CLI connector connects to {} which gets full control of this application (development use only)", + where()); + if ("ws".equals(scheme) && !loopback) { Review Comment: The docs say it clearly: over `ws://`, anyone on the network path can pose as the tool and get full control (`load`, `eval`, sending messages), and the bearer token also travels in clear text. The feature is opt-in and refused under `prod`, so this is outside the threat model. Still, would it be safer to *refuse* `ws://` to a non-loopback host unless the user explicitly opts in (for example `camel.cli.websocket.allow-insecure=true`)? Users would then end up on `wss://` or a tunnel, which is what the docs already recommend. ########## dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java: ########## @@ -0,0 +1,757 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.camel.cli.connector; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.WebSocket; +import java.net.http.WebSocketHandshakeException; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import org.apache.camel.CamelContext; +import org.apache.camel.ExtendedStartupListener; +import org.apache.camel.spi.PropertiesComponent; +import org.apache.camel.support.service.ServiceSupport; +import org.apache.camel.util.StringHelper; +import org.apache.camel.util.concurrent.CamelThreadFactory; +import org.apache.camel.util.json.JsonArray; +import org.apache.camel.util.json.JsonObject; +import org.apache.camel.util.json.Jsoner; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Transport that dials out to a developer tool over a WebSocket (JDK client), for tools that do not share the + * filesystem of the integration or need events pushed instead of polled. + * <p/> + * Every frame is a JSON envelope <tt>{"v":1,"type":...}</tt>. The tool sends <tt>action</tt> frames holding the same + * action JSON the Camel CLI writes to its action file, and gets <tt>result</tt> frames back correlated by + * <tt>requestId</tt>. The transport sends a <tt>hello</tt> frame once the integration is started, and <tt>snapshot</tt> + * frames (the content of the status, trace, ... files) while connected. + * <p/> + * The connection is kept alive with pings, and re-established with an exponential backoff when it is lost. + * <p/> + * Threads: actions run one at a time on their own thread (the dispatcher is not thread-safe, and an action can block + * for a long time); connecting, snapshots, heartbeats and every write to the socket run on a second thread, so the JDK + * WebSocket never has two sends in flight. + */ +public class WebSocketCliConnectorTransport extends ServiceSupport implements CliConnectorTransport { + + static final int VERSION = 1; + + private static final Logger LOG = LoggerFactory.getLogger(WebSocketCliConnectorTransport.class); + private static final long SEND_TIMEOUT = 10000; + private static final long STABLE_CONNECTION = 10000; + private static final int MAX_FRAME_SIZE = 16 * 1024 * 1024; + private static final int MAX_PENDING_ACTIONS = 64; + // WebSocket servers commonly refuse messages over 256 KB (Vert.x, Quarkus): trace and receive snapshots can be + // much larger (after a reconnect they hold every retained message), so they are split into frames of this size + private static final int MAX_SNAPSHOT_SIZE = 128 * 1024; + + private CamelContext camelContext; + private CliActionDispatcher dispatcher; + private CliSnapshotProducer snapshots; + private Runnable shutdown; + + private URI url; + private String token; + private long reconnectDelay; + private long reconnectMaxDelay; + private long snapshotInterval; + private long heartbeatInterval; + + private String threadNamePattern; + private HttpClient client; + private ThreadPoolExecutor actions; + private ScheduledExecutorService scheduler; + // fields below are only used from the scheduler thread, except the volatile ones + private volatile Connection connection; + private volatile boolean ready; + private volatile boolean stopping; + private boolean listenerAdded; + // the snapshot task, only (re)scheduled from the scheduler thread or before it runs anything + private ScheduledFuture<?> snapshotFuture; + private long debugInterval; + private ScheduledFuture<?> debugFuture; + private int failures; + private long ticks; + + @Override + public void configure( + CamelContext camelContext, CliActionDispatcher dispatcher, CliSnapshotProducer snapshots, Runnable shutdown) { + this.camelContext = camelContext; + this.dispatcher = dispatcher; + this.snapshots = snapshots; + this.shutdown = shutdown; + } + + @Override + protected void doStart() throws Exception { + if ("prod".equals(camelContext.getCamelContextExtension().getProfile())) { + throw new IllegalStateException( + "The Camel CLI connector websocket transport gives the connected tool full control of this" + + " application and cannot be used with the prod profile." + + " Remove camel.cli.transport=websocket."); + } + String value = property("camel.cli.websocket.url", null); + if (value == null || value.isBlank()) { + throw new IllegalArgumentException("camel.cli.websocket.url must be set when camel.cli.transport=websocket"); + } + url = URI.create(value); + String scheme = url.getScheme() != null ? url.getScheme().toLowerCase(Locale.ROOT) : ""; + if (!"ws".equals(scheme) && !"wss".equals(scheme)) { + throw new IllegalArgumentException("camel.cli.websocket.url must be a ws:// or wss:// url: " + value); + } + boolean loopback = isLoopback(url.getHost()); + token = property("camel.cli.websocket.token", null); + if ((token == null || token.isBlank()) && !loopback) { + throw new IllegalArgumentException( + "camel.cli.websocket.token must be set when camel.cli.websocket.url is not a loopback address"); + } + reconnectDelay = Long.parseLong(property("camel.cli.websocket.reconnect-delay", "1000")); + reconnectMaxDelay = Long.parseLong(property("camel.cli.websocket.reconnect-max-delay", "30000")); + snapshotInterval = Long.parseLong(property("camel.cli.websocket.snapshot-interval", "1000")); + // when debugging, the debug snapshot (breakpoints) is sent faster than the others, as the file transport does + debugInterval = camelContext.isDebugging() ? 100 : 0; + heartbeatInterval = Long.parseLong(property("camel.cli.websocket.heartbeat-interval", "10000")); + + LOG.warn("Camel CLI connector connects to {} which gets full control of this application (development use only)", + where()); + if ("ws".equals(scheme) && !loopback) { + LOG.warn("Camel CLI connector uses an unencrypted connection to a remote host; use wss:// or a tunnel"); + } + + client = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(10)).build(); + // Camel's thread factory (naming, virtual threads when enabled) but not Camel's thread pools: these threads must + // keep running while Camel is stopping, to report it and to send the close frame + threadNamePattern = camelContext.getExecutorServiceManager().getThreadNamePattern(); + actions = new ThreadPoolExecutor( + 1, 1, 0, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>(MAX_PENDING_ACTIONS), + new CamelThreadFactory(threadNamePattern, "CliConnectorActions", true)); + scheduler = Executors.newSingleThreadScheduledExecutor( + new CamelThreadFactory(threadNamePattern, "CliConnectorWebSocket", true)); + + stopping = false; + ready = camelContext.isStarted(); + if (!listenerAdded) { + // Camel cannot remove a startup listener, so it is added once even if this transport is restarted + listenerAdded = true; + camelContext.addStartupListener(new ExtendedStartupListener() { + @Override + public void onCamelContextStarted(CamelContext context, boolean alreadyStarted) { + // wait until fully started + } + + @Override + public void onCamelContextFullyStarted(CamelContext context, boolean alreadyStarted) { + ready = true; + execute(() -> sayHello(connection)); + } + }); + } + + // a periodic task that throws is never run again, hence safely() + scheduleSnapshots(); + scheduler.scheduleWithFixedDelay(() -> safely(this::heartbeatTask), heartbeatInterval, heartbeatInterval, + TimeUnit.MILLISECONDS); + scheduler.execute(this::connect); + } + + @Override + public void updateDelay(int delay) { + // camel-cli-debug makes it faster so breakpoints show up quickly: only the debug snapshot needs that + debugInterval = delay; + execute(this::scheduleSnapshots); + } + + private void scheduleSnapshots() { + if (snapshotFuture != null) { + snapshotFuture.cancel(false); + } + snapshotFuture = scheduler.scheduleWithFixedDelay(() -> safely(this::snapshotTask), snapshotInterval, + snapshotInterval, TimeUnit.MILLISECONDS); + if (debugFuture != null) { + debugFuture.cancel(false); + debugFuture = null; + } + if (debugInterval > 0) { + debugFuture = scheduler.scheduleWithFixedDelay(() -> safely(this::debugTask), debugInterval, + debugInterval, TimeUnit.MILLISECONDS); + } + } + + private void debugTask() { + Connection c = connection; + if (c != null && c.helloSent) { + // only sent when it changed, so the regular snapshot task does not send it again + snapshot(c, "debug", true); + } + } + + @Override + protected void doStop() throws Exception { + stopping = true; + if (scheduler != null) { + try { + scheduler.submit(() -> { + Connection c = connection; + if (c != null) { + connection = null; + c.close(WebSocket.NORMAL_CLOSURE, "stopping"); + } + }).get(SEND_TIMEOUT, TimeUnit.MILLISECONDS); + } catch (Exception e) { + LOG.debug("Error closing websocket due to: {}. This exception is ignored.", e.getMessage(), e); + } + scheduler.shutdownNow(); + scheduler = null; + } + if (actions != null) { + actions.shutdownNow(); + actions = null; + } + client = null; + } + + // ---- connection lifecycle (scheduler thread) ---- + + private void connect() { + if (stopping) { + return; + } + try { + WebSocket.Builder builder = client.newWebSocketBuilder().connectTimeout(Duration.ofSeconds(10)); + if (token != null && !token.isBlank()) { + builder.header("Authorization", "Bearer " + token); + } + builder.buildAsync(url, new Connection()).whenComplete((ws, e) -> { + if (e != null) { + execute(() -> connectFailed(e)); + } + }); + } catch (Exception e) { + connectFailed(e); + } + } + + private void connectFailed(Throwable e) { + failures++; + Throwable cause = e.getCause() != null ? e.getCause() : e; + if (cause instanceof WebSocketHandshakeException he) { + int code = he.getResponse().statusCode(); + if (code == 401 || code == 403) { + LOG.warn("Camel CLI connector was rejected by {} (HTTP {}): check camel.cli.websocket.token", where(), code); + } else { + LOG.warn("Camel CLI connector was rejected by {} (HTTP {})", where(), code); + } + } else if (failures == 1) { + LOG.warn("Camel CLI connector cannot connect to {} due to: {}. Will keep retrying.", where(), describe(cause)); + } else { + LOG.debug("Camel CLI connector cannot connect to {} due to: {}", where(), describe(cause)); + } + reconnect(); + } + + private void opened(Connection c) { + if (stopping) { + c.ws.abort(); + return; + } + connection = c; + LOG.info("Camel CLI connector connected to {}", where()); + sayHello(c); + } + + private void closed(Connection c, String reason) { + if (connection != c) { + // an old connection, or already handled + return; + } + connection = null; + c.ws.abort(); + // a connection that drops right away counts as a failure, so a tool that keeps closing us is not hammered + if (System.currentTimeMillis() - c.openedAt < STABLE_CONNECTION) { + failures++; + } else { + failures = 0; + } + if (!stopping) { + LOG.info("Camel CLI connector disconnected from {} ({})", where(), reason); + } + reconnect(); + } + + private void reconnect() { + if (stopping) { + return; + } + long delay = reconnectDelay; + for (int i = 1; i < failures && delay < reconnectMaxDelay; i++) { + delay *= 2; + } + delay = Math.min(delay, reconnectMaxDelay); + // jitter, so many applications restarted together do not reconnect in lockstep + delay += ThreadLocalRandom.current().nextLong(delay / 5 + 1); + scheduler.schedule(this::connect, delay, TimeUnit.MILLISECONDS); + } + + private void heartbeatTask() { + Connection c = connection; + if (c == null) { + return; + } + if (System.currentTimeMillis() - c.lastSeen > 3 * heartbeatInterval) { + // the tool went away without closing the connection (sleep, network change, crash) + closed(c, "no heartbeat"); + } else { + c.ws.sendPing(ByteBuffer.allocate(0)); + } + } + + // ---- outgoing frames (scheduler thread) ---- + + private void sayHello(Connection c) { + if (c == null || c.helloSent || !ready || c != connection) { + return; + } + JsonObject frame = envelope("hello"); + frame.put("camelVersion", camelContext.getVersion()); + frame.put("name", camelContext.getName()); + frame.put("transport", "jdk"); + try { + JsonObject status = snapshots.status(); + frame.put("runtime", status.get("runtime")); + c.helloSent = true; + send(c, frame); + sendSnapshot(c, "status", status); + } catch (Exception e) { + LOG.debug("Error sending hello due to: {}. Will retry.", e.getMessage(), e); + } + } + + private void snapshotTask() { + Connection c = connection; + if (c == null) { + return; + } + if (!c.helloSent) { + // not started yet, or the hello failed + sayHello(c); + return; + } + snapshot(c, "status", false); + // as the file transport: these have more overhead and are only needed when tracing/debugging/receiving + if (++ticks % 2 == 0) { + snapshot(c, "trace", false); + snapshot(c, "receive", false); + snapshot(c, "debug", true); + snapshot(c, "history", true); + snapshot(c, "error", true); + snapshot(c, "activity", true); + } + } + + private void snapshot(Connection c, String kind, boolean onlyIfChanged) { + try { + JsonObject data = switch (kind) { + case "status" -> snapshots.status(); + case "trace" -> c.newMessages(snapshots.trace(), "traces"); + case "receive" -> c.newMessages(snapshots.receive(), "messages"); + case "debug" -> snapshots.debug(); + case "history" -> snapshots.messageHistory(); + case "error" -> snapshots.errors(); + default -> snapshots.activity(); + }; + if (data == null || data.isEmpty()) { + return; + } + if (onlyIfChanged && !c.changed(kind, data)) { + return; + } + if ("trace".equals(kind)) { + sendInBatches(c, kind, data, "traces"); + } else if ("receive".equals(kind)) { + sendInBatches(c, kind, data, "messages"); + } else { + sendSnapshot(c, kind, data); + } + } catch (Exception e) { + LOG.trace("Error sending {} snapshot due to: {}. This exception is ignored.", kind, e.getMessage(), e); + } + } + + /** + * Sends the messages in as many snapshots as needed to keep each under {@link #MAX_SNAPSHOT_SIZE} (a single message + * larger than that is sent on its own). + */ + private void sendInBatches(Connection c, String kind, JsonObject data, String key) { + List<JsonObject> messages = data.getCollection(key); + List<JsonObject> batch = new ArrayList<>(); + int size = 0; + for (JsonObject m : messages) { + // servers limit the size in bytes: non-ASCII characters are not escaped and take up to 3 bytes + int length = m.toJson().getBytes(StandardCharsets.UTF_8).length; + if (!batch.isEmpty() && size + length > MAX_SNAPSHOT_SIZE) { Review Comment: A possible edge case: a single message larger than `MAX_SNAPSHOT_SIZE` is sent in a frame of its own. If the tool's server limit (256 KB in Vert.x/Quarkus) closes the connection on it, the next connection starts again from uid 0 (`lastTraceUid` / `lastReceiveUid` are kept per `Connection`), sends the same message again, and is closed again, with backoff. `BODY_MAX_CHARS` makes this unlikely, but large headers could still hit it. Two possible fixes: drop or truncate a single oversized message (with a marker), or keep the uid cursor in the transport rather than per connection. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
