Compare commits

...

2 Commits
0.0.4 ... 0.0.5

Author SHA1 Message Date
a71be5001a add missing utilities.go 2026-05-05 12:07:31 +02:00
T. von Dein
fcfeaf66a9 Add nodes and cluster settings support (#9) 2026-05-05 12:06:34 +02:00
8 changed files with 644 additions and 335 deletions

View File

@@ -12,13 +12,14 @@ USAGE:
esctl [global options] [command [command options]] esctl [global options] [command [command options]]
VERSION: VERSION:
v0.0.3 v0.0.4
COMMANDS: COMMANDS:
search, / search within an index search, / search within an index
index, i manage indicies index, i manage indicies
snapshot, snap manage snapshots snapshot, snap manage snapshots
cluster, c manage cluster[s] cluster, c manage cluster[s]
node, snap manage nodes
help, h Shows a list of commands or help for one command help, h Shows a list of commands or help for one command
GLOBAL OPTIONS: GLOBAL OPTIONS:

View File

@@ -33,21 +33,22 @@ func Cluster(conf *cfg.Config) *cli.Command {
Usage: "manage cluster[s]", Usage: "manage cluster[s]",
Commands: []*cli.Command{ Commands: []*cli.Command{
Compare(conf), ClusterCompare(conf),
Status(conf), ClusterStatus(conf),
List(conf), ClusterList(conf),
ClusterSettings(conf),
}, },
} }
} }
func List(conf *cfg.Config) *cli.Command { func ClusterList(conf *cfg.Config) *cli.Command {
return &cli.Command{ return &cli.Command{
Name: "list", Name: "list",
Usage: "list configured clusters", Usage: "list configured clusters",
Aliases: []string{"ls"}, Aliases: []string{"ls"},
Action: func(ctx context.Context, cmd *cli.Command) error { 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 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{ return &cli.Command{
Name: "status", Name: "status",
Usage: "show cluster 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 { 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 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{ return &cli.Command{
Name: "compare", Name: "compare",
Aliases: []string{"c"}, 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
},
}
}

72
cmd/node.go Normal file
View File

@@ -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 <http://www.gnu.org/licenses/>.
*/
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] <node>",
Action: func(ctx context.Context, cmd *cli.Command) error {
// if err := es.NodeShow(conf, cmd.Args().Get(0)); err != nil {
// return err
// }
return nil
},
}
}

View File

@@ -76,6 +76,7 @@ func Main() int {
Index(conf), Index(conf),
Snapshot(conf), Snapshot(conf),
Cluster(conf), Cluster(conf),
Node(conf),
}, },
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) { Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {

View File

@@ -30,7 +30,7 @@ import (
) )
const ( const (
Version string = `v0.0.4` Version string = `v0.0.5`
) )
type Cluster struct { type Cluster struct {
@@ -52,6 +52,7 @@ type Config struct {
Filter []string // search: -F Filter []string // search: -F
Exclude string // cluster compare: -e (regexp) Exclude string // cluster compare: -e (regexp)
All bool // cluster status: -a All bool // cluster status: -a
Persistent, Transient bool // -p -t cluster settings set
} }
func NewConfig() *Config { func NewConfig() *Config {

View File

@@ -18,15 +18,14 @@ package es
import ( import (
"context" "context"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"log"
"log/slog" "log/slog"
"regexp"
"codeberg.org/scip/esctl/pkg/cfg" "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/elastic/go-elasticsearch/v9/typedapi/types"
"github.com/urfave/cli/v3"
) )
const ( const (
@@ -75,7 +74,7 @@ func ClusterCompare(conf *cfg.Config, leader, follower string) error {
return nil return nil
} }
func List(conf *cfg.Config) error { func ClusterList(conf *cfg.Config) error {
table := NewTable(2, len(conf.Clusters)) table := NewTable(2, len(conf.Clusters))
table.Addheaders("cluster", "uri") table.Addheaders("cluster", "uri")
@@ -91,7 +90,7 @@ func List(conf *cfg.Config) error {
return nil return nil
} }
func Status(conf *cfg.Config) error { func ClusterStatus(conf *cfg.Config) error {
clusters := []string{} clusters := []string{}
if conf.All { if conf.All {
@@ -113,7 +112,7 @@ func Status(conf *cfg.Config) error {
Header("accept", "application/json"). Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { 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) slog.Debug("ES result", "cluster health", res)
@@ -135,340 +134,79 @@ func Status(conf *cfg.Config) error {
return nil return nil
} }
// look for indicies only on leader func ClusterSettingsList(conf *cfg.Config) error {
func findIndicesOnlyOnMaster(conf *cfg.Config, indices ClusterIndices, leader, follower string) { res, err := conf.DefaultCluster.ES.Cluster.GetSettings().
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("content-type", "application/json").
Header("accept", "application/json"). Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
continue // ignore it then return fmt.Errorf("failed to get cluster settings: %s", err)
} }
_, defined := res[name] table := NewTable(2, 0)
if !defined { table.Addheaders("setting", "value")
// json response map didn't contain the index entries := [][]string{}
continue
}
isWritable := false for topic, val := range res.Persistent {
for _, alias := range res[name].Aliases { data := map[string]map[string]any{}
if *alias.IsWriteIndex {
isWritable = true
break
}
}
if isWritable { err := json.Unmarshal(val, &data)
// 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().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil { if err != nil {
fmt.Printf("failed to get ccr stats from %s: %s", leader, err) return fmt.Errorf("failed to unmarshall setting for topic %s: %s", topic, err)
return false
} }
if len(stats.AutoFollowStats.AutoFollowedClusters) == 0 { for key, settings := range data {
// is not following anyone for setting, value := range settings {
return true entries = append(entries, []string{
} fmt.Sprintf("%s.%s.%s", topic, key, setting),
fmt.Sprintf("%v", value),
})
if stats.AutoFollowStats.AutoFollowedClusters[0].ClusterName != "" {
fmt.Println("leader/follower attribution is invalid, reverse cluster attribution and retry")
return false
}
return true
}
// 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
}
status[cluster] = st
}
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.PrintMarkdown()
if status[leader].Status.Name == "green" && status[follower].Status.Name == "green" {
return true
}
return false
}
// finds indices on both clusters which have ilm errors
func findIlmErrors(conf *cfg.Config, leader, follower string) bool {
failed := map[string]map[string]string{}
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
}
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 { table.entries = entries
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())
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
}
}
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.Sort()
table.PrintMarkdown() table.PrintMarkdown()
return false return nil
} }
// find indices only present on follower func ClusterSettingsSet(conf *cfg.Config, args cli.Args) error {
func findOrphanedIndices(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool { put := conf.Clusters[conf.CurrentCluster].ES.Cluster.PutSettings()
orphaned := map[string]*types.IndicesRecord{}
for name, index := range indices[follower] { for _, arg := range args.Slice() {
_, leaderHasIt := indices[leader][name] setting, value := splitArg(arg)
if !leaderHasIt {
// fetch index details switch {
res, err := conf.Clusters[follower].ES.Indices.Get(name). case conf.Transient:
message, err := json.Marshal(value)
if err != nil {
return fmt.Errorf("failed to marshall transient value <%v> to valid JSON: %s", value, err)
}
put.AddTransient(setting, message)
case conf.Persistent:
fallthrough
default:
message, err := json.Marshal(value)
if err != nil {
return fmt.Errorf("failed to marshall persistent value <%v> to valid JSON: %s", value, err)
}
put.AddPersistent(setting, message)
}
}
_, err := put.
Header("content-type", "application/json"). Header("content-type", "application/json").
Header("accept", "application/json"). Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
continue // ignore it then return fmt.Errorf("failed to set settings: %s", err)
} }
_, defined := res[name] return nil
if !defined {
// json response map didn't contain the index
continue
}
orphaned[name] = index
}
}
if len(orphaned) == 0 {
return true
}
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
} }

