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

Croway pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new 6cc52b80a5dd CAMEL-25197: camel-cli-connector - do not let a slow 
status starve the WebSocket heartbeat (#27279)
6cc52b80a5dd is described below

commit 6cc52b80a5dd15e8f105bd74d9b4b583f530342c
Author: Federico Mariani <[email protected]>
AuthorDate: Fri Oct 2 15:27:51 2026 +0200

    CAMEL-25197: camel-cli-connector - do not let a slow status starve the 
WebSocket heartbeat (#27279)
    
    * CAMEL-25197: camel-cli-connector - do not let a slow status starve the 
WebSocket heartbeat
    
    The status of a large integration can take a long time to collect: with 316
    routes (about 4000 processors) the route dev console took about 30 s, as it
    creates a JMX proxy per processor. The WebSocket transport collected it on 
the
    thread that also sends heartbeats, so the connection was declared dead ("no
    heartbeat") and dropped about every 65 s, and hello (which collected a full
    status for its runtime section) took 30 s after each reconnect.
    
    - snapshots are collected on their own thread and handed to the scheduler,
      which alone sends (heartbeats and results are never held up)
    - hello uses the new CliSnapshotProducer.runtime(), cheap to collect
    - the next snapshot round waits at least as long as the previous one took,
      so a slow status does not keep a CPU busy
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
    
    * CAMEL-25197: camel-cli-connector - review: regression test, snapshot 
backpressure, thread comments
    
    - test: a status that takes 2 s holds up neither hello nor the heartbeats 
(fails
      without the previous commit: hello after 2.2 s, reconnect after about 6 s)
    - the next snapshot round is scheduled once the scheduler has sent the 
frames of
      the current one, so a slow link does not queue frames without end; a round
      counter keeps updateDelay from starting a second chain
    - comments describe the three threads and which fields the collector owns; 
note
      that the debug snapshot shares the collector with the status
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
    
    ---------
    
    Co-authored-by: Claude Opus 5.5 <[email protected]>
---
 .../camel/cli/connector/CliSnapshotProducer.java   |  9 +++
 .../camel/cli/connector/LocalCliConnector.java     | 57 ++++++++--------
 .../connector/WebSocketCliConnectorTransport.java  | 78 +++++++++++++++++-----
 .../WebSocketCliConnectorTransportTest.java        | 23 +++++++
 4 files changed, 125 insertions(+), 42 deletions(-)

diff --git 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliSnapshotProducer.java
 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliSnapshotProducer.java
index 673d9becc3ac..cbdb6b0b1a6c 100644
--- 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliSnapshotProducer.java
+++ 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/CliSnapshotProducer.java
@@ -36,6 +36,15 @@ public interface CliSnapshotProducer {
      */
     JsonObject status() throws Exception;
 
+    /**
+     * The runtime section of the status (pid, platform, Java), cheap to 
collect.
+     *
+     * @since 4.23
+     */
+    default JsonObject runtime() throws Exception {
+        return status().getMap("runtime");
+    }
+
     /**
      * Traced messages (<tt>traces</tt> array, each with a <tt>uid</tt>).
      */
diff --git 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/LocalCliConnector.java
 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/LocalCliConnector.java
index ec9645068cb3..85a069e12400 100644
--- 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/LocalCliConnector.java
+++ 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/LocalCliConnector.java
@@ -1365,33 +1365,7 @@ public class LocalCliConnector extends ServiceSupport
     @Override
     public JsonObject status() throws Exception {
         JsonObject root = new JsonObject();
-
-        // what runtime are in use
-        JsonObject rc = new JsonObject();
-        String dir = new File(".").getAbsolutePath();
-        dir = FileUtil.onlyPath(dir);
-        rc.put("pid", ProcessHandle.current().pid());
-        rc.put("directory", dir);
-        ProcessHandle.current().info().user().ifPresent(u -> rc.put("user", 
u));
-        rc.put("platform", platform);
-        if (platformVersion != null) {
-            rc.put("platformVersion", platformVersion);
-        }
-        if (mainClass != null) {
-            rc.put("mainClass", mainClass);
-        }
-        RuntimeMXBean mb = ManagementFactory.getRuntimeMXBean();
-        if (mb != null) {
-            rc.put("javaVersion", mb.getVmVersion());
-            rc.put("javaVendor", mb.getVmVendor());
-            rc.put("javaVmName", mb.getVmName());
-        }
-        String readmeFiles = camelContext.getPropertiesComponent()
-                .resolveProperty("camel.jbang.readmeFiles").orElse(null);
-        if (readmeFiles != null) {
-            rc.put("readmeFiles", readmeFiles);
-        }
-        root.put("runtime", rc);
+        root.put("runtime", runtime());
 
         DevConsoleRegistry dcr = 
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class);
         if (dcr != null) {
@@ -1654,6 +1628,35 @@ public class LocalCliConnector extends ServiceSupport
         return root;
     }
 
+    @Override
+    public JsonObject runtime() throws Exception {
+        JsonObject rc = new JsonObject();
+        String dir = new File(".").getAbsolutePath();
+        dir = FileUtil.onlyPath(dir);
+        rc.put("pid", ProcessHandle.current().pid());
+        rc.put("directory", dir);
+        ProcessHandle.current().info().user().ifPresent(u -> rc.put("user", 
u));
+        rc.put("platform", platform);
+        if (platformVersion != null) {
+            rc.put("platformVersion", platformVersion);
+        }
+        if (mainClass != null) {
+            rc.put("mainClass", mainClass);
+        }
+        RuntimeMXBean mb = ManagementFactory.getRuntimeMXBean();
+        if (mb != null) {
+            rc.put("javaVersion", mb.getVmVersion());
+            rc.put("javaVendor", mb.getVmVendor());
+            rc.put("javaVmName", mb.getVmName());
+        }
+        String readmeFiles = camelContext.getPropertiesComponent()
+                .resolveProperty("camel.jbang.readmeFiles").orElse(null);
+        if (readmeFiles != null) {
+            rc.put("readmeFiles", readmeFiles);
+        }
+        return rc;
+    }
+
     @Override
     public JsonObject trace() throws Exception {
         return callConsole("trace", Map.of("dump", "true"));
diff --git 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
index 7ce0700d54d7..b806fa38826c 100644
--- 
a/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
+++ 
b/dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java
@@ -61,8 +61,9 @@ import org.slf4j.LoggerFactory;
  * client (<tt>camel.cli.websocket.client=jdk</tt> always uses the JDK client).
  * <p/>
  * Threads: actions run one at a time on their own thread (the dispatcher is 
not thread-safe, and an action can block
- * for a long time); connecting, parsing the incoming frames, snapshots, 
heartbeats and every write to the socket run on
- * a second thread, so the client never has two sends in flight, and its 
callbacks never wait.
+ * for a long time); snapshots are collected on a second thread (the status of 
a large integration can take seconds);
+ * connecting, parsing the incoming frames, heartbeats and every write to the 
socket run on a third thread (the
+ * scheduler), so the client never has two sends in flight, and its callbacks 
never wait.
  */
 public class WebSocketCliConnectorTransport extends ServiceSupport implements 
CliConnectorTransport {
 
@@ -94,17 +95,22 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
     private CliWebSocketClient client;
     private ThreadPoolExecutor actions;
     private ScheduledExecutorService scheduler;
-    // fields below are only used from the scheduler thread, except the 
volatile ones
+    // collects snapshots: the status of a large integration can take seconds, 
which must not hold up the scheduler
+    // (heartbeats, sends); snapshots are handed to the scheduler to be sent
+    private ScheduledExecutorService collector;
+    // fields below are only used from the scheduler thread, except the 
volatile ones and the snapshot ones below
     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 int failures;
+    // snapshot fields, only used on the collector thread (or before it runs 
anything)
     private ScheduledFuture<?> snapshotFuture;
     private long debugInterval;
     private ScheduledFuture<?> debugFuture;
-    private int failures;
     private long ticks;
+    // bumped when the rounds are rescheduled, so a round that was already 
running does not start a second chain
+    private long rounds;
 
     @Override
     public void configure(
@@ -166,6 +172,8 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
                 new CamelThreadFactory(threadNamePattern, 
"CliConnectorActions", true));
         scheduler = Executors.newSingleThreadScheduledExecutor(
                 new CamelThreadFactory(threadNamePattern, 
"CliConnectorWebSocket", true));
+        collector = Executors.newSingleThreadScheduledExecutor(
+                new CamelThreadFactory(threadNamePattern, 
"CliConnectorSnapshots", true));
 
         stopping = false;
         ready = camelContext.isStarted();
@@ -215,26 +223,60 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
     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);
+        ScheduledExecutorService c = collector;
+        if (c != null) {
+            // on the collector, which also reschedules the snapshot rounds
+            c.execute(this::scheduleSnapshots);
+        }
     }
 
     private void scheduleSnapshots() {
         if (snapshotFuture != null) {
             snapshotFuture.cancel(false);
         }
-        snapshotFuture = scheduler.scheduleWithFixedDelay(() -> 
safely(this::snapshotTask), snapshotInterval,
-                snapshotInterval, TimeUnit.MILLISECONDS);
+        long round = ++rounds;
+        snapshotFuture = collector.schedule(() -> snapshotRound(round), 
snapshotInterval, TimeUnit.MILLISECONDS);
         if (debugFuture != null) {
             debugFuture.cancel(false);
             debugFuture = null;
         }
         if (debugInterval > 0) {
-            debugFuture = scheduler.scheduleWithFixedDelay(() -> 
safely(this::debugTask), debugInterval,
+            debugFuture = collector.scheduleWithFixedDelay(() -> 
safely(this::debugTask), debugInterval,
                     debugInterval, TimeUnit.MILLISECONDS);
         }
     }
 
+    /**
+     * One round of snapshots, then the next one: once the scheduler has sent 
the frames of this round (a slow link does
+     * not queue frames without end), at least the snapshot interval later, 
and at least as long as collecting took (a
+     * slow status does not keep a CPU busy).
+     */
+    private void snapshotRound(long round) {
+        if (round != rounds) {
+            return;
+        }
+        long start = System.nanoTime();
+        safely(this::snapshotTask);
+        long delay = Math.max(snapshotInterval, 
TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start));
+        // queued behind the frames of this round
+        execute(() -> {
+            ScheduledExecutorService c = collector;
+            if (c != null) {
+                try {
+                    c.execute(() -> {
+                        if (round == rounds) {
+                            snapshotFuture = c.schedule(() -> 
snapshotRound(round), delay, TimeUnit.MILLISECONDS);
+                        }
+                    });
+                } catch (RejectedExecutionException e) {
+                    // stopping
+                }
+            }
+        });
+    }
+
     private void debugTask() {
+        // note: shares the collector with the status, so a slow status delays 
it (as when both ran on the scheduler)
         Connection c = connection;
         if (c != null && c.helloSent) {
             // only sent when it changed, so the regular snapshot task does 
not send it again
@@ -264,6 +306,10 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
             actions.shutdownNow();
             actions = null;
         }
+        if (collector != null) {
+            collector.shutdownNow();
+            collector = null;
+        }
         client = null;
     }
 
@@ -391,16 +437,16 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
         frame.put("name", camelContext.getName());
         frame.put("transport", client.getName());
         try {
-            JsonObject status = snapshots.status();
-            frame.put("runtime", status.get("runtime"));
+            frame.put("runtime", snapshots.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);
         }
     }
 
