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");