55
pkg/es/node.go Normal file
View File

@@ -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 <http://www.gnu.org/licenses/>.
*/
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
}

379
pkg/es/utilities.go Normal file
View File

@@ -0,0 +1,379 @@
/*
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 <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"fmt"
"regexp"
"strings"
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/elastic/go-elasticsearch/v9/typedapi/cluster/health"
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
)
// 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().
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
}
if len(stats.AutoFollowStats.AutoFollowedClusters) == 0 {
// is not following anyone
return true
}
if stats.AutoFollowStats.AutoFollowedClusters[0].ClusterName != "" {
fmt.Println("leader/follower attribution is invalid, reverse cluster attribution and retry")
return false
}
return true
}
// 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
}
status[cluster] = st
}
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.PrintMarkdown()
if status[leader].Status.Name == "green" && status[follower].Status.Name == "green" {
return true
}
return false
}
// finds indices on both clusters which have ilm errors
func findIlmErrors(conf *cfg.Config, leader, follower string) bool {
failed := map[string]map[string]string{}
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
}
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())
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
}
}
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())
if err != nil {
continue // ignore it then
}
_, defined := res[name]
if !defined {
// json response map didn't contain the index
continue
}
orphaned[name] = index
}
}
if len(orphaned) == 0 {
return true
}
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
}
func splitArg(arg string) (string, string) {
parts := strings.Split(arg, ":")
switch len(parts) {
case 0:
fallthrough
case 1:
return arg, ""
default:
return parts[0], parts[1]
}
}