/* Copyright © 2026 Thomas von Dein This program is free software: you can redistribute it and/or modify it under the terms of the GNU General Public License as published by the Free Software Foundation, either version 3 of the License, or (at your option) any later version. This program is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License for more details. You should have received a copy of the GNU General Public License along with this program. If not, see . */ package es import ( "context" "fmt" "log/slog" "slices" "strings" "sync" "codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/printer" "github.com/dustin/go-humanize" "github.com/elastic/go-elasticsearch/v9/typedapi/ccr/stats" "github.com/elastic/go-elasticsearch/v9/typedapi/cluster/health" clusterstats "github.com/elastic/go-elasticsearch/v9/typedapi/cluster/stats" "github.com/elastic/go-elasticsearch/v9/typedapi/core/info" "github.com/elastic/go-elasticsearch/v9/typedapi/types" ) const ( ResponseHealth = iota ResponseInfo ResponseCcr ResponseStats ) type ClusterIndices map[string]map[string]*types.IndicesRecord type apiResponse struct { error error info *info.Response health *health.Response ccr *stats.Response stats *clusterstats.Response which int } func ClusterList(conf *cfg.Config) error { table := printer.NewTable(conf, 3, len(conf.Clusters)) table.Addheaders("cluster", "uri", "default") idx := 0 names := make([]string, len(conf.Clusters)) for name := range conf.Clusters { names[idx] = name idx++ } slices.Sort(names) for idx, name := range names { current := name == "default" || name == conf.CurrentCluster cluster := conf.Clusters[name] _, err := cluster.ES.Cluster.Health(). Header("content-type", "application/json"). Header("accept", "application/json"). Do(context.Background()) if err == nil { name = printer.Colorize(conf, "green", name) } table.Entries[idx] = []string{name, cluster.Uri, fmt.Sprintf("%t", current)} } if err := table.Print(); err != nil { return err } return nil } // We're using goroutines here to parallelize API requests, since we // have to do 3 of'em for each cluster. This speeds things up. func ClusterStatus(conf *cfg.Config) error { clusters := []string{} gocount := 3 if conf.Verbose { gocount++ } if conf.All { for key := range conf.Clusters { clusters = append(clusters, key) } } else { clusters = []string{"default"} } for _, cluster := range clusters { es := conf.DefaultCluster.ES if cluster != "default" { es = conf.Clusters[cluster].ES } responses := make(chan apiResponse, gocount) wg := &sync.WaitGroup{} wg.Add(gocount) go getClusterData(es, wg, responses, "health") go getClusterData(es, wg, responses, "info") go getClusterData(es, wg, responses, "ccrstats") if conf.Verbose { go getClusterData(es, wg, responses, "stats") } wg.Wait() var clusterhealth *health.Response var info *info.Response var ccrstats *stats.Response var clusterstats *clusterstats.Response for i := 0; i < gocount; i++ { r := <-responses if r.error != nil { return r.error } switch r.which { case ResponseHealth: clusterhealth = r.health case ResponseCcr: ccrstats = r.ccr case ResponseInfo: info = r.info case ResponseStats: clusterstats = r.stats } } slog.Debug("ES result", "cluster health", clusterhealth) ccrfollowing := "" if len(ccrstats.AutoFollowStats.AutoFollowedClusters) > 0 { // is following another cluster ccrfollowing = fmt.Sprintf("%s (%d/%d)", ccrstats.AutoFollowStats.AutoFollowedClusters[0].ClusterName, ccrstats.AutoFollowStats.NumberOfSuccessfulFollowIndices, ccrstats.AutoFollowStats.NumberOfFailedFollowIndices, ) } table := printer.NewTable(conf, 2, 5) table.Addheaders(cluster, "status") table.Entries = [][]string{ {"Cluster Name", printer.Colorize(conf, clusterhealth.Status.Name, clusterhealth.ClusterName)}, {"ES Version", info.Version.Int}, {"Active Shards", fmt.Sprintf("%d", clusterhealth.ActiveShards)}, {"Active Primary Shards", fmt.Sprintf("%d", clusterhealth.ActivePrimaryShards)}, {"Nodes", fmt.Sprintf("%d", clusterhealth.NumberOfNodes)}, {"AutoFollow (success/failed indices)", ccrfollowing}, } if conf.Verbose { table = gatherClusterStats(conf, clusterstats, table) } if err := table.Print(); err != nil { return err } } return nil } func gatherClusterStats(conf *cfg.Config, clusterstats *clusterstats.Response, table *printer.Table) *printer.Table { var querycount int64 var vmversion string for _, count := range clusterstats.Indices.Search.Queries { querycount += count } if len(clusterstats.Nodes.Jvm.Versions) > 0 { vmversion = strings.Join([]string{ clusterstats.Nodes.Jvm.Versions[0].VmName, clusterstats.Nodes.Jvm.Versions[0].VmVersion}, " ") } isleader := checkClusterIsLeader(conf, conf.CurrentCluster) table.Entries = append(table.Entries, [][]string{ {"Indicies", fmt.Sprintf("%d", clusterstats.Indices.Count)}, {"Is Leader", fmt.Sprintf("%t", isleader)}, {"Docs", fmt.Sprintf("%d", clusterstats.Indices.Docs.Count)}, {"Total Size", humanize.Bytes(uint64(clusterstats.Indices.Docs.TotalSizeInBytes))}, {"Total Queries", fmt.Sprintf("%d", querycount)}, {"Shards Primaries", fmt.Sprintf("%d", clusterstats.Indices.Shards.Primaries)}, {"Shards Total", fmt.Sprintf("%d", clusterstats.Indices.Shards.Total)}, {"Storage", fmt.Sprintf( "%s/%s", humanize.Bytes(uint64(clusterstats.Indices.Store.SizeInBytes)), humanize.Bytes(uint64(*clusterstats.Indices.Store.TotalDataSetSizeInBytes)), )}, {"JVM Heap", fmt.Sprintf( "%s/%s", humanize.Bytes(uint64(clusterstats.Nodes.Jvm.Mem.HeapUsedInBytes)), humanize.Bytes(uint64(clusterstats.Nodes.Jvm.Mem.HeapMaxInBytes)), )}, {"JVM Threads", fmt.Sprintf("%d", clusterstats.Nodes.Jvm.Threads)}, {"JVM Version", vmversion}, {"CPUs", fmt.Sprintf("%d", clusterstats.Nodes.Os.AllocatedProcessors)}, {"CPU Usage", fmt.Sprintf("%d%%", clusterstats.Nodes.Process.Cpu.Percent)}, {"Open FDs", fmt.Sprintf("%d", clusterstats.Nodes.Process.OpenFileDescriptors.Avg)}, }...) return table }