Croway commented on code in PR #27145:
URL: https://github.com/apache/camel/pull/27145#discussion_r4153652318
##########
dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/LocalCliConnector.java:
##########
@@ -1457,441 +1349,336 @@ private void doActionProcessorTask(JsonObject root) {
}
}
- JsonObject loadAction() {
- return loadAction(actionFile);
- }
+ @Override
+ public JsonObject status() throws Exception {
+ JsonObject root = new JsonObject();
- JsonObject loadAction(File file) {
- try {
- if (file != null && file.exists()) {
- FileInputStream fis = new FileInputStream(file);
- String text = IOHelper.loadText(fis);
- IOHelper.close(fis);
- if (!text.isEmpty()) {
- return (JsonObject) Jsoner.deserialize(text);
- }
- }
- } catch (Exception e) {
- // ignore
+ // 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());
}
- return null;
- }
-
- protected void statusTask() {
- try {
- // even during termination then collect status as we want to see
status changes during stopping
- 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);
-
- DevConsoleRegistry dcr =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class);
- if (dcr != null) {
- // collect details via console
- DevConsole dc = dcr.resolveById("context");
- DevConsole dc2 = dcr.resolveById("route");
- if (dc != null && dc2 != null) {
- JsonObject json = (JsonObject)
dc.call(DevConsole.MediaType.JSON);
- JsonObject json2 = (JsonObject)
dc2.call(DevConsole.MediaType.JSON, Map.of("processors", "true"));
- if (json != null && json2 != null) {
- root.put("context", json);
- root.put("routes", json2.get("routes"));
- }
- }
- DevConsole dc2b = dcr.resolveById("route-group");
- if (dc2b != null) {
- JsonObject json = (JsonObject)
dc2b.call(DevConsole.MediaType.JSON);
- if (json != null) {
- root.put("routeGroups", json.get("routeGroups"));
- }
- }
- DevConsole dc3 = dcr.resolveById("endpoint");
- if (dc3 != null) {
- JsonObject json = (JsonObject)
dc3.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("endpoints", json);
- }
- }
- DevConsole dc4 = dcr.resolveById("health");
- if (dc4 != null) {
- // include full details in health checks
- JsonObject json = (JsonObject)
dc4.call(DevConsole.MediaType.JSON, Map.of("exposureLevel", "full"));
- if (json != null && !json.isEmpty()) {
- root.put("healthChecks", json);
- }
- }
- DevConsole dc5 = dcr.resolveById("event");
- if (dc5 != null) {
- JsonObject json = (JsonObject)
dc5.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("events", json);
- }
- }
- DevConsole dc6 = dcr.resolveById("log");
- if (dc6 != null) {
- JsonObject json = (JsonObject)
dc6.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("logger", json);
- }
- }
- DevConsole dc7 = dcr.resolveById("inflight");
- if (dc7 != null) {
- JsonObject json = (JsonObject)
dc7.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("inflight", json);
- }
- }
- DevConsole dc8 = dcr.resolveById("blocked");
- if (dc8 != null) {
- JsonObject json = (JsonObject)
dc8.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("blocked", json);
- }
- }
- DevConsole dc9 = dcr.resolveById("micrometer");
- if (dc9 != null) {
- JsonObject json = (JsonObject)
dc9.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("micrometer", json);
- }
- }
- DevConsole dc10 = dcr.resolveById("resilience4j");
- if (dc10 != null) {
- JsonObject json = (JsonObject)
dc10.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("resilience4j", json);
- }
- }
- DevConsole dc11 = dcr.resolveById("fault-tolerance");
- if (dc11 != null) {
- JsonObject json = (JsonObject)
dc11.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("fault-tolerance", json);
- }
- }
- DevConsole dc12 = dcr.resolveById("circuit-breaker");
- if (dc12 != null) {
- JsonObject json = (JsonObject)
dc12.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("circuit-breaker", json);
- }
- }
- DevConsole dc13 = dcr.resolveById("trace");
- if (dc13 != null) {
- JsonObject json = (JsonObject)
dc13.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("trace", json);
- }
- }
- DevConsole dc14 = dcr.resolveById("consumer");
- if (dc14 != null) {
- JsonObject json = (JsonObject)
dc14.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("consumers", json);
- }
+ String readmeFiles = camelContext.getPropertiesComponent()
+ .resolveProperty("camel.jbang.readmeFiles").orElse(null);
+ if (readmeFiles != null) {
+ rc.put("readmeFiles", readmeFiles);
+ }
+ root.put("runtime", rc);
+
+ DevConsoleRegistry dcr =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class);
+ if (dcr != null) {
+ // collect details via console
+ DevConsole dc = dcr.resolveById("context");
+ DevConsole dc2 = dcr.resolveById("route");
+ if (dc != null && dc2 != null) {
+ JsonObject json = (JsonObject)
dc.call(DevConsole.MediaType.JSON);
+ JsonObject json2 = (JsonObject)
dc2.call(DevConsole.MediaType.JSON, Map.of("processors", "true"));
+ if (json != null && json2 != null) {
+ root.put("context", json);
+ root.put("routes", json2.get("routes"));
}
- DevConsole dc14b = dcr.resolveById("producer");
- if (dc14b != null) {
- JsonObject json = (JsonObject)
dc14b.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("producers", json);
- }
+ }
+ DevConsole dc2b = dcr.resolveById("route-group");
+ if (dc2b != null) {
+ JsonObject json = (JsonObject)
dc2b.call(DevConsole.MediaType.JSON);
+ if (json != null) {
+ root.put("routeGroups", json.get("routeGroups"));
}
- DevConsole dc14c = dcr.resolveById("route-controller");
- if (dc14c != null) {
- JsonObject json
- = (JsonObject)
dc14c.call(DevConsole.MediaType.JSON, Map.of("stacktrace", "false"));
- if (json != null && !json.isEmpty()) {
- root.put("routeController", json);
- }
+ }
+ DevConsole dc3 = dcr.resolveById("endpoint");
+ if (dc3 != null) {
+ JsonObject json = (JsonObject)
dc3.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("endpoints", json);
}
- DevConsole dc15 = dcr.resolveById("variables");
- if (dc15 != null) {
- JsonObject json = (JsonObject)
dc15.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("variables", json);
- }
+ }
+ DevConsole dc4 = dcr.resolveById("health");
+ if (dc4 != null) {
+ // include full details in health checks
+ JsonObject json = (JsonObject)
dc4.call(DevConsole.MediaType.JSON, Map.of("exposureLevel", "full"));
+ if (json != null && !json.isEmpty()) {
+ root.put("healthChecks", json);
}
- DevConsole dc16 = dcr.resolveById("transformers");
- if (dc16 != null) {
- JsonObject json = (JsonObject)
dc16.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("transformers", json);
- }
+ }
+ DevConsole dc5 = dcr.resolveById("event");
+ if (dc5 != null) {
+ JsonObject json = (JsonObject)
dc5.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("events", json);
}
- DevConsole dc17 = dcr.resolveById("service");
- if (dc17 != null) {
- JsonObject json = (JsonObject)
dc17.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("services", json);
- }
+ }
+ DevConsole dc6 = dcr.resolveById("log");
+ if (dc6 != null) {
+ JsonObject json = (JsonObject)
dc6.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("logger", json);
}
- DevConsole dc18 = dcr.resolveById("platform-http");
- if (dc18 != null) {
- JsonObject json = (JsonObject)
dc18.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("platform-http", json);
- }
+ }
+ DevConsole dc7 = dcr.resolveById("inflight");
+ if (dc7 != null) {
+ JsonObject json = (JsonObject)
dc7.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("inflight", json);
}
- DevConsole dc19 = dcr.resolveById("rest");
- if (dc19 != null) {
- JsonObject json = (JsonObject)
dc19.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("rests", json);
- }
+ }
+ DevConsole dc8 = dcr.resolveById("blocked");
+ if (dc8 != null) {
+ JsonObject json = (JsonObject)
dc8.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("blocked", json);
}
- DevConsole dc20 = dcr.resolveById("kafka");
- if (dc20 != null) {
- JsonObject json = (JsonObject)
dc20.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("kafka", json);
- }
+ }
+ DevConsole dc9 = dcr.resolveById("micrometer");
+ if (dc9 != null) {
+ JsonObject json = (JsonObject)
dc9.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("micrometer", json);
}
- DevConsole dc21 = dcr.resolveById("properties");
- if (dc21 != null) {
- JsonObject json = (JsonObject)
dc21.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("properties", json);
- }
+ }
+ DevConsole dc10 = dcr.resolveById("resilience4j");
+ if (dc10 != null) {
+ JsonObject json = (JsonObject)
dc10.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("resilience4j", json);
}
- DevConsole dc22 = dcr.resolveById("main-configuration");
- if (dc22 != null) {
- JsonObject json = (JsonObject)
dc22.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("main-configuration", json);
- }
+ }
+ DevConsole dc11 = dcr.resolveById("fault-tolerance");
+ if (dc11 != null) {
+ JsonObject json = (JsonObject)
dc11.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("fault-tolerance", json);
}
- DevConsole dc23 = dcr.resolveById("receive");
- if (dc23 != null) {
- JsonObject json = (JsonObject)
dc23.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("receive", json);
- }
+ }
+ DevConsole dc12 = dcr.resolveById("circuit-breaker");
+ if (dc12 != null) {
+ JsonObject json = (JsonObject)
dc12.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("circuit-breaker", json);
}
- DevConsole dc24 = dcr.resolveById("internal-tasks");
- if (dc24 != null) {
- JsonObject json = (JsonObject)
dc24.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("internal-tasks", json);
- }
+ }
+ DevConsole dc13 = dcr.resolveById("trace");
+ if (dc13 != null) {
+ JsonObject json = (JsonObject)
dc13.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("trace", json);
}
- DevConsole dc25 = dcr.resolveById("groovy");
- if (dc25 != null) {
- JsonObject json = (JsonObject)
dc25.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("groovy", json);
- }
+ }
+ DevConsole dc14 = dcr.resolveById("consumer");
+ if (dc14 != null) {
+ JsonObject json = (JsonObject)
dc14.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("consumers", json);
}
- DevConsole dc26 = dcr.resolveById("errors");
- if (dc26 != null) {
- JsonObject json = (JsonObject)
dc26.call(DevConsole.MediaType.JSON,
- Map.of("stackTrace", "true"));
- if (json != null && !json.isEmpty()) {
- // only include metadata in status file (full error
data is in the error file)
- JsonObject summary = new JsonObject();
- summary.put("enabled", json.get("enabled"));
- summary.put("size", json.get("size"));
- summary.put("maximumEntries",
json.get("maximumEntries"));
- summary.put("timeToLive", json.get("timeToLive"));
- root.put("errors", summary);
- }
+ }
+ DevConsole dc14b = dcr.resolveById("producer");
+ if (dc14b != null) {
+ JsonObject json = (JsonObject)
dc14b.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("producers", json);
}
- DevConsole dc27 = dcr.resolveById("datasource");
- if (dc27 != null) {
- JsonObject json = (JsonObject)
dc27.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("dataSources", json);
- }
+ }
+ DevConsole dc14c = dcr.resolveById("route-controller");
+ if (dc14c != null) {
+ JsonObject json
+ = (JsonObject) dc14c.call(DevConsole.MediaType.JSON,
Map.of("stacktrace", "false"));
+ if (json != null && !json.isEmpty()) {
+ root.put("routeController", json);
}
- DevConsole dc28 = dcr.resolveById("sql-trace");
- if (dc28 != null) {
- JsonObject json = (JsonObject)
dc28.call(DevConsole.MediaType.JSON);
- if (json != null && !json.isEmpty()) {
- root.put("sqlTrace", json);
- }
+ }
+ DevConsole dc15 = dcr.resolveById("variables");
+ if (dc15 != null) {
+ JsonObject json = (JsonObject)
dc15.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("variables", json);
}
- JsonArray consoleIds = new JsonArray();
- consoleIds.addAll(dcr.getConsoleIDs());
- root.put("devConsoles", consoleIds);
- }
- boolean hasBrowseable = false;
- for (Endpoint ep : getCamelContext().getEndpoints()) {
- if (ep instanceof BrowsableEndpoint) {
- hasBrowseable = true;
- break;
+ }
+ DevConsole dc16 = dcr.resolveById("transformers");
+ if (dc16 != null) {
+ JsonObject json = (JsonObject)
dc16.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("transformers", json);
}
}
- root.put("hasBrowseableEndpoints", hasBrowseable);
- // various details
- JsonObject mem = collectMemory();
- if (mem != null) {
- root.put("memory", mem);
+ DevConsole dc17 = dcr.resolveById("service");
+ if (dc17 != null) {
+ JsonObject json = (JsonObject)
dc17.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("services", json);
+ }
}
- JsonObject cl = collectClassLoading();
- if (cl != null) {
- root.put("classLoading", cl);
+ DevConsole dc18 = dcr.resolveById("platform-http");
+ if (dc18 != null) {
+ JsonObject json = (JsonObject)
dc18.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("platform-http", json);
+ }
}
- JsonObject threads = collectThreads();
- if (threads != null) {
- root.put("threads", threads);
+ DevConsole dc19 = dcr.resolveById("rest");
+ if (dc19 != null) {
+ JsonObject json = (JsonObject)
dc19.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("rests", json);
+ }
}
- JsonObject gc = collectGC();
- if (gc != null) {
- root.put("gc", gc);
+ DevConsole dc20 = dcr.resolveById("kafka");
+ if (dc20 != null) {
+ JsonObject json = (JsonObject)
dc20.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("kafka", json);
+ }
}
- JsonObject vaults = collectVaults();
- if (!vaults.isEmpty()) {
- root.put("vaults", vaults);
+ DevConsole dc21 = dcr.resolveById("properties");
+ if (dc21 != null) {
+ JsonObject json = (JsonObject)
dc21.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("properties", json);
+ }
}
- LOG.trace("Updating status file: {}", statusFile);
- IOHelper.writeText(root.toJson(), statusFile);
- } catch (Exception e) {
- // ignore
- LOG.trace("Error updating status file: {} due to: {}. This
exception is ignored.",
- statusFile, e.getMessage(), e);
- }
- }
-
- protected void traceTask() {
- try {
- DevConsole dc12 =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class)
- .resolveById("trace");
- if (dc12 != null) {
- JsonObject json = (JsonObject)
dc12.call(DevConsole.MediaType.JSON, Map.of("dump", "true"));
- JsonArray arr = json.getCollection("traces");
- // filter based on last uid
- if (traceFilePos > 0) {
- arr.removeIf(r -> {
- JsonObject jo = (JsonObject) r;
- return jo.getLong("uid") <= traceFilePos;
- });
+ DevConsole dc22 = dcr.resolveById("main-configuration");
+ if (dc22 != null) {
+ JsonObject json = (JsonObject)
dc22.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("main-configuration", json);
}
- if (arr != null && !arr.isEmpty()) {
- // store traces in a special file
- LOG.trace("Updating trace file: {}", traceFile);
- String data = json.toJson() + System.lineSeparator();
- IOHelper.appendText(data, traceFile);
- json = arr.getMap(arr.size() - 1);
- traceFilePos = json.getLong("uid");
+ }
+ DevConsole dc23 = dcr.resolveById("receive");
+ if (dc23 != null) {
+ JsonObject json = (JsonObject)
dc23.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("receive", json);
}
}
- } catch (Exception e) {
- // ignore
- LOG.trace("Error updating trace file: {} due to: {}. This
exception is ignored.",
- traceFile, e.getMessage(), e);
- }
- try {
- DevConsole dc13 =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class)
- .resolveById("debug");
- if (dc13 != null) {
- JsonObject json = (JsonObject)
dc13.call(DevConsole.MediaType.JSON);
- // store debugs in a special file
- LOG.trace("Updating debug file: {}", debugFile);
- String data = json.toJson() + System.lineSeparator();
- IOHelper.writeText(data, debugFile);
+ DevConsole dc24 = dcr.resolveById("internal-tasks");
+ if (dc24 != null) {
+ JsonObject json = (JsonObject)
dc24.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("internal-tasks", json);
+ }
}
- } catch (Exception e) {
- // ignore
- LOG.trace("Error updating debug file: {} due to: {}. This
exception is ignored.",
- debugFile, e.getMessage(), e);
- }
- try {
- DevConsole dc13b =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class)
- .resolveById("message-history");
- if (dc13b != null) {
- JsonObject json = (JsonObject)
dc13b.call(DevConsole.MediaType.JSON);
- // store replays in a special file
- LOG.trace("Updating message-history file: {}",
messageHistoryFile);
- String data = json.toJson() + System.lineSeparator();
- IOHelper.writeText(data, messageHistoryFile);
+ DevConsole dc25 = dcr.resolveById("groovy");
+ if (dc25 != null) {
+ JsonObject json = (JsonObject)
dc25.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("groovy", json);
+ }
}
- } catch (Exception e) {
- // ignore
- LOG.trace("Error updating message-history file: {} due to: {}.
This exception is ignored.",
- messageHistoryFile, e.getMessage(), e);
- }
- try {
- DevConsole dc13c =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class)
- .resolveById("errors");
- if (dc13c != null) {
- JsonObject json = (JsonObject)
dc13c.call(DevConsole.MediaType.JSON,
+ DevConsole dc26 = dcr.resolveById("errors");
+ if (dc26 != null) {
+ JsonObject json = (JsonObject)
dc26.call(DevConsole.MediaType.JSON,
Map.of("stackTrace", "true"));
if (json != null && !json.isEmpty()) {
- LOG.trace("Updating error file: {}", errorFile);
- String data = json.toJson() + System.lineSeparator();
- IOHelper.writeText(data, errorFile);
+ // only include metadata in status file (full error data
is in the error file)
+ JsonObject summary = new JsonObject();
+ summary.put("enabled", json.get("enabled"));
+ summary.put("size", json.get("size"));
+ summary.put("maximumEntries", json.get("maximumEntries"));
+ summary.put("timeToLive", json.get("timeToLive"));
+ root.put("errors", summary);
}
}
- } catch (Exception e) {
- // ignore
- LOG.trace("Error updating error file: {} due to: {}. This
exception is ignored.",
- errorFile, e.getMessage(), e);
- }
- try {
- DevConsole dc14 =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class)
- .resolveById("receive");
- if (dc14 != null) {
- JsonObject json = (JsonObject)
dc14.call(DevConsole.MediaType.JSON, Map.of("dump", "true"));
- JsonArray arr = json.getCollection("messages");
- // filter based on last uid
- if (receiveFilePos > 0) {
- arr.removeIf(r -> {
- JsonObject jo = (JsonObject) r;
- return jo.getLong("uid") <= receiveFilePos;
- });
+ DevConsole dc27 = dcr.resolveById("datasource");
+ if (dc27 != null) {
+ JsonObject json = (JsonObject)
dc27.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("dataSources", json);
}
- if (arr != null && !arr.isEmpty()) {
- // store messages in a special file
- LOG.trace("Updating receive file: {}", receiveFile);
- String data = json.toJson() + System.lineSeparator();
- IOHelper.appendText(data, receiveFile);
- json = arr.getMap(arr.size() - 1);
- receiveFilePos = json.getLong("uid");
+ }
+ DevConsole dc28 = dcr.resolveById("sql-trace");
+ if (dc28 != null) {
+ JsonObject json = (JsonObject)
dc28.call(DevConsole.MediaType.JSON);
+ if (json != null && !json.isEmpty()) {
+ root.put("sqlTrace", json);
}
}
- } catch (Exception e) {
- // ignore
- LOG.trace("Error updating receive file: {} due to: {}. This
exception is ignored.",
- receiveFile, e.getMessage(), e);
+ JsonArray consoleIds = new JsonArray();
+ consoleIds.addAll(dcr.getConsoleIDs());
+ root.put("devConsoles", consoleIds);
}
- try {
- DevConsole dc15 =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class)
- .resolveById("activity");
- if (dc15 != null) {
- JsonObject json = (JsonObject)
dc15.call(DevConsole.MediaType.JSON);
- LOG.trace("Updating activity file: {}", activityFile);
- String data = json.toJson() + System.lineSeparator();
- IOHelper.writeText(data, activityFile);
+ boolean hasBrowseable = false;
+ for (Endpoint ep : getCamelContext().getEndpoints()) {
+ if (ep instanceof BrowsableEndpoint) {
+ hasBrowseable = true;
+ break;
}
- } catch (Exception e) {
- // ignore
- LOG.trace("Error updating activity file: {} due to: {}. This
exception is ignored.",
- activityFile, e.getMessage(), e);
}
+ root.put("hasBrowseableEndpoints", hasBrowseable);
+ // various details
+ JsonObject mem = collectMemory();
+ if (mem != null) {
+ root.put("memory", mem);
+ }
+ JsonObject cl = collectClassLoading();
+ if (cl != null) {
+ root.put("classLoading", cl);
+ }
+ JsonObject threads = collectThreads();
+ if (threads != null) {
+ root.put("threads", threads);
+ }
+ JsonObject gc = collectGC();
+ if (gc != null) {
+ root.put("gc", gc);
+ }
+ JsonObject vaults = collectVaults();
+ if (!vaults.isEmpty()) {
+ root.put("vaults", vaults);
+ }
+ return root;
+ }
+
+ @Override
+ public JsonObject trace() throws Exception {
+ return callConsole("trace", Map.of("dump", "true"));
+ }
+
+ @Override
+ public JsonObject debug() throws Exception {
+ return callConsole("debug", Map.of());
+ }
+
+ @Override
+ public JsonObject messageHistory() throws Exception {
+ return callConsole("message-history", Map.of());
+ }
+
+ @Override
+ public JsonObject errors() throws Exception {
+ return callConsole("errors", Map.of("stackTrace", "true"));
+ }
+
+ @Override
+ public JsonObject receive() throws Exception {
+ return callConsole("receive", Map.of("dump", "true"));
+ }
+
+ @Override
+ public JsonObject activity() throws Exception {
+ return callConsole("activity", Map.of());
+ }
+
+ private JsonObject callConsole(String id, Map<String, Object> options) {
+ DevConsole dc =
camelContext.getCamelContextExtension().getContextPlugin(DevConsoleRegistry.class)
Review Comment:
Applied in 478cc27a7f5f: `callConsole` now guards the registry as `status()`
does. (It is installed whenever camel-console is on the classpath, which
camel-cli-connector depends on, and snapshot calls are wrapped in a catch in
both transports, so this is consistency rather than a reachable crash.)
_Claude Code on behalf of Croway_
##########
dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/WebSocketCliConnectorTransport.java:
##########
@@ -0,0 +1,732 @@
+/*
+ * 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.support.service.ServiceSupport;
+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 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.reconnectDelay", "1000"));
+ reconnectMaxDelay =
Long.parseLong(property("camel.cli.websocket.reconnectMaxDelay", "30000"));
+ // as the file transport: faster when debugging
+ snapshotInterval = Long.parseLong(
+ property("camel.cli.websocket.snapshotInterval",
camelContext.isDebugging() ? "100" : "1000"));
+ heartbeatInterval =
Long.parseLong(property("camel.cli.websocket.heartbeatInterval", "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) {
+ // e.g. camel-cli-debug makes it faster so breakpoints show up quickly
+ snapshotInterval = delay;
+ execute(this::scheduleSnapshots);
+ }
+
+ private void scheduleSnapshots() {
+ if (snapshotFuture != null) {
+ snapshotFuture.cancel(false);
+ }
+ snapshotFuture = scheduler.scheduleWithFixedDelay(() ->
safely(this::snapshotTask), snapshotInterval,
+ snapshotInterval, TimeUnit.MILLISECONDS);
+ }
+
+ @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));
Review Comment:
Not applied: a pending ping does not block `sendText`. Per the
`java.net.http.WebSocket` javadoc, `sendText` fails with
`IllegalStateException` only "if there is a pending text or binary send
operation", and `sendPing` only "if there is a pending ping or pong send
operation". Awaiting the ping would block the scheduler thread (snapshots,
heartbeats) for up to 10 s on a slow network. The only overlap is a new ping
while the previous one, 10 s earlier, is still pending: that ping fails without
affecting the connection.
_Claude Code on behalf of Croway_
##########
dsl/camel-cli-connector/src/main/java/org/apache/camel/cli/connector/FileCliConnectorTransport.java:
##########
@@ -0,0 +1,391 @@
+/*
+ * 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.io.File;
+import java.io.FileInputStream;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.util.FileUtil;
+import org.apache.camel.util.HomeHelper;
+import org.apache.camel.util.IOHelper;
+import org.apache.camel.util.concurrent.ThreadHelper;
+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 used by the Camel CLI: polls <tt>~/.camel/{pid}-action*.json</tt>
files for actions, writes their results
+ * to <tt>{pid}-output*.json</tt>, and the snapshots to
<tt>{pid}-status.json</tt>, <tt>{pid}-trace.json</tt> and so on.
+ * Deleting the <tt>{pid}</tt> lock file shuts down the integration.
+ */
+public class FileCliConnectorTransport extends ServiceSupport implements
CliConnectorTransport {
+
+ // log as the connector, which is where users know to find these messages
+ private static final Logger LOG =
LoggerFactory.getLogger(LocalCliConnector.class);
+
+ private CamelContext camelContext;
+ private CliActionDispatcher dispatcher;
+ private CliSnapshotProducer snapshots;
+ private Runnable shutdown;
+ private ScheduledExecutorService executor;
+ private ScheduledFuture<?> scheduledFuture;
+ private int delay = 1000;
+ private long counter;
+ private final AtomicBoolean terminating = new AtomicBoolean();
+ private File lockFile;
+ private File statusFile;
+ private File actionFile;
+ private File outputFile;
+ private File traceFile;
+ private long traceFilePos; // keep track of trace offset
+ private File messageHistoryFile;
+ private File debugFile;
+ private File errorFile;
+ private File receiveFile;
+ private long receiveFilePos; // keep track of receive offset
+ private File activityFile;
+
+ @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 {
+ terminating.set(false);
+
+ // create thread from JDK so it is not managed by Camel because we
want the pool to be independent when
+ // camel is being stopped which otherwise can lead to stopping the
thread pool while the task is running
+ executor = Executors.newSingleThreadScheduledExecutor(r -> {
+ String threadName = ThreadHelper.resolveThreadName(null,
"LocalCliConnector");
+ return new Thread(r, threadName);
+ });
+
+ // make it go faster in debug mode
+ if (camelContext.isDebugging()) {
+ delay = 100;
+ }
+
+ lockFile = createLockFile(getPid());
+ if (lockFile != null) {
+ statusFile = createLockFile(lockFile.getName() + "-status.json");
+ actionFile = createLockFile(lockFile.getName() + "-action.json");
+ outputFile = createLockFile(lockFile.getName() + "-output.json");
+ traceFile = createLockFile(lockFile.getName() + "-trace.json");
+ messageHistoryFile = createLockFile(lockFile.getName() +
"-history.json");
+ errorFile = createLockFile(lockFile.getName() + "-error.json");
+ debugFile = createLockFile(lockFile.getName() + "-debug.json");
+ receiveFile = createLockFile(lockFile.getName() + "-receive.json");
+ activityFile = createLockFile(lockFile.getName() +
"-activity.json");
+ scheduledFuture = executor.scheduleWithFixedDelay(this::task, 0,
delay, TimeUnit.MILLISECONDS);
+ LOG.info("Camel CLI connector enabled");
+ } else {
+ LOG.warn("Cannot create PID file: {}. This integration cannot be
managed by Camel CLI connector.", getPid());
+ }
+ }
+
+ @Override
+ public void updateDelay(int delay) {
+ if (this.delay == delay) {
+ return;
+ }
+ if (scheduledFuture != null) {
+ try {
+ scheduledFuture.cancel(true);
+ } catch (Exception e) {
+ // ignore
+ }
+ }
+ boolean done = scheduledFuture == null || scheduledFuture.isDone();
+ if (done) {
+ this.delay = delay;
+ scheduledFuture = executor.scheduleWithFixedDelay(this::task, 0,
delay, TimeUnit.MILLISECONDS);
+ }
+ }
+
+ protected void task() {
+ if (!lockFile.exists() && terminating.compareAndSet(false, true)) {
+ // if the lock file is deleted then trigger termination
+ shutdown.run();
+ return;
+ }
+ if (!statusFile.exists()) {
+ return;
+ }
+
+ actionTask();
+ statusTask();
+ // only run this every 2nd time as gathering this data has more
overhead
+ // and are only needed when doing tracing/debugging/receive
+ if (++counter % 2 == 0) {
+ traceTask();
+ }
+ }
+
+ protected void actionTask() {
+ // scan for all action files: {pid}-action.json (legacy) and
{pid}-action-{requestId}.json (multi-client)
+ File dir = lockFile.getParentFile();
+ String prefix = lockFile.getName() + "-action";
+ File[] actionFiles = dir.listFiles((d, name) ->
name.startsWith(prefix) && name.endsWith(".json"));
+ if (actionFiles == null || actionFiles.length == 0) {
+ return;
+ }
+ for (File af : actionFiles) {
+ String suffix = af.getName().substring(prefix.length());
+ // suffix is either ".json" (legacy) or "-{requestId}.json"
(multi-client)
+ String requestId = suffix.startsWith("-")
+ ? suffix.substring(1, suffix.length() - 5) // strip
leading "-" and trailing ".json"
+ : null;
+ File of = requestId != null
+ ? new File(dir, lockFile.getName() + "-output-" +
requestId + ".json")
+ : this.outputFile;
+ processAction(af, of);
+ }
+ }
+
+ private void processAction(File af, File of) {
+ String action = null;
+ try {
+ JsonObject root = loadAction(af);
+ if (root == null || root.isEmpty()) {
+ return;
+ }
+ action = root.getString("action");
+ dispatcher.dispatch(root, result -> {
+ LOG.trace("Updating output file: {}", of);
+ IOHelper.writeText(result.toJson(), of);
+ });
+ } catch (Exception e) {
+ LOG.warn("Error executing action: {} due to: {}. This exception is
ignored.", action != null ? action : af,
+ e.getMessage(),
+ e);
+ } finally {
+ FileUtil.deleteFile(af);
+ }
+ }
+
+ private static JsonObject loadAction(File file) {
+ try {
+ if (file != null && file.exists()) {
+ FileInputStream fis = new FileInputStream(file);
+ String text = IOHelper.loadText(fis);
+ IOHelper.close(fis);
Review Comment:
Applied in 478cc27a7f5f (try-with-resources).
_Claude Code on behalf of Croway_
--
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]