From fcfeaf66a959e937c056151c2c639d974d0e21bb Mon Sep 17 00:00:00 2001 From: "T. von Dein" Date: Tue, 5 May 2026 12:06:34 +0200 Subject: [PATCH] Add nodes and cluster settings support (#9) --- README.md | 3 +- cmd/cluster.go | 78 +++++++++- cmd/node.go | 72 +++++++++ cmd/root.go | 1 + pkg/cfg/config.go | 29 ++-- pkg/es/cluster.go | 362 +++++++--------------------------------------- pkg/es/node.go | 55 +++++++ 7 files changed, 265 insertions(+), 335 deletions(-) create mode 100644 cmd/node.go create mode 100644 pkg/es/node.go diff --git a/README.md b/README.md index 108cbf5..d30a288 100644 --- a/README.md +++ b/README.md @@ -12,13 +12,14 @@ USAGE: esctl [global options] [command [command options]] VERSION: - v0.0.3 + v0.0.4 COMMANDS: search, / search within an index index, i manage indicies snapshot, snap manage snapshots cluster, c manage cluster[s] + node, snap manage nodes help, h Shows a list of commands or help for one command GLOBAL OPTIONS: diff --git a/cmd/cluster.go b/cmd/cluster.go index 00f1cb9..7b086d4 100644 --- a/cmd/cluster.go +++ b/cmd/cluster.go @@ -33,21 +33,22 @@ func Cluster(conf *cfg.Config) *cli.Command { Usage: "manage cluster[s]", Commands: []*cli.Command{ - Compare(conf), - Status(conf), - List(conf), + ClusterCompare(conf), + ClusterStatus(conf), + ClusterList(conf), + ClusterSettings(conf), }, } } -func List(conf *cfg.Config) *cli.Command { +func ClusterList(conf *cfg.Config) *cli.Command { return &cli.Command{ Name: "list", Usage: "list configured clusters", Aliases: []string{"ls"}, Action: func(ctx context.Context, cmd *cli.Command) error { - if err := es.List(conf); err != nil { + if err := es.ClusterList(conf); err != nil { return err } @@ -56,7 +57,7 @@ func List(conf *cfg.Config) *cli.Command { } } -func Status(conf *cfg.Config) *cli.Command { +func ClusterStatus(conf *cfg.Config) *cli.Command { return &cli.Command{ Name: "status", Usage: "show cluster status", @@ -72,7 +73,7 @@ func Status(conf *cfg.Config) *cli.Command { }, Action: func(ctx context.Context, cmd *cli.Command) error { - if err := es.Status(conf); err != nil { + if err := es.ClusterStatus(conf); err != nil { return err } @@ -81,7 +82,7 @@ func Status(conf *cfg.Config) *cli.Command { } } -func Compare(conf *cfg.Config) *cli.Command { +func ClusterCompare(conf *cfg.Config) *cli.Command { return &cli.Command{ Name: "compare", Aliases: []string{"c"}, @@ -119,3 +120,64 @@ func Compare(conf *cfg.Config) *cli.Command { }, } } + +func ClusterSettings(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "settings", + Usage: "cluster settings management", + Aliases: []string{"config"}, + + Commands: []*cli.Command{ + ClusterSettingsList(conf), + ClusterSettingsSet(conf), + }, + } +} + +func ClusterSettingsList(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "list", + Usage: "show cluster settings", + Aliases: []string{"ls", "get"}, + + Action: func(ctx context.Context, cmd *cli.Command) error { + if err := es.ClusterSettingsList(conf); err != nil { + return err + } + + return nil + }, + } +} + +func ClusterSettingsSet(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "set", + Usage: "set|update cluster settings", + Aliases: []string{"set", "update"}, + UsageText: "set [options] setting:value [setting:value ...]", + + Flags: []cli.Flag{ + &cli.BoolFlag{ + Name: "persistent", + Usage: "add persistent setting[s] (default)", + Destination: &conf.Persistent, + Aliases: []string{"p"}, + }, + &cli.BoolFlag{ + Name: "transient", + Usage: "add transient setting[s]", + Destination: &conf.Transient, + Aliases: []string{"t"}, + }, + }, + + Action: func(ctx context.Context, cmd *cli.Command) error { + if err := es.ClusterSettingsSet(conf, cmd.Args()); err != nil { + return err + } + + return nil + }, + } +} diff --git a/cmd/node.go b/cmd/node.go new file mode 100644 index 0000000..73c823c --- /dev/null +++ b/cmd/node.go @@ -0,0 +1,72 @@ +/* +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 cmd + +import ( + "context" + + "codeberg.org/scip/esctl/pkg/cfg" + "codeberg.org/scip/esctl/pkg/es" + + "github.com/urfave/cli/v3" +) + +func Node(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "node", + Aliases: []string{"snap"}, + Usage: "manage nodes", + + Commands: []*cli.Command{ + NodeList(conf), + NodeShow(conf), + }, + } +} + +func NodeList(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "list", + Aliases: []string{"ls"}, + Usage: "list nodes", + + Action: func(ctx context.Context, cmd *cli.Command) error { + if err := es.NodeList(conf); err != nil { + return err + } + + return nil + }, + } +} + +func NodeShow(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "show", + Aliases: []string{"sh"}, + Usage: "show details about a node", + UsageText: "show [options] ", + + Action: func(ctx context.Context, cmd *cli.Command) error { + // if err := es.NodeShow(conf, cmd.Args().Get(0)); err != nil { + // return err + // } + + return nil + }, + } +} diff --git a/cmd/root.go b/cmd/root.go index c7ea297..5097d02 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -76,6 +76,7 @@ func Main() int { Index(conf), Snapshot(conf), Cluster(conf), + Node(conf), }, Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) { diff --git a/pkg/cfg/config.go b/pkg/cfg/config.go index 3785e96..a6f0514 100644 --- a/pkg/cfg/config.go +++ b/pkg/cfg/config.go @@ -30,7 +30,7 @@ import ( ) const ( - Version string = `v0.0.4` + Version string = `v0.0.5` ) type Cluster struct { @@ -39,19 +39,20 @@ type Cluster struct { } type Config struct { - ConfigFile string // -c - CurrentCluster string // -C - Debug bool // -d - Clusters map[string]*Cluster - DefaultCluster *Cluster - Index string // index: -i - Failed, Partials bool // index: flags - Shards, Replicas int // index create: -s -r - Wait bool // index create: -w - From, To, MaxItems int // search: flags - Filter []string // search: -F - Exclude string // cluster compare: -e (regexp) - All bool // cluster status: -a + ConfigFile string // -c + CurrentCluster string // -C + Debug bool // -d + Clusters map[string]*Cluster + DefaultCluster *Cluster + Index string // index: -i + Failed, Partials bool // index: flags + Shards, Replicas int // index create: -s -r + Wait bool // index create: -w + From, To, MaxItems int // search: flags + Filter []string // search: -F + Exclude string // cluster compare: -e (regexp) + All bool // cluster status: -a + Persistent, Transient bool // -p -t cluster settings set } func NewConfig() *Config { diff --git a/pkg/es/cluster.go b/pkg/es/cluster.go index 9c5ccff..3860a56 100644 --- a/pkg/es/cluster.go +++ b/pkg/es/cluster.go @@ -18,15 +18,14 @@ package es import ( "context" + "encoding/json" "errors" "fmt" - "log" "log/slog" - "regexp" "codeberg.org/scip/esctl/pkg/cfg" - "github.com/elastic/go-elasticsearch/v9/typedapi/cluster/health" "github.com/elastic/go-elasticsearch/v9/typedapi/types" + "github.com/urfave/cli/v3" ) const ( @@ -75,7 +74,7 @@ func ClusterCompare(conf *cfg.Config, leader, follower string) error { return nil } -func List(conf *cfg.Config) error { +func ClusterList(conf *cfg.Config) error { table := NewTable(2, len(conf.Clusters)) table.Addheaders("cluster", "uri") @@ -91,7 +90,7 @@ func List(conf *cfg.Config) error { return nil } -func Status(conf *cfg.Config) error { +func ClusterStatus(conf *cfg.Config) error { clusters := []string{} if conf.All { @@ -113,7 +112,7 @@ func Status(conf *cfg.Config) error { Header("accept", "application/json"). Do(context.Background()) if err != nil { - log.Fatalf("Error getting health: %s", err) + return fmt.Errorf("failed to getcluster health: %s", err) } slog.Debug("ES result", "cluster health", res) @@ -135,340 +134,79 @@ func Status(conf *cfg.Config) error { return nil } -// look for indicies only on leader -func findIndicesOnlyOnMaster(conf *cfg.Config, indices ClusterIndices, leader, follower string) { - exclude := regexp.MustCompile(DefaultExclude) - if conf.Exclude != "" { - exclude = regexp.MustCompile(conf.Exclude) - } - - indexOnlyOnLeader := map[string]*types.IndicesRecord{} - for name, index := range indices[leader] { - if exclude.MatchString(name) { - continue - } - - _, followerHasIt := indices[follower][name] - if !followerHasIt { - // fetch index details - res, err := conf.Clusters[leader].ES.Indices.Get(name). - Header("content-type", "application/json"). - Header("accept", "application/json"). - Do(context.Background()) - if err != nil { - continue // ignore it then - } - - _, defined := res[name] - if !defined { - // json response map didn't contain the index - continue - } - - isWritable := false - for _, alias := range res[name].Aliases { - if *alias.IsWriteIndex { - isWritable = true - break - } - } - - if isWritable { - // ignore index if associated alias index is writing - continue - } - - indexOnlyOnLeader[name] = index - } - } - - idx := 0 - table := NewTable(3, len(indexOnlyOnLeader)) - table.Addheaders("index only on leader", "size", "docscount") - - for name, index := range indexOnlyOnLeader { - name := Colorize("red", name) - - table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount} - idx++ - } - - table.Sort() - table.PrintMarkdown() -} - -func checkClusterFollower(conf *cfg.Config, leader string) bool { - stats, err := conf.Clusters[leader].ES.Ccr.Stats(). +func ClusterSettingsList(conf *cfg.Config) error { + res, err := conf.DefaultCluster.ES.Cluster.GetSettings(). Header("content-type", "application/json"). Header("accept", "application/json"). Do(context.Background()) if err != nil { - fmt.Printf("failed to get ccr stats from %s: %s", leader, err) - return false + return fmt.Errorf("failed to get cluster settings: %s", err) } - if len(stats.AutoFollowStats.AutoFollowedClusters) == 0 { - // is not following anyone - return true - } + table := NewTable(2, 0) + table.Addheaders("setting", "value") + entries := [][]string{} - if stats.AutoFollowStats.AutoFollowedClusters[0].ClusterName != "" { - fmt.Println("leader/follower attribution is invalid, reverse cluster attribution and retry") - return false - } + for topic, val := range res.Persistent { + data := map[string]map[string]any{} - return true -} + err := json.Unmarshal(val, &data) -// checks if both clusters are green -func checkClusterStatus(conf *cfg.Config, leader, follower string) bool { - status := map[string]*health.Response{} - - for _, cluster := range []string{leader, follower} { - st, err := conf.Clusters[leader].ES.Cluster.Health(). - Header("content-type", "application/json"). - Header("accept", "application/json"). - Do(context.Background()) if err != nil { - fmt.Printf("failed to get health from %s: %s", cluster, err) - return false + return fmt.Errorf("failed to unmarshall setting for topic %s: %s", topic, err) } - status[cluster] = st + for key, settings := range data { + for setting, value := range settings { + entries = append(entries, []string{ + fmt.Sprintf("%s.%s.%s", topic, key, setting), + fmt.Sprintf("%v", value), + }) + + } + } } - table := NewTable(3, 5) - - table.Addheaders("setting", "leader:"+leader, "follower:"+follower) - - table.entries = [][]string{ - {"Cluster Name", - Colorize(*&status[leader].Status.Name, status[leader].ClusterName), - Colorize(*&status[follower].Status.Name, status[follower].ClusterName), - }, - {"Active Shards", - fmt.Sprintf("%d", status[leader].ActiveShards), - fmt.Sprintf("%d", status[follower].ActiveShards), - }, - {"Active Primary Shards", - fmt.Sprintf("%d", status[leader].ActivePrimaryShards), - fmt.Sprintf("%d", status[follower].ActivePrimaryShards), - }, - {"Indicies", - fmt.Sprintf("%d", len(status[leader].Indices)), - fmt.Sprintf("%d", len(status[follower].Indices)), - }, - {"Nodes", - fmt.Sprintf("%d", status[leader].NumberOfNodes), - fmt.Sprintf("%d", status[follower].NumberOfNodes), - }, - } + table.entries = entries + table.Sort() table.PrintMarkdown() - if status[leader].Status.Name == "green" && status[follower].Status.Name == "green" { - return true - } - - return false + return nil } -// finds indices on both clusters which have ilm errors -func findIlmErrors(conf *cfg.Config, leader, follower string) bool { - failed := map[string]map[string]string{} +func ClusterSettingsSet(conf *cfg.Config, args cli.Args) error { + put := conf.Clusters[conf.CurrentCluster].ES.Cluster.PutSettings() - for _, cluster := range []string{leader, follower} { - ilm, err := conf.Clusters[cluster].ES.Ilm.ExplainLifecycle("_all"). - OnlyManaged(true). - Header("content-type", "application/json"). - Header("accept", "application/json"). - Do(context.Background()) - if err != nil { - fmt.Printf("failed to get ilm status from %s: %s", cluster, err) - return false - } + for _, arg := range args.Slice() { + setting, value := splitArg(arg) - failed[cluster] = map[string]string{} - - for name, ilmstate := range ilm.Indices { - count := ilmstate.(*types.LifecycleExplainManaged).FailedStepRetryCount - - if count != nil && *count > 0 { - failed[cluster][name] = fmt.Sprintf("%d", *count) - } - } - } - - if len(failed[leader]) == 0 && len(failed[follower]) == 0 { - return true - } - - for idx, cluster := range []string{leader, follower} { - which := "leader" - if idx > 0 { - which = "follower" - } - - if len(failed[cluster]) > 0 { - idx := 0 - table := NewTable(2, len(failed[cluster])) - table.Addheaders("ilm errors on "+which, "errors") - - for name, count := range failed[cluster] { - table.entries[idx] = []string{name, count} - idx++ - } - - table.Sort() - table.PrintMarkdown() - } - } - - return false -} - -// find unsynchronized indicies only present on leader -func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool { - exclude := regexp.MustCompile(DefaultExclude) - if conf.Exclude != "" { - exclude = regexp.MustCompile(conf.Exclude) - } - - indexOnlyOnLeader := map[string]*types.IndicesRecord{} - for name, index := range indices[leader] { - if exclude.MatchString(name) { - continue - } - - _, followerHasIt := indices[follower][name] - if !followerHasIt { - // fetch index details - res, err := conf.Clusters[leader].ES.Indices.Get(name). - Header("content-type", "application/json"). - Header("accept", "application/json"). - Do(context.Background()) + switch { + case conf.Transient: + message, err := json.Marshal(value) if err != nil { - continue // ignore it then + return fmt.Errorf("failed to marshall transient value <%v> to valid JSON: %s", value, err) } - - _, defined := res[name] - if !defined { - // json response map didn't contain the index - continue - } - - isWritable := false - for _, alias := range res[name].Aliases { - if *alias.IsWriteIndex { - isWritable = true - break - } - } - - if isWritable { - // ignore index if associated alias index is writing - continue - } - - indexOnlyOnLeader[name] = index - } - } - - if len(indexOnlyOnLeader) == 0 { - return true - } - - idx := 0 - table := NewTable(3, len(indexOnlyOnLeader)) - table.Addheaders("index only on leader", "size", "docscount") - - for name, index := range indexOnlyOnLeader { - name := Colorize("red", name) - - table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount} - idx++ - } - - table.Sort() - table.PrintMarkdown() - - return false -} - -// find indices only present on follower -func findOrphanedIndices(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool { - orphaned := map[string]*types.IndicesRecord{} - - for name, index := range indices[follower] { - _, leaderHasIt := indices[leader][name] - if !leaderHasIt { - // fetch index details - res, err := conf.Clusters[follower].ES.Indices.Get(name). - Header("content-type", "application/json"). - Header("accept", "application/json"). - Do(context.Background()) + put.AddTransient(setting, message) + case conf.Persistent: + fallthrough + default: + message, err := json.Marshal(value) if err != nil { - continue // ignore it then + return fmt.Errorf("failed to marshall persistent value <%v> to valid JSON: %s", value, err) } - - _, defined := res[name] - if !defined { - // json response map didn't contain the index - continue - } - - orphaned[name] = index + put.AddPersistent(setting, message) } } - if len(orphaned) == 0 { - return true + _, err := put. + Header("content-type", "application/json"). + Header("accept", "application/json"). + Do(context.Background()) + + if err != nil { + return fmt.Errorf("failed to set settings: %s", err) } - idx := 0 - table := NewTable(3, len(orphaned)) - table.Addheaders("orphaned index on follower", "size", "docscount") - - for name, index := range orphaned { - name := Colorize("red", name) - - table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount} - idx++ - } - - table.Sort() - table.PrintMarkdown() - - return false -} - -// find red indices on follower -func findFailedFollowerIndices(conf *cfg.Config, indices ClusterIndices, follower string) bool { - red := map[string]*types.IndicesRecord{} - - for name, index := range indices[follower] { - if *index.Health == "red" { - red[name] = index - } - } - - if len(red) == 0 { - return true - } - - idx := 0 - table := NewTable(3, len(red)) - table.Addheaders("red index on follower", "size", "docscount") - - for name, index := range red { - name := Colorize("red", name) - - table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount} - idx++ - } - - table.Sort() - table.PrintMarkdown() - - return false + return nil } diff --git a/pkg/es/node.go b/pkg/es/node.go new file mode 100644 index 0000000..9f8a094 --- /dev/null +++ b/pkg/es/node.go @@ -0,0 +1,55 @@ +/* +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" + + "codeberg.org/scip/esctl/pkg/cfg" +) + +func NodeList(conf *cfg.Config) error { + // get nodes + nodes, err := conf.DefaultCluster.ES.Cat.Nodes().Do(context.Background()) + if err != nil { + return fmt.Errorf("Error getting nodes: %s", err) + } + + slog.Debug("ES result", "nodes", nodes) + + table := NewTable(7, len(nodes)) + table.Addheaders("name", "ip", "load1m", "load5m", "load15m", "ram %", "heap %") + + for idx, node := range nodes { + table.entries[idx] = []string{ + *node.Name, + *node.Ip, + *node.Load1M, + *node.Load5M, + *node.Load15M, + node.RamPercent.(string), + node.HeapPercent.(string), + } + } + + table.Sort() + table.PrintMarkdown() + + return nil +}