From 8f3a2f9ac8e7619d6d7f7463e0cd6918c8c131b2 Mon Sep 17 00:00:00 2001
From: ash9146 <52308288+ash9146@users.noreply.github.com>
Date: Wed, 29 Jul 2026 13:01:10 +0900
Subject: [PATCH] Add opt-in 7-day CPU/heap trend summary from monitoring data
Point-in-time snapshots of CPU/heap usage don't show whether a cluster
has been trending toward trouble, which limits their usefulness for
after-the-fact investigations. This adds an opt-in --includeTrends flag
that queries .monitoring-es-* (when present) for a 7-day percentile
summary (p50/p95/p99) of CPU and heap usage, broken down both as an
overall summary and daily buckets, written to monitoring-trends.json.
Filters on type: node_stats to avoid scanning unrelated monitoring doc
types sharing the same source_node.name field - verified this brought
scanned docs from ~10M/node down to ~55k/node over 7 days on a test
cluster, in line with the expected count at a 10s collection interval.
Scoped to legacy self-monitoring indices for now; 8.x Metricbeat-based
monitoring data streams use a different schema and aren't covered.
---
README.md | 5 +
.../support/diagnostics/DiagnosticInputs.java | 3 +
.../chain/DiagnosticChainExec.java | 4 +
.../commands/CollectMonitoringTrends.java | 94 +++++++++++++++++++
4 files changed, 106 insertions(+)
create mode 100644 src/main/java/co/elastic/support/diagnostics/commands/CollectMonitoringTrends.java
diff --git a/README.md b/README.md
index d17f69a2..89ee4635 100644
--- a/README.md
+++ b/README.md
@@ -252,6 +252,11 @@ Elasticsearch, Kibana, and Logstash each have three distinct execution modes ava
Option only - no value. |
+
+ | --includeTrends |
+ Collect a 7-day CPU/heap trend summary from monitoring data (.monitoring-es-*), broken down by node and by day, if monitoring is enabled on the target cluster. Adds one additional query against the monitored cluster. Default value is false. |
+ Option only - no value. |
+
#### PKI Authentication Options
diff --git a/src/main/java/co/elastic/support/diagnostics/DiagnosticInputs.java b/src/main/java/co/elastic/support/diagnostics/DiagnosticInputs.java
index 469f0e44..c9242c37 100644
--- a/src/main/java/co/elastic/support/diagnostics/DiagnosticInputs.java
+++ b/src/main/java/co/elastic/support/diagnostics/DiagnosticInputs.java
@@ -106,6 +106,7 @@ public class DiagnosticInputs extends ElasticRestClientInputs {
public final static String knownHostsDescription = "Known hosts file to search for target server. Default is ~/.ssh/known_hosts for Linux/Mac. Windows users should always set this explicitly.";
public final static String sudoDescription = "Use sudo for remote commands? If not used, log retrieval and some system calls may fail.";
public final static String remotePortDescription = "SSH port for the host being queried.";
+ public final static String includeTrendsDescription = "Collect a 7-day CPU/heap trend summary from monitoring data (.monitoring-es-*), if present. Adds one additional query against the monitored cluster.";
// Input Fields
@Parameter(names = {
@@ -133,6 +134,8 @@ public class DiagnosticInputs extends ElasticRestClientInputs {
public String knownHostsFile = "";
@Parameter(names = { "--sudo" }, description = sudoDescription)
public boolean isSudo = false;
+ @Parameter(names = { "--includeTrends" }, description = includeTrendsDescription)
+ public boolean includeTrends = false;
@Parameter(names = { "--remotePort" }, description = remotePortDescription)
public int remotePort = 22;
// End Input Fields
diff --git a/src/main/java/co/elastic/support/diagnostics/chain/DiagnosticChainExec.java b/src/main/java/co/elastic/support/diagnostics/chain/DiagnosticChainExec.java
index 533fdc7d..87f63794 100644
--- a/src/main/java/co/elastic/support/diagnostics/chain/DiagnosticChainExec.java
+++ b/src/main/java/co/elastic/support/diagnostics/chain/DiagnosticChainExec.java
@@ -14,6 +14,7 @@
import co.elastic.support.diagnostics.commands.CheckPlatformDetails;
import co.elastic.support.diagnostics.commands.CheckUserAuthLevel;
import co.elastic.support.diagnostics.commands.CollectDockerInfo;
+import co.elastic.support.diagnostics.commands.CollectMonitoringTrends;
import co.elastic.support.diagnostics.commands.CollectKibanaLogs;
import co.elastic.support.diagnostics.commands.CollectLogs;
import co.elastic.support.diagnostics.commands.CollectSystemCalls;
@@ -39,6 +40,7 @@ public static void runDiagnostic(DiagnosticContext context, String type) throws
// Removed temporarily due to issues with finding and accessing cloud master
// new CheckPlatformDetails().execute(context);
new RunClusterQueries().execute(context);
+ new CollectMonitoringTrends().execute(context);
break;
case Constants.local:
@@ -46,6 +48,7 @@ public static void runDiagnostic(DiagnosticContext context, String type) throws
new CheckUserAuthLevel().execute(context);
new CheckPlatformDetails().execute(context);
new RunClusterQueries().execute(context);
+ new CollectMonitoringTrends().execute(context);
if (context.runSystemCalls) {
new CollectSystemCalls().execute(context);
new CollectLogs().execute(context);
@@ -61,6 +64,7 @@ public static void runDiagnostic(DiagnosticContext context, String type) throws
new CheckUserAuthLevel().execute(context);
new CheckPlatformDetails().execute(context);
new RunClusterQueries().execute(context);
+ new CollectMonitoringTrends().execute(context);
if (context.runSystemCalls) {
new CollectSystemCalls().execute(context);
new CollectLogs().execute(context);
diff --git a/src/main/java/co/elastic/support/diagnostics/commands/CollectMonitoringTrends.java b/src/main/java/co/elastic/support/diagnostics/commands/CollectMonitoringTrends.java
new file mode 100644
index 00000000..5d4da3f5
--- /dev/null
+++ b/src/main/java/co/elastic/support/diagnostics/commands/CollectMonitoringTrends.java
@@ -0,0 +1,94 @@
+/*
+ * Copyright Elasticsearch B.V. and/or licensed to Elasticsearch B.V. under one
+ * or more contributor license agreements. Licensed under the Elastic License
+ * 2.0; you may not use this file except in compliance with the Elastic License
+ * 2.0.
+ */
+package co.elastic.support.diagnostics.commands;
+
+import co.elastic.support.Constants;
+import co.elastic.support.diagnostics.chain.Command;
+import co.elastic.support.diagnostics.chain.DiagnosticContext;
+import co.elastic.support.rest.RestClient;
+import co.elastic.support.rest.RestResult;
+import co.elastic.support.util.JsonYamlUtils;
+import com.fasterxml.jackson.databind.JsonNode;
+import org.apache.commons.io.FileUtils;
+import org.apache.http.HttpEntity;
+import org.apache.http.HttpResponse;
+import org.apache.http.util.EntityUtils;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.io.File;
+
+public class CollectMonitoringTrends implements Command {
+
+ private static final Logger logger = LogManager.getLogger(CollectMonitoringTrends.class);
+
+ private static final String AGG_QUERY = "{"
+ + "\"size\": 0,"
+ + "\"timeout\": \"10s\","
+ + "\"query\": { \"bool\": { \"filter\": ["
+ + " { \"term\": { \"type\": \"node_stats\" } },"
+ + " { \"range\": { \"timestamp\": { \"gte\": \"now-7d\" } } }"
+ + "] } },"
+ + "\"aggs\": {"
+ + " \"by_node\": {"
+ + " \"terms\": { \"field\": \"source_node.name\", \"size\": 50 },"
+ + " \"aggs\": {"
+ + " \"cpu_pct\": { \"percentiles\": { \"field\": \"node_stats.process.cpu.percent\", \"percents\": [50, 95, 99] } },"
+ + " \"heap_pct\": { \"percentiles\": { \"field\": \"node_stats.jvm.mem.heap_used_percent\", \"percents\": [50, 95, 99] } },"
+ + " \"by_day\": {"
+ + " \"date_histogram\": { \"field\": \"timestamp\", \"fixed_interval\": \"1d\" },"
+ + " \"aggs\": {"
+ + " \"cpu_pct\": { \"percentiles\": { \"field\": \"node_stats.process.cpu.percent\", \"percents\": [50, 95, 99] } },"
+ + " \"heap_pct\": { \"percentiles\": { \"field\": \"node_stats.jvm.mem.heap_used_percent\", \"percents\": [50, 95, 99] } }"
+ + " }"
+ + " }"
+ + " }"
+ + " }"
+ + "}"
+ + "}";
+
+ public void execute(DiagnosticContext context) {
+ if (!context.diagnosticInputs.includeTrends) {
+ return;
+ }
+
+ try {
+ RestClient client = context.resourceCache.getRestClient(Constants.restInputHost);
+
+ // 1) Check whether monitoring indices exist at all - skip quietly if not.
+ RestResult checkResult = client.execQuery("/.monitoring-es-*/_search?size=0");
+ JsonNode checkNode = JsonYamlUtils.createJsonNodeFromString(checkResult.toString());
+ long totalShards = checkNode.path("_shards").path("total").asLong(0);
+
+ if (totalShards == 0) {
+ logger.info(Constants.CONSOLE, "No monitoring indices found - skipping trend summary.");
+ return;
+ }
+
+ // 2) Run the aggregation query.
+ HttpResponse response = client.execPost("/.monitoring-es-*/_search", AGG_QUERY);
+ int status = response.getStatusLine().getStatusCode();
+ HttpEntity entity = response.getEntity();
+ String body = entity != null ? EntityUtils.toString(entity) : "";
+
+ if (status < 200 || status >= 300) {
+ logger.info(Constants.CONSOLE, "Monitoring trend query failed (status {}) - skipping.", status);
+ return;
+ }
+
+ // 3) Write the result to the diagnostic output directory.
+ File outFile = new File(context.tempDir, "monitoring-trends.json");
+ FileUtils.writeStringToFile(outFile, body, "UTF-8");
+ logger.info(Constants.CONSOLE, "Monitoring trend summary written to: {}", outFile.getName());
+
+ } catch (Exception e) {
+ // This feature failing should never block the rest of the diagnostic.
+ logger.info(Constants.CONSOLE, "Could not collect monitoring trend summary - bypassing.");
+ logger.error("Error collecting monitoring trends", e);
+ }
+ }
+}