+    // ---- snapshots (collector thread) ----
+
     private void snapshotTask() {
         Connection c = connection;
         if (c == null) {
@@ -408,7 +454,7 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
         }
         if (!c.helloSent) {
             // not started yet, or the hello failed
-            sayHello(c);
+            execute(() -> sayHello(c));
             return;
         }
         snapshot(c, "status", false);
@@ -491,7 +537,8 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
         JsonObject frame = envelope("snapshot");
         frame.put("kind", kind);
         frame.put("data", data);
-        send(c, frame);
+        // only the scheduler sends, one frame at a time
+        execute(() -> send(c, frame));
     }
 
     private void send(Connection c, JsonObject frame) {
@@ -715,7 +762,8 @@ public class WebSocketCliConnectorTransport extends 
ServiceSupport implements Cl
         private boolean lost;
         private long openedAt = System.currentTimeMillis();
         private volatile long lastSeen = openedAt;
-        private boolean helloSent;
+        private volatile boolean helloSent;
+        // what was sent on this connection, only used on the collector thread
         private long lastTraceUid;
         private long lastReceiveUid;
         private final Map<String, String> lastSent = new HashMap<>();
diff --git 
a/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
 
b/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
index f285369027a0..3115c8bc2a5d 100644
--- 
a/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
+++ 
b/dsl/camel-cli-connector/src/test/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransportTest.java
@@ -257,6 +257,29 @@ class WebSocketCliConnectorTransportTest extends 
CamelTestSupport {
         tool.awaitFrame(f -> "hello".equals(f.getString("type")));
     }
 
+    @Test
+    void keepsTheConnectionWhileASlowStatusIsCollected() throws Exception {
+        // the status of a large integration (thousands of processors) can 
take seconds to collect: it must hold up
+        // neither hello nor the heartbeats
+        property("camel.cli.websocket.heartbeat-interval", "200");
+        connector = new LocalCliConnector(new DefaultCliConnectorFactory()) {
+            @Override
+            public JsonObject status() throws Exception {
+                Thread.sleep(2000);
+                return super.status();
+            }
+        };
+        connector.setCamelContext(context);
+        long start = System.currentTimeMillis();
+        connector.start();
+
+        tool.awaitFrame(f -> "hello".equals(f.getString("type")));
+        assertThat(System.currentTimeMillis() - start).isLessThan(1500);
+        // longer than 3 heartbeat intervals behind several status 
collections: still the first connection
+        await().during(6, TimeUnit.SECONDS).atMost(8, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertThat(tool.handshakes).hasValue(1));
+    }
+
     @Test
     void reconnectsWhenTheToolStopsAnswering() throws Exception {
         property("camel.cli.websocket.heartbeatInterval", "200");

Reply via email to