Compare commits

..

9 Commits

28 changed files with 1668 additions and 375 deletions

View File

@@ -113,6 +113,7 @@ Configure `esctl` with environment variables:
- `ES_URI`: elasticsearch uri - `ES_URI`: elasticsearch uri
- `ES_USER`: username - `ES_USER`: username
- `ES_PASS`: password - `ES_PASS`: password
- `ES_TOKEN`: API token, instead of user+password
Or create a config file such as this: Or create a config file such as this:
@@ -121,11 +122,10 @@ clusters:
foobar: foobar:
uri: https://es.foo.bar:9200/ uri: https://es.foo.bar:9200/
user: elastic user: elastic
pass: 123456 pass: ******
other: other:
uri: https://myes.foo:9200/ uri: https://myes.foo:9200/
user: elastic token: ******
pass: asdasdasd
``` ```
and specify it with `-c configfile`. You may also put clusters into a and specify it with `-c configfile`. You may also put clusters into a
@@ -433,6 +433,40 @@ $ esctl search -i hyperdrive -l 2 | jq
} }
``` ```
There are also some uniq features which are not directly available via
the ES API or GUI. Here's one example: say you have a number of
indices with ILM policies each. At some day there were a surge in
incoming data in some of them and you want to know, when the excess
data will be rolled over to cold storage:
```console
$ esctl ilm forecast sh -f warm -w 3d
1.1 TB bytes of data in warm phase will be rolled within 72h0m0s to the next phase
```
You may also look at a detailed ilm forecast list:
```console
$ esctl ilm forecast ls -f warm -w 3d
INDEX CURRENT-SIZE CURRENT-AGE VIRTUAL-AGE MIN-AGE MIN-SIZE CURRENT-PHASE NEXT-PHASE
foobar-n1-p01-elastic-000126 12 GB 2d:8h:12m 3d:5h:16m 7d 25 GB hot warm
foobar-n1-p01-misc-000070 16 GB 5d:16h:31m 4d:10h:31m 7d 25 GB hot warm
foobar-n1-q01-elastic-000122 23 GB 5d:0h:11m 6d:10h:33m 7d 25 GB hot warm
foobar-n1-q01-misc-000072 18 GB 6d:15h:21m 4d:21h:36m 7d 25 GB hot warm
delaware-f1-p01-kafka-000023 1.3 MB 15h:12m 0s 7d 25 GB hot warm
delaware-f1-p01-misc-000024 17 GB 6d:16h:1m 4d:16h:33m 7d 25 GB hot warm
delaware-f1-q01-kafka-000023 1.5 MB 15h:12m 0s 7d 25 GB hot warm
delaware-f1-q01-misc-000024 16 GB 6d:15h:31m 4d:13h:52m 7d 25 GB hot warm
delaware-n1-p01-elastic-000417 16 GB 12h:42m 4d:10h:31m 7d 25 GB hot warm
```
So you can see, which index will be rolled over when. Note the
**VIRTUAL-AGE** field however: it is calculated from the current
storage usage of the index in relation to rollover max shard size. So
you can see, when an index will be rolled over either because it aged
out or because its storage exceeded the limit.
---
Please note, that `esctl` is still in its early stages and things are Please note, that `esctl` is still in its early stages and things are
changing heavily every now and then. New commands are being added changing heavily every now and then. New commands are being added
constantly as well. constantly as well.
@@ -464,6 +498,12 @@ cluster - manage cluster[s]
settings - cluster settings management settings - cluster settings management
list - show cluster settings list - show cluster settings
set - set|update cluster settings set - set|update cluster settings
reroute - manually change the allocation of individual shards in the cluster.
move - move shard to another node
allocate-replica - allocate-replica replica to another node
cancel - cancel a reroute operation
allocate-empty-primary - allocate an empty primary shard to a node
allocate-stale-primary - allocate a stale primary shard to a node
datastream - manage data streams datastream - manage data streams
list - list indicies list - list indicies
show - show details about an data stream show - show details about an data stream
@@ -479,15 +519,17 @@ ilm - manage index lifecycle
status - get the current index lifecycle management status status - get the current index lifecycle management status
list - list index lifecycle policies list - list index lifecycle policies
show - show details about an index lifecycle policy show - show details about an index lifecycle policy
create - create a index lifecycle policy create - create a new lifecycle policy
update - update an lifecycle policy
forecast - calculate index phase movements forecast - calculate index phase movements
list - list index rollover config list - list index rollover config
show - show rollover forecast over all indices show - show rollover forecast over all indices
explain - explain ilm condition of an index
index - manage indicies index - manage indicies
list - list indicies list - list indicies
show - show details about an index show - show details about an index
create - create a new index create - create a new index
modify - modify anindex update - update an index
delete - delete an index delete - delete an index
close - close an index close - close an index
allocation - explain index allocation allocation - explain index allocation
@@ -502,11 +544,12 @@ index - manage indicies
list - list index templates list - list index templates
show - show details about an index template show - show details about an index template
create - create a new index template create - create a new index template
modify - modify a new index template update - update a new index template
delete - delete an index template delete - delete an index template
node - manage nodes node - manage nodes
list - list nodes list - list nodes
show - show details about a node show - show details about a node
clients - show node http clients
role - manage roles role - manage roles
list - list roles list - list roles
show - show details about a role show - show details about a role
@@ -524,11 +567,7 @@ task - manage tasks
version - show esctl version information version - show esctl version information
debug - developer only debug - developer only
help-jsonpath - show jsonpath help help-jsonpath - show jsonpath help
completion - Output shell completion script for bash, zsh, fish, or Powershell help-command-overview - show overview of all available commands
pwsh - Output pwsh completion script
bash - Output bash completion script
zsh - Output zsh completion script
fish - Output fish completion script
``` ```
# Development # Development

View File

@@ -37,6 +37,7 @@ func Cluster(conf *cfg.Config) *cli.Command {
ClusterSwitch(conf), ClusterSwitch(conf),
ClusterList(conf), ClusterList(conf),
ClusterSettings(conf), ClusterSettings(conf),
ClusterReroute(conf),
}, },
} }
} }
@@ -60,12 +61,6 @@ func ClusterStatus(conf *cfg.Config) *cli.Command {
Aliases: []string{"s"}, Aliases: []string{"s"},
Flags: []cli.Flag{ Flags: []cli.Flag{
&cli.BoolFlag{
Name: "all",
Usage: "show status of all clusters",
Destination: &conf.All,
Aliases: []string{"a"},
},
&cli.BoolFlag{ &cli.BoolFlag{
Name: "verbose", Name: "verbose",
Usage: "include verbose statistics", Usage: "include verbose statistics",

225
cmd/cluster_reroute.go Normal file
View File

@@ -0,0 +1,225 @@
/*
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"
"errors"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
"github.com/urfave/cli/v3"
)
func ClusterReroute(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "reroute",
Usage: "manually change the allocation of individual shards in the cluster.",
UsageText: "reroute <cluster-name>",
Aliases: []string{"ctx"},
Commands: []*cli.Command{
ClusterRerouteMove(conf),
ClusterRerouteAllocateReplica(conf),
ClusterRerouteCancel(conf),
ClusterRerouteAllocatePrimary(conf, false),
ClusterRerouteAllocatePrimary(conf, true),
},
}
}
func ClusterRerouteMove(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "move",
Usage: "move shard to another node",
UsageText: "move [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex)
},
Flags: []cli.Flag{
&cli.IntFlag{
Name: "shard",
Usage: "shard to move",
Destination: &conf.Shards,
Aliases: []string{"s"},
Required: true,
},
&cli.StringFlag{
Name: "from-node",
Usage: "current node",
Destination: &conf.FromNode,
Aliases: []string{"f"},
Required: true,
},
&cli.StringFlag{
Name: "to-node",
Usage: "node to move to",
Destination: &conf.ToNode,
Aliases: []string{"t"},
Required: true,
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
index := cmd.Args().Get(0)
if index == "" {
return errors.New("no index specified")
}
return es.ClusterRerouteMove(conf, index)
},
}
}
func ClusterRerouteAllocateReplica(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "allocate-replica",
Usage: "allocate-replica replica to another node",
UsageText: "allocate-replica [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex)
},
Flags: []cli.Flag{
&cli.IntFlag{
Name: "shard",
Usage: "shard to move",
Destination: &conf.Shards,
Aliases: []string{"s"},
Required: true,
},
&cli.StringFlag{
Name: "to-node",
Usage: "node to move to",
Destination: &conf.ToNode,
Aliases: []string{"t"},
Required: true,
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
index := cmd.Args().Get(0)
if index == "" {
return errors.New("no index specified")
}
return es.ClusterRerouteAllocateReplica(conf, index)
},
}
}
func ClusterRerouteCancel(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "cancel",
Usage: "cancel a reroute operation",
UsageText: "cancel [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex)
},
Flags: []cli.Flag{
&cli.IntFlag{
Name: "shard",
Usage: "shard to move",
Destination: &conf.Shards,
Aliases: []string{"s"},
Required: true,
},
&cli.StringFlag{
Name: "to-node",
Usage: "node to move to",
Destination: &conf.ToNode,
Aliases: []string{"t"},
Required: true,
},
&cli.BoolFlag{
Name: "allow-primary",
Usage: "allow primary shard to cancel",
Destination: &conf.AllowPrimary,
Aliases: []string{"a"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
index := cmd.Args().Get(0)
if index == "" {
return errors.New("no index specified")
}
return es.ClusterRerouteCancel(conf, index)
},
}
}
func ClusterRerouteAllocatePrimary(conf *cfg.Config, stale bool) *cli.Command {
name := "allocate-empty-primary"
usage := "allocate an empty primary shard to a node"
if stale {
name = "allocate-stale-primary"
usage = "allocate a stale primary shard to a node"
}
return &cli.Command{
Name: name,
Usage: usage,
UsageText: name + " [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex)
},
Flags: []cli.Flag{
&cli.IntFlag{
Name: "shard",
Usage: "shard to move",
Destination: &conf.Shards,
Aliases: []string{"s"},
Required: true,
},
&cli.StringFlag{
Name: "to-node",
Usage: "node to move to",
Destination: &conf.ToNode,
Aliases: []string{"t"},
Required: true,
},
&cli.BoolFlag{
Name: "accept-data-loss",
Usage: "",
Destination: &conf.AcceptDataLoss,
Aliases: []string{"a"},
Required: true,
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
index := cmd.Args().Get(0)
if index == "" {
return errors.New("no index specified")
}
return es.ClusterRerouteAllocatePrimary(conf, index, stale)
},
}
}

View File

@@ -32,6 +32,7 @@ const (
Ccluster Ccluster
Capi Capi
Cilm Cilm
Cnode
) )
func complete(cmd *cli.Command, what int) { func complete(cmd *cli.Command, what int) {
@@ -64,6 +65,8 @@ func complete(cmd *cli.Command, what int) {
list = es.ApiPathNames() list = es.ApiPathNames()
case Cilm: case Cilm:
list, err = es.IlmNames(conf) list, err = es.IlmNames(conf)
case Cnode:
list, err = es.NodeNames(conf)
} }
if err != nil { if err != nil {

View File

@@ -36,7 +36,8 @@ func Ilm(conf *cfg.Config) *cli.Command {
IlmStatus(conf), IlmStatus(conf),
IlmList(conf), IlmList(conf),
IlmShow(conf), IlmShow(conf),
IlmCreate(conf), IlmCreate(conf, false),
IlmCreate(conf, true),
IlmForecast(conf), IlmForecast(conf),
IlmExplain(conf), IlmExplain(conf),
}, },
@@ -103,6 +104,15 @@ func IlmShow(conf *cfg.Config) *cli.Command {
Usage: "show details about an index lifecycle policy", Usage: "show details about an index lifecycle policy",
UsageText: "show <policy>", UsageText: "show <policy>",
Flags: []cli.Flag{
&cli.BoolFlag{
Name: "tree",
Usage: "display policy as a tree",
Destination: &conf.Ilm.Tree,
Aliases: []string{"t"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
policy := cmd.Args().Get(0) policy := cmd.Args().Get(0)
if policy == "" { if policy == "" {
@@ -118,12 +128,22 @@ func IlmShow(conf *cfg.Config) *cli.Command {
} }
} }
func IlmCreate(conf *cfg.Config) *cli.Command { func IlmCreate(conf *cfg.Config, modify bool) *cli.Command {
name := "create"
alias := "+"
usage := "create a new lifecycle policy"
if modify {
name = "update"
usage = "update an lifecycle policy"
alias = "upd"
}
return &cli.Command{ return &cli.Command{
Name: "create", Name: name,
Aliases: []string{"+"}, Aliases: []string{alias},
Usage: "create a index lifecycle policy", Usage: usage,
UsageText: "create [options] <policy>", UsageText: name + " [options] <policy>",
Flags: []cli.Flag{ Flags: []cli.Flag{
&cli.StringFlag{ &cli.StringFlag{

View File

@@ -154,8 +154,8 @@ func IndexCreate(conf *cfg.Config, modify bool) *cli.Command {
usage := "create a new index" usage := "create a new index"
if modify { if modify {
name = "modify" name = "update"
usage = "modify anindex" usage = "update an index"
} }
return &cli.Command{ return &cli.Command{

View File

@@ -83,8 +83,8 @@ func IndexTemplateCreate(conf *cfg.Config, modify bool) *cli.Command {
required := true required := true
if modify { if modify {
name = "modify" name = "update"
alias = "mod" alias = "upd"
required = false required = false
} }

View File

@@ -18,6 +18,7 @@ package cmd
import ( import (
"context" "context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es" "codeberg.org/scip/esctl/pkg/es"
@@ -34,6 +35,7 @@ func Node(conf *cfg.Config) *cli.Command {
Commands: []*cli.Command{ Commands: []*cli.Command{
NodeList(conf), NodeList(conf),
NodeShow(conf), NodeShow(conf),
NodeClients(conf),
}, },
} }
} }
@@ -57,10 +59,47 @@ func NodeShow(conf *cfg.Config) *cli.Command {
Usage: "show details about a node", Usage: "show details about a node",
UsageText: "show [options] <node>", UsageText: "show [options] <node>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cnode)
},
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
// FIXME: implement es.NodeShow() node := cmd.Args().Get(0)
// return es.NodeShow(conf, cmd.Args().Get(0)) if node == "" {
return nil return errors.New("no node specified")
}
return es.NodeShow(conf, node)
},
}
}
func NodeClients(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "clients",
Usage: "show node http clients",
UsageText: "clients [options] <node>",
Flags: []cli.Flag{
&cli.BoolFlag{
Name: "query",
Usage: "include query parameters",
Destination: &conf.All,
Aliases: []string{"q"},
},
},
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cnode)
},
Action: func(ctx context.Context, cmd *cli.Command) error {
node := cmd.Args().Get(0)
if node == "" {
return errors.New("no node specified")
}
return es.NodeClients(conf, node)
}, },
} }
} }

View File

@@ -42,7 +42,6 @@ func Finish(err error) int {
func Main() int { func Main() int {
conf := cfg.NewConfig() conf := cfg.NewConfig()
tree := false
cmd := &cli.Command{ cmd := &cli.Command{
Name: "esctl", Name: "esctl",
@@ -64,13 +63,6 @@ func Main() int {
Usage: "enable HTTP debugging", Usage: "enable HTTP debugging",
Destination: &conf.DebugHTTP, Destination: &conf.DebugHTTP,
}, },
&cli.BoolFlag{
Name: "show-command-tree",
Value: false,
Usage: "generate a command tree",
Destination: &tree,
Hidden: true,
},
&cli.BoolFlag{ &cli.BoolFlag{
Name: "align-ints", Name: "align-ints",
Aliases: []string{"I"}, Aliases: []string{"I"},
@@ -125,17 +117,10 @@ func Main() int {
Version(conf), Version(conf),
Debug(conf), Debug(conf),
HelpJsonPath(conf), HelpJsonPath(conf),
HelpUsage(conf),
}, },
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) { Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {
if tree {
if err := Tree(cmd); err != nil {
return nil, err
}
os.Exit(0)
}
if err := conf.Init(); err != nil { if err := conf.Init(); err != nil {
if len(os.Args) > 1 { if len(os.Args) > 1 {
return nil, err return nil, err
@@ -261,10 +246,42 @@ func Debug(conf *cfg.Config) *cli.Command {
} }
} }
func Tree(cmd *cli.Command) error { func HelpUsage(conf *cfg.Config) *cli.Command {
max := 20 return &cli.Command{
Name: "help-command-overview",
Usage: "show overview of all available commands",
Aliases: []string{"usage"},
return cmd.Walk(func(cmd *cli.Command) error { Flags: []cli.Flag{
&cli.BoolFlag{
Name: "hidden",
Usage: "include hidden commands",
Destination: &conf.Hidden,
Aliases: []string{"H"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
maxCommandWidth := 0
// first pass, determine max command width
if err := walkVisible(conf, cmd.Root(), func(cmd *cli.Command) error {
path := cmd.Path()
size := len(path[len(path)-1])
if size > maxCommandWidth {
maxCommandWidth = size
}
return nil
}); err != nil {
return err
}
maxCommandWidth += 4 // account for indent width
// second pass, build tree
return walkVisible(conf, cmd.Root(), func(cmd *cli.Command) error {
path := cmd.Path() path := cmd.Path()
command := path[len(path)-1] command := path[len(path)-1]
@@ -273,10 +290,37 @@ func Tree(cmd *cli.Command) error {
} }
indent := strings.Repeat(" ", len(path[1:])-1) indent := strings.Repeat(" ", len(path[1:])-1)
space := strings.Repeat(" ", max-(len(command)+len(indent))) space := strings.Repeat(" ", maxCommandWidth-(len(command)+len(indent)))
fmt.Printf("%s%s %s - %s\n", indent, command, space, cmd.Usage) fmt.Printf("%s%s %s - %s\n", indent, command, space, cmd.Usage)
return nil return nil
}) })
},
}
}
// copy of cmd.Walk() with the exception to skip hidden commands and its siblings
// see: https://github.com/urfave/cli/issues/2372
func walkVisible(conf *cfg.Config, cmd *cli.Command, fn func(*cli.Command) error) error {
if fn == nil {
return nil
}
if !conf.Hidden && cmd.Hidden {
return nil
}
if err := fn(cmd); err != nil {
return err
}
for _, sub := range cmd.Commands {
if err := walkVisible(conf, sub, fn); err != nil {
return err
}
}
return nil
} }

1
go.mod
View File

@@ -85,6 +85,7 @@ require (
github.com/olekukonko/errors v1.2.0 // indirect github.com/olekukonko/errors v1.2.0 // indirect
github.com/olekukonko/ll v0.1.8 // indirect github.com/olekukonko/ll v0.1.8 // indirect
github.com/rivo/uniseg v0.4.7 // indirect github.com/rivo/uniseg v0.4.7 // indirect
github.com/seeruk/go-wordwrap v0.0.0-20191208221741-14ec4aac9550 // indirect
github.com/tidwall/match v1.1.1 // indirect github.com/tidwall/match v1.1.1 // indirect
github.com/tidwall/pretty v1.2.0 // indirect github.com/tidwall/pretty v1.2.0 // indirect
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect

2
go.sum
View File

@@ -171,6 +171,8 @@ github.com/rivo/uniseg v0.4.7 h1:WUdvkW8uEhrYfLC4ZzdpI2ztxP1I582+49Oc5Mq64VQ=
github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88= github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88=
github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII= github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII=
github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o= github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o=
github.com/seeruk/go-wordwrap v0.0.0-20191208221741-14ec4aac9550 h1:C3CfUXH/qmWuQFRqnPm3Sx8PFxa+pqACjhV5CaNO8pw=
github.com/seeruk/go-wordwrap v0.0.0-20191208221741-14ec4aac9550/go.mod h1:Sl541M2Em6rRG3V9WObycR7MYFZiERVkd/TJg0Gt0U4=
github.com/sergi/go-diff v1.0.0 h1:Kpca3qRNrduNnOQeazBd0ysaKrUJiIuISHxogkT9RPQ= github.com/sergi/go-diff v1.0.0 h1:Kpca3qRNrduNnOQeazBd0ysaKrUJiIuISHxogkT9RPQ=
github.com/sergi/go-diff v1.0.0/go.mod h1:0CfEIISq7TuYL3j771MWULgwwjU+GofnZX9QAmXWZgo= github.com/sergi/go-diff v1.0.0/go.mod h1:0CfEIISq7TuYL3j771MWULgwwjU+GofnZX9QAmXWZgo=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=

174
pkg/cfg/automate.go Normal file
View File

@@ -0,0 +1,174 @@
/*
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 cfg
import (
"bytes"
"context"
"errors"
"fmt"
"log/slog"
"os"
"os/exec"
"path/filepath"
"strings"
"time"
"gopkg.in/yaml.v3"
)
type Automator struct {
IsReachable bool
Error error
Env []string
Cluster *Cluster
// all bash commands, .name must succeed first
Name string `yaml:"name"`
User string `yaml:"user"`
Pass string `yaml:"pass"`
Uri string `yaml:"uri"`
}
type AutomatorRunner struct {
Output string
Error error
}
func NewAutoEnv() []string {
return []string{
"PATH=" + os.Getenv("PATH"),
"HOME=" + os.Getenv("HOME"),
"SHELL=" + os.Getenv("SHELL"),
"KUBECONFIG=" + os.Getenv("KUBECONFIG"),
}
}
func NewAutomator() Automator {
autocfg := filepath.Join([]string{os.Getenv("HOME"), ".config", "esctl", "automate.yaml"}...)
auto := Automator{Env: NewAutoEnv()}
if !fileExists(autocfg) {
return auto
}
data, err := os.ReadFile(autocfg)
if err != nil {
auto.Error = fmt.Errorf("failed to read config file: %w", err)
return auto
}
err = yaml.Unmarshal(data, &auto)
if err != nil {
auto.Error = fmt.Errorf("failed to unmarshal config file: %w", err)
return auto
}
hasname := execute(auto.Name, auto.Env)
if hasname.Error != nil {
auto.Error = hasname.Error
return auto
}
auto.Name = hasname.Output
hasuser := execute(auto.User, auto.Env)
if hasuser.Error != nil {
auto.Error = hasuser.Error
return auto
}
auto.User = hasuser.Output
haspass := execute(auto.Pass, auto.Env)
if haspass.Error != nil {
auto.Error = haspass.Error
return auto
}
auto.Pass = haspass.Output
hasuri := execute(auto.Uri, auto.Env)
if hasuri.Error != nil {
auto.Error = hasuri.Error
return auto
}
auto.Uri = hasuri.Output
auto.IsReachable = true
auto.Cluster = &Cluster{
Name: auto.Name,
Uri: auto.Uri,
User: auto.User,
Pass: auto.Pass,
}
return auto
}
// FIXME: execute auto.Exec, if defined, in a go routine and let it run forever, might be a tunnel
func (auto *Automator) Exec() {}
// FIXME: cache automator results, only check auto.Name and if it matches the cache use those vars, but run auto.Exec anyway
func (auto *Automator) Cache() {}
func execute(code string, env []string) *AutomatorRunner {
timeoutCtx, cancel := context.WithTimeout(context.Background(),
time.Duration(10)*time.Second)
defer cancel()
var cmd *exec.Cmd
// pipe code into bash
cmd = exec.CommandContext(timeoutCtx, "bash")
cmd.Stdin = strings.NewReader(code)
cmd.Env = env
errbuf := &bytes.Buffer{}
cmd.Stderr = errbuf
done := make(chan bool)
out := AutomatorRunner{}
go func() {
output, err := cmd.Output()
out.Output = strings.TrimSpace(string(output))
switch {
case err != nil:
out.Error = err
case errbuf.Len() > 0:
out.Error = fmt.Errorf(errbuf.String())
case timeoutCtx.Err() == context.DeadlineExceeded:
out.Error = errors.New("timed out")
}
slog.Debug("executed automator",
"code", code,
"output", out.Output,
"error", out.Error)
done <- true
}()
for {
select {
case <-done:
return &out
}
}
}

View File

@@ -17,26 +17,32 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
package cfg package cfg
import ( import (
"context"
"crypto/tls"
"errors" "errors"
"fmt" "fmt"
"net/http" "net/http"
"os" "os"
"strings"
"syscall"
"time"
"github.com/elastic/elastic-transport-go/v8/elastictransport" "github.com/elastic/elastic-transport-go/v8/elastictransport"
"github.com/elastic/go-elasticsearch/v9" "github.com/elastic/go-elasticsearch/v9"
"golang.org/x/term"
"gopkg.in/yaml.v3" "gopkg.in/yaml.v3"
) )
// used in general config struct // used in general config struct
type Cluster struct { type Cluster struct {
Uri, User, Pass string Name, Uri, User, Pass, Token string
client *elasticsearch.TypedClient client *elasticsearch.TypedClient
Default bool Default, DebugHTTP bool
} }
// used just for writing back to the config file // used just for writing back to the config file
type ClusterConfig struct { type ClusterConfig struct {
Uri, User, Pass string Uri, User, Pass, Token string
Default bool Default bool
} }
@@ -45,17 +51,96 @@ type WriteConfig struct {
Clusters map[string]*ClusterConfig Clusters map[string]*ClusterConfig
} }
func (cluster *Cluster) SetClient(client *elasticsearch.TypedClient) {
cluster.client = client
}
func (cluster *Cluster) getTransport() elastictransport.Option {
transport := &http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
}
if cluster.DebugHTTP {
return elastictransport.WithTransport(
&DebugTransport{Transport: transport},
)
}
return elastictransport.WithTransport(transport)
}
func (cluster *Cluster) getDefaultOptions() []elasticsearch.Option {
// These headers are not needed with ES 9, but with ES 8, we set
// them here so every API call uses it. The only exception being
// the api repl, which does it on its own.
headers := http.Header{}
headers.Add("content-type", "application/json")
headers.Add("Accept", "application/json")
return []elasticsearch.Option{
elasticsearch.WithAddresses(cluster.Uri),
elasticsearch.WithTransportOptions(
cluster.getTransport(),
elastictransport.WithHeader(headers),
),
}
}
// return the go-elasticsearch client object but before doing that,
// check if we need to tune auth
func (cluster *Cluster) ES() *elasticsearch.TypedClient { func (cluster *Cluster) ES() *elasticsearch.TypedClient {
if cluster.client == nil { if cluster.client == nil {
fmt.Println("no current cluster, use 'esctl cluster switch <name>' to set one") fmt.Println("no current cluster, use 'esctl cluster switch <name>' to set one")
os.Exit(1) os.Exit(1)
} }
if err := cluster.CheckAuth(); err != nil {
fmt.Printf("Error: %s", err)
os.Exit(1)
}
return cluster.client return cluster.client
} }
func (cluster *Cluster) SetClient(client *elasticsearch.TypedClient) { // add authentication to es client, if not yet done
cluster.client = client func (cluster *Cluster) CheckAuth() error {
if cluster.Pass == "" && cluster.User != "" && cluster.Token == "" && cluster.Default {
// no token - user is set, but no password.
// check if the env var is set
pass := os.Getenv("ES_PASS")
if pass != "" {
cluster.Pass = pass
} else {
// k, try interactively
fmt.Printf("Enter password for elasticsearch user %s@%s: ", cluster.User, cluster.Name)
pass, err := term.ReadPassword(int(syscall.Stdin))
if err != nil {
return err
}
passwd := strings.TrimSpace(string(pass))
if passwd == "" {
return errors.New("password empty")
}
cluster.Pass = string(pass)
fmt.Println()
}
opts := cluster.getDefaultOptions()
opts = append(opts, elasticsearch.WithBasicAuth(cluster.User, cluster.Pass))
es, err := elasticsearch.NewTyped(opts...)
if err != nil {
return fmt.Errorf("failed to setup elasticsearch connection: %w", err)
}
cluster.SetClient(es)
}
return nil
} }
// set Default=true for the given cluster in the config (if exists) // set Default=true for the given cluster in the config (if exists)
@@ -72,6 +157,7 @@ func (conf *Config) SwitchCluster(name string) error {
Uri: cluster.Uri, Uri: cluster.Uri,
User: cluster.User, User: cluster.User,
Pass: cluster.Pass, Pass: cluster.Pass,
Token: cluster.Token,
Default: false, Default: false,
} }
@@ -97,29 +183,75 @@ func (conf *Config) SwitchCluster(name string) error {
return nil return nil
} }
func (conf *Config) SetupES() error { // We do NOT use go-elasticsearch to check for cluster reachability,
// These headers are not needed with ES 9, but with ES 8, we set // because at this stage, auth may not have been configured. So
// them here so every API call uses it. The only exception being // instead we just connect to the cluster using plan net/http, ignore
// the api repl, which does it on its own. // HTTP response status and return true if we could just reach ith
headers := http.Header{} func (cluster *Cluster) IsReachable() (bool, error) {
headers.Add("content-type", "application/json") ctx, cancel := context.WithTimeout(
headers.Add("Accept", "application/json") context.Background(),
time.Duration(500)*time.Millisecond)
defer cancel()
for _, cluster := range conf.Clusters { req, err := http.NewRequestWithContext(
es, err := elasticsearch.NewTyped( ctx,
elasticsearch.WithAddresses(cluster.Uri), "GET",
elasticsearch.WithBasicAuth(cluster.User, cluster.Pass), cluster.Uri,
elasticsearch.WithTransportOptions( nil,
conf.getTransport(),
elastictransport.WithHeader(headers),
),
) )
if err != nil {
return false, err
}
client := &http.Client{Transport: &http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
}}
resp, err := client.Do(req)
if err != nil {
return false, err
}
if resp != nil {
// at this stage we do not care if the elasticsearch cluster
// accepts our request or if it's misconfigured in some way
return true, nil
}
return false, nil
}
func (conf *Config) SetupElasticClient(name string, cluster *Cluster) error {
cluster.Name = name
cluster.DebugHTTP = conf.DebugHTTP
opts := cluster.getDefaultOptions()
switch {
case cluster.Pass != "" && cluster.User != "":
opts = append(opts, elasticsearch.WithBasicAuth(cluster.User, cluster.Pass))
case cluster.Token != "":
opts = append(opts, elasticsearch.WithAPIKey(cluster.Token))
}
es, err := elasticsearch.NewTyped(opts...)
if err != nil { if err != nil {
return fmt.Errorf("failed to setup elasticsearch connection: %w", err) return fmt.Errorf("failed to setup elasticsearch connection: %w", err)
} }
cluster.SetClient(es) cluster.SetClient(es)
return nil
}
func (conf *Config) SetupElasticClients() error {
for name, cluster := range conf.Clusters {
if err := conf.SetupElasticClient(name, cluster); err != nil {
return err
}
} }
return nil return nil

View File

@@ -28,7 +28,7 @@ import (
) )
const ( const (
Version string = `v0.0.22` Version string = `v0.0.25`
) )
var ( var (
@@ -103,6 +103,9 @@ type Config struct {
Tag string // api ls: -t Tag string // api ls: -t
Ilm Ilm // ilm create Ilm Ilm // ilm create
FromNode, ToNode string // cluster reroute move: -f + -t
AllowPrimary, AcceptDataLoss bool // cluster reroute cancel: -p,-a
} }
func NewConfig() *Config { func NewConfig() *Config {
@@ -130,7 +133,7 @@ func (conf *Config) Init() error {
} }
} }
if err := conf.SetupES(); err != nil { if err := conf.SetupElasticClients(); err != nil {
return err return err
} }
@@ -150,6 +153,9 @@ func (conf *Config) Init() error {
conf.DefaultCluster.Default = true conf.DefaultCluster.Default = true
} }
} else { } else {
// load auto config, if any
auto := NewAutomator()
// we need to determine ourselfes // we need to determine ourselfes
if len(conf.Clusters) == 1 { if len(conf.Clusters) == 1 {
// ok, just one cluster configured, use this, of course // ok, just one cluster configured, use this, of course
@@ -158,6 +164,14 @@ func (conf *Config) Init() error {
conf.CurrentCluster = name conf.CurrentCluster = name
conf.DefaultCluster.Default = true conf.DefaultCluster.Default = true
} }
} else if auto.Error == nil && auto.IsReachable {
// use auto conf
if err := conf.SetupElasticClient(auto.Name, auto.Cluster); err != nil {
return err
}
conf.DefaultCluster = auto.Cluster
conf.CurrentCluster = auto.Name
} else { } else {
// multiple ones exists, look if one is set as default // multiple ones exists, look if one is set as default
for name, cluster := range conf.Clusters { for name, cluster := range conf.Clusters {
@@ -213,15 +227,14 @@ func (conf *Config) LoadEnv() error {
Uri: os.Getenv("ES_URI"), Uri: os.Getenv("ES_URI"),
User: os.Getenv("ES_USER"), User: os.Getenv("ES_USER"),
Pass: os.Getenv("ES_PASS"), Pass: os.Getenv("ES_PASS"),
Token: os.Getenv("ES_TOKEN"),
} }
switch { switch {
case cluster.Uri == "": case cluster.Uri == "":
return errors.New("ES_URI unset") return errors.New("ES_URI unset")
case cluster.User == "": case cluster.User == "" || cluster.Token == "":
return errors.New("ES_USER unset") return errors.New("ES_USER and ES_TOKEN unset")
case cluster.Pass == "":
return errors.New("ES_PASS unset")
} }
conf.Clusters["default"] = &cluster conf.Clusters["default"] = &cluster

View File

@@ -42,6 +42,7 @@ type Ilm struct {
DeleteSearchableSnapshots bool DeleteSearchableSnapshots bool
MinAge, FromPhase, Within string // forecast MinAge, FromPhase, Within string // forecast
Tree bool
} }
func (cfg *Ilm) HaveHot() bool { func (cfg *Ilm) HaveHot() bool {

38
pkg/cfg/term.go Normal file
View File

@@ -0,0 +1,38 @@
/*
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 cfg
import (
"os"
"golang.org/x/term"
)
const (
DefaultMargin = 4
)
func GetTermWidth() int {
if term.IsTerminal(int(os.Stdout.Fd())) {
width, _, err := term.GetSize(int(os.Stdout.Fd()))
if err == nil {
return width - DefaultMargin
}
}
return 80
}

View File

@@ -18,13 +18,10 @@ package cfg
import ( import (
"bytes" "bytes"
"crypto/tls"
"encoding/json" "encoding/json"
"fmt" "fmt"
"log/slog" "log/slog"
"net/http" "net/http"
"github.com/elastic/elastic-transport-go/v8/elastictransport"
) )
// used to print uri, path and body of a request made by the go-client // used to print uri, path and body of a request made by the go-client
@@ -62,17 +59,3 @@ func (t *DebugTransport) RoundTrip(req *http.Request) (*http.Response, error) {
return t.Transport.RoundTrip(req) return t.Transport.RoundTrip(req)
} }
func (conf *Config) getTransport() elastictransport.Option {
transport := &http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
}
if conf.DebugHTTP {
return elastictransport.WithTransport(
&DebugTransport{Transport: transport},
)
}
return elastictransport.WithTransport(transport)
}

View File

@@ -41,11 +41,9 @@ import (
"github.com/charmbracelet/lipgloss" "github.com/charmbracelet/lipgloss"
"github.com/chzyer/readline" "github.com/chzyer/readline"
"github.com/go-openapi/spec" "github.com/go-openapi/spec"
"golang.org/x/term"
) )
const ( const (
DefaultMargin = 4
intro = `Input format: verb path [data]" intro = `Input format: verb path [data]"
Example: Example:
@@ -172,7 +170,17 @@ func CallAPI(conf *cfg.Config, verb, path, data string) ([]byte, error) {
req.Header.Add("Content-Type", "application/json") req.Header.Add("Content-Type", "application/json")
req.Header.Add("accept", "application/json") req.Header.Add("accept", "application/json")
// make sure we have got all we need
if err := conf.DefaultCluster.CheckAuth(); err != nil {
return nil, err
}
if conf.DefaultCluster.Token != "" {
req.Header.Add("Authorization", "APIKey "+conf.DefaultCluster.Token)
} else {
req.Header.Add("Authorization", "Basic "+encodeAuth(conf.DefaultCluster.User, conf.DefaultCluster.Pass)) req.Header.Add("Authorization", "Basic "+encodeAuth(conf.DefaultCluster.User, conf.DefaultCluster.Pass))
}
// actually execute the request // actually execute the request
resp, err := client.Do(req) resp, err := client.Do(req)
@@ -318,19 +326,19 @@ func ApiShow(conf *cfg.Config, showpath, verb string) error {
cleanMarkup := regexp.MustCompile(`<[^<>]+>`) cleanMarkup := regexp.MustCompile(`<[^<>]+>`)
width := getTermWidth() width := cfg.GetTermWidth()
params := getApiParameters(op, showpath, width) params := getApiParameters(op, showpath, width)
sample := getApiExample(op) sample := getApiExample(op)
description := markdown.Render(cleanMarkup.ReplaceAllString(op.Op.Description, ""), width, DefaultMargin) description := markdown.Render(cleanMarkup.ReplaceAllString(op.Op.Description, ""), width, cfg.DefaultMargin)
var bold = lipgloss.NewStyle(). var bold = lipgloss.NewStyle().
Bold(true) Bold(true)
var paragraph = lipgloss.NewStyle(). var paragraph = lipgloss.NewStyle().
MarginBottom(1). MarginBottom(1).
MarginLeft(DefaultMargin) MarginLeft(cfg.DefaultMargin)
var boldparagraph = lipgloss.NewStyle(). var boldparagraph = lipgloss.NewStyle().
MarginBottom(1). MarginBottom(1).
MarginLeft(DefaultMargin). MarginLeft(cfg.DefaultMargin).
Bold(true) Bold(true)
var indentparagraph = lipgloss.NewStyle(). var indentparagraph = lipgloss.NewStyle().
MarginBottom(1). MarginBottom(1).
@@ -413,7 +421,7 @@ func getApiParameters(op *Op, path string, width int) Params {
} }
} else { } else {
par := Param{ par := Param{
Description: string(markdown.Render(param.Description, width, DefaultMargin)), Description: string(markdown.Render(param.Description, width, cfg.DefaultMargin)),
Param: param.Name, Param: param.Name,
} }
@@ -517,14 +525,3 @@ func findOperation(item spec.PathItemProps) []*Op {
return ops return ops
} }
func getTermWidth() int {
if term.IsTerminal(int(os.Stdout.Fd())) {
width, _, err := term.GetSize(int(os.Stdout.Fd()))
if err == nil {
return width - DefaultMargin
}
}
return 80
}

View File

@@ -0,0 +1,59 @@
/*
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 (
"encoding/json"
"fmt"
"codeberg.org/scip/esctl/pkg/cfg"
)
type HealthReport struct {
Indicators map[string]HealthReportIndicator
}
type HealthReportIndicator struct {
Status, Symptom string
Diagnosis []HealthReportDiagnosis
}
type HealthReportDiagnosis struct {
Id, Cause, Action string
AffectedResources map[string][]string `json:"affected_resources"`
}
// we do not use the go-elasticsearch client API here but call the ES
// API directly, because the returned structure (a
// healthreport.Response) is not iterable, you'd have to explicitly
// call every indicator type and every cause etc which also have
// different types each. To check which is !green would result in a
// gigantic function.
func getHealthReport(conf *cfg.Config) (*HealthReport, error) {
raw, err := CallAPI(conf, "GET", "/_health_report", "")
if err != nil {
return nil, err
}
report := HealthReport{}
if err := json.Unmarshal(raw, &report); err != nil {
return nil, fmt.Errorf("failed to unmarshal healthreport response: %w", err)
}
return &report, nil
}

View File

@@ -17,7 +17,6 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
package es package es
import ( import (
"context"
"fmt" "fmt"
"log/slog" "log/slog"
"strings" "strings"
@@ -26,44 +25,40 @@ import (
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer" "codeberg.org/scip/esctl/pkg/printer"
"github.com/dustin/go-humanize" "github.com/dustin/go-humanize"
"github.com/elastic/go-elasticsearch/v9/typedapi/cat/indices"
"github.com/elastic/go-elasticsearch/v9/typedapi/cat/tasks"
"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" 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" "github.com/elastic/go-elasticsearch/v9/typedapi/types"
) )
type ClusterIndices map[string]map[string]*types.IndicesRecord type ClusterIndices map[string]map[string]*types.IndicesRecord
func ClusterList(conf *cfg.Config) error { func ClusterList(conf *cfg.Config) error {
table := printer.NewTable(conf, 4, len(conf.Clusters)) table := printer.NewTable(conf, 5, len(conf.Clusters))
table.Addheaders("cluster", "uri", "reachable", "current") table.Addheaders("cluster", "uri", "reachable", "current", "error")
idx := 0 idx := 0
for name, cluster := range conf.Clusters { for name, cluster := range conf.Clusters {
reachable := "no" reachable := "no"
current := "no" current := "no"
errmsg := ""
_, err := cluster.ES().Cluster.Health(). online, err := cluster.IsReachable()
Do(context.Background())
if err == nil { if online {
reachable = printer.Colorize(conf, "green", "reachable") reachable = printer.Colorize(conf, "green", "reachable")
} }
if cluster.Default { if cluster.Default {
current = printer.Colorize(conf, "green", "yes") current = printer.Colorize(conf, "green", "yes")
if err != nil { if !online {
reachable = printer.Colorize(conf, "red", "no") reachable = printer.Colorize(conf, "red", "no")
errmsg = err.Error()
} }
} }
table.Entries[idx] = []string{name, cluster.Uri, reachable, current} table.Entries[idx] = []string{name, cluster.Uri, reachable, current, errmsg}
idx++ idx++
} }
@@ -78,90 +73,84 @@ func ClusterList(conf *cfg.Config) error {
// We're using goroutines here to parallelize API requests, since we // We're using goroutines here to parallelize API requests, since we
// have to do 3 of'em for each cluster. This speeds things up. // have to do 3 of'em for each cluster. This speeds things up.
func ClusterStatus(conf *cfg.Config) error { func getClusterStatus(conf *cfg.Config) (*apiResponse, error) {
clusters := []string{} gocount := 6
gocount := 5
if conf.Verbose { if conf.Verbose {
gocount++ 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() es := conf.DefaultCluster.ES()
if cluster != "default" {
es = conf.Clusters[cluster].ES()
}
responses := make(chan apiResponse, gocount) responses := make(chan apiResponse, gocount)
wg := &sync.WaitGroup{} wg := &sync.WaitGroup{}
wg.Add(gocount) wg.Add(gocount)
go getApiData(es, wg, responses, "health") go getApiData(conf, es, wg, responses, "health")
go getApiData(es, wg, responses, "info") go getApiData(conf, es, wg, responses, "healthreport")
go getApiData(es, wg, responses, "ccrstats") go getApiData(conf, es, wg, responses, "info")
go getApiData(es, wg, responses, "indices") go getApiData(conf, es, wg, responses, "ccr")
go getApiData(es, wg, responses, "tasks") go getApiData(conf, es, wg, responses, "indices")
go getApiData(conf, es, wg, responses, "tasks")
if conf.Verbose { if conf.Verbose {
go getApiData(es, wg, responses, "stats") go getApiData(conf, es, wg, responses, "stats")
} }
wg.Wait() wg.Wait()
var clusterhealth *health.Response all := apiResponse{}
var info *info.Response
var ccrstats *stats.Response
var clusterstats *clusterstats.Response
var indexstats *indices.Response
var taskstatus *tasks.Response
for i := 0; i < gocount; i++ { for i := 0; i < gocount; i++ {
r := <-responses r := <-responses
if r.error != nil { if r.error != nil {
return r.error return nil, r.error
} }
switch r.which { switch r.which {
case ResponseHealth: case ResponseHealth:
clusterhealth = r.health all.health = r.health
case ResponseCcr: case ResponseCcr:
ccrstats = r.ccr all.ccr = r.ccr
case ResponseInfo: case ResponseInfo:
info = r.info all.info = r.info
case ResponseStats: case ResponseStats:
clusterstats = r.stats all.stats = r.stats
case ResponseIndices: case ResponseIndices:
indexstats = r.indices all.indices = r.indices
case ResponseTasks: case ResponseTasks:
taskstatus = r.tasks all.tasks = r.tasks
case ResponseHealthReport:
all.healthreport = r.healthreport
} }
} }
slog.Debug("ES result", "cluster health", clusterhealth) return &all, nil
}
isleader := len(ccrstats.AutoFollowStats.AutoFollowedClusters) == 0 func ClusterStatus(conf *cfg.Config) error {
res, err := getClusterStatus(conf)
if err != nil {
return err
}
slog.Debug("ES result", "cluster health", res.health)
isleader := len(res.ccr.AutoFollowStats.AutoFollowedClusters) == 0
ccrfollowing := "" ccrfollowing := ""
if len(ccrstats.AutoFollowStats.AutoFollowedClusters) > 0 { if len(res.ccr.AutoFollowStats.AutoFollowedClusters) > 0 {
// is following another cluster // is following another cluster
ccrfollowing = fmt.Sprintf("%s (%d/%d)", ccrfollowing = fmt.Sprintf("%s (%d/%d)",
ccrstats.AutoFollowStats.AutoFollowedClusters[0].ClusterName, res.ccr.AutoFollowStats.AutoFollowedClusters[0].ClusterName,
ccrstats.AutoFollowStats.NumberOfSuccessfulFollowIndices, res.ccr.AutoFollowStats.NumberOfSuccessfulFollowIndices,
ccrstats.AutoFollowStats.NumberOfFailedFollowIndices, res.ccr.AutoFollowStats.NumberOfFailedFollowIndices,
) )
} }
// look for red indices, if any // look for red indices, if any
redindices := 0 redindices := 0
for _, index := range *indexstats { for _, index := range *res.indices {
if *index.Health == "red" { if *index.Health == "red" {
redindices++ redindices++
} }
@@ -169,26 +158,26 @@ func ClusterStatus(conf *cfg.Config) error {
// look for long running tasks // look for long running tasks
longtasks := 0 longtasks := 0
for _, task := range *taskstatus { for _, task := range *res.tasks {
if strings.Contains(*task.RunningTime, "d") { if strings.Contains(*task.RunningTime, "d") {
longtasks++ longtasks++
} }
} }
table := printer.NewTable(conf, 2, 7) table := printer.NewTable(conf, 2, 7)
table.Addheaders(cluster, "status") table.Addheaders(conf.DefaultCluster.Name, "status")
table.Entries = [][]string{ table.Entries = [][]string{
{"Cluster Name", clusterhealth.ClusterName}, {"Cluster Name", res.health.ClusterName},
{"ES Status", printer.Colorize(conf, clusterhealth.Status.Name, clusterhealth.Status.Name)}, {"ES Status", printer.Colorize(conf, res.health.Status.Name, res.health.Status.Name)},
{"ES Version", info.Version.Int}, {"ES Version", res.info.Version.Int},
{"Is Leader", fmt.Sprintf("%t", isleader)}, {"Is Leader", fmt.Sprintf("%t", isleader)},
{"Active Shards", fmt.Sprintf("%d", clusterhealth.ActiveShards)}, {"Active Shards", fmt.Sprintf("%d", res.health.ActiveShards)},
{"Active Primary Shards", fmt.Sprintf("%d", clusterhealth.ActivePrimaryShards)}, {"Active Primary Shards", fmt.Sprintf("%d", res.health.ActivePrimaryShards)},
{"Unassigned Shards", fmt.Sprintf("%d", clusterhealth.UnassignedShards)}, {"Unassigned Shards", fmt.Sprintf("%d", res.health.UnassignedShards)},
{"Unassigned Primary Shards", fmt.Sprintf("%d", clusterhealth.UnassignedPrimaryShards)}, {"Unassigned Primary Shards", fmt.Sprintf("%d", res.health.UnassignedPrimaryShards)},
{"Pending Tasks", fmt.Sprintf("%d", clusterhealth.NumberOfPendingTasks)}, {"Pending Tasks", fmt.Sprintf("%d", res.health.NumberOfPendingTasks)},
{"Nodes", fmt.Sprintf("%d", clusterhealth.NumberOfNodes)}, {"Nodes", fmt.Sprintf("%d", res.health.NumberOfNodes)},
{"Red Indices", fmt.Sprintf("%d", redindices)}, {"Red Indices", fmt.Sprintf("%d", redindices)},
{"Long Running Tasks", fmt.Sprintf("%d", longtasks)}, {"Long Running Tasks", fmt.Sprintf("%d", longtasks)},
} }
@@ -196,18 +185,35 @@ func ClusterStatus(conf *cfg.Config) error {
if !isleader { if !isleader {
table.Entries = append(table.Entries, [][]string{ table.Entries = append(table.Entries, [][]string{
{"AutoFollow (success/failed indices)", ccrfollowing}, {"AutoFollow (success/failed indices)", ccrfollowing},
{"Followed Indices", fmt.Sprintf("%d", len(ccrstats.FollowStats.Indices))}, {"Followed Indices", fmt.Sprintf("%d", len(res.ccr.FollowStats.Indices))},
}...) }...)
} }
if conf.Verbose { if conf.Verbose {
table = gatherClusterStats(conf, clusterstats, table) table = gatherClusterStats(conf, res.stats, table)
}
if res.health.Status.Name != "green" {
for name, indicator := range res.healthreport.Indicators {
if indicator.Status != "green" {
table.Entries = append(table.Entries, []string{
printer.Colorize(conf, indicator.Status, "Bad health "+name), indicator.Symptom,
})
for _, diag := range indicator.Diagnosis {
table.Entries = append(table.Entries, []string{" -> cause", diag.Cause})
for resource, items := range diag.AffectedResources {
table.Entries = append(table.Entries, []string{" -> affected " + resource, strings.Join(items, ",")})
}
}
}
}
} }
if err := table.Print(); err != nil { if err := table.Print(); err != nil {
return err return err
} }
}
return nil return nil
} }

125
pkg/es/cluster_reroute.go Normal file
View File

@@ -0,0 +1,125 @@
/*
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"
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/elastic/go-elasticsearch/v9/typedapi/esdsl"
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
)
func ClusterRerouteMove(conf *cfg.Config, index string) error {
move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand()
moveCommand := &types.CommandMoveAction{
Shard: conf.Shards,
FromNode: conf.FromNode,
ToNode: conf.ToNode,
Index: index,
}
commands.CommandCaster().Move = moveCommand
move.Commands(commands)
_, err := move.Do(context.Background())
if err != nil {
return fmt.Errorf("failed to reroute move: %w", err)
}
return nil
}
func ClusterRerouteAllocateReplica(conf *cfg.Config, index string) error {
move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand()
allocCommand := &types.CommandAllocateReplicaAction{
Shard: conf.Shards,
Node: conf.ToNode,
Index: index,
}
commands.CommandCaster().AllocateReplica = allocCommand
move.Commands(commands)
_, err := move.Do(context.Background())
if err != nil {
return fmt.Errorf("failed to allocate a replica shard: %w", err)
}
return nil
}
func ClusterRerouteCancel(conf *cfg.Config, index string) error {
move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand()
cancelCommand := &types.CommandCancelAction{
Shard: conf.Shards,
Node: conf.ToNode,
Index: index,
AllowPrimary: &conf.AllowPrimary,
}
commands.CommandCaster().Cancel = cancelCommand
move.Commands(commands)
_, err := move.Do(context.Background())
if err != nil {
return fmt.Errorf("failed to cancel a reroute process: %w", err)
}
return nil
}
func ClusterRerouteAllocatePrimary(conf *cfg.Config, index string, stale bool) error {
move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand()
allocCommand := &types.CommandAllocatePrimaryAction{
Shard: conf.Shards,
Node: conf.ToNode,
Index: index,
AcceptDataLoss: conf.AcceptDataLoss,
}
if stale {
commands.CommandCaster().AllocateStalePrimary = allocCommand
} else {
commands.CommandCaster().AllocateEmptyPrimary = allocCommand
}
move.Commands(commands)
_, err := move.Do(context.Background())
if err != nil {
which := "empty"
if stale {
which = "stale"
}
return fmt.Errorf("failed to allocate %s primary shard: %w", which, err)
}
return nil
}

48
pkg/es/debug.go Normal file
View File

@@ -0,0 +1,48 @@
/*
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 (
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/alecthomas/repr"
)
func Debug(conf *cfg.Config) error {
report, err := getHealthReport(conf)
if err != nil {
return err
}
repr.Println(report)
/*
res, err := conf.DefaultCluster.ES().Search().
Index(conf.Index).
Size(0).
Aggregations(map[string]types.Aggregations{
"min_ts": *esdsl.NewMinAggregation().Field("@timestamp").AggregationsCaster(),
"max_ts": *esdsl.NewMaxAggregation().Field("@timestamp").AggregationsCaster(),
}).
Do(context.Background())
if err != nil {
return err
}
repr.Println(res)
*/
return nil
}

View File

@@ -121,6 +121,10 @@ func IlmShow(conf *cfg.Config, policy string) error {
return errors.New("no ilm policy retrieved") return errors.New("no ilm policy retrieved")
} }
if conf.Ilm.Tree {
return IlmShowTree(conf, ilm.Policy)
}
table := printer.NewTable(conf, 2, 5) table := printer.NewTable(conf, 2, 5)
table.Addheaders("ilm policy setting", "value") table.Addheaders("ilm policy setting", "value")
@@ -135,6 +139,101 @@ func IlmShow(conf *cfg.Config, policy string) error {
return table.Print() return table.Print()
} }
// Visualize phases or a ILM policy, taking into account that the
// minAge for the current phases is set in the next phase
func IlmShowTree(conf *cfg.Config, ilm types.IlmPolicy) error {
indent := ""
table := printer.NewTable(conf, 4, 0)
table.Addheaders("phase", "min age", "min size", "snapshot repo")
for _, phase := range IlmPhaseOrder {
if phase == "hot" {
fmt.Printf("%s%s phase:\n%s rollover after %s\n",
indent, phase,
indent, formatDuration(parseDuration(ilm.Phases.Hot.Actions.Rollover.MaxAge.(string))),
)
if ilm.Phases.Hot.Actions.Rollover.MaxPrimaryShardSize != "" {
fmt.Printf("%s rollover when storage > %s\n",
indent, ilm.Phases.Hot.Actions.Rollover.MaxPrimaryShardSize)
}
if ilm.Phases.Hot.Actions.Rollover.MaxDocs != nil {
fmt.Printf("%s rollover docs > %d\n",
indent, *ilm.Phases.Hot.Actions.Rollover.MaxDocs)
}
indent += " "
continue
}
current := getIlmCurrentPhase(ilm, phase)
if current == nil {
continue
}
next := findNextPhase(ilm, phase)
fmt.Printf("%s%s phase:\n", indent, phase)
if next != nil {
fmt.Printf("%s rollover after %s\n",
indent, formatDuration(next.minage),
)
if current.Actions.SearchableSnapshot != nil {
fmt.Printf("%s roll to snapshot repo: %s\n",
indent, current.Actions.SearchableSnapshot.SnapshotRepository)
}
if current.Actions.Forcemerge != nil {
fmt.Printf("%s force merge segments: %d\n",
indent, current.Actions.Forcemerge.MaxNumSegments)
}
if current.Actions.Shrink != nil {
fmt.Printf("%s shrink shards: %d\n",
indent, *current.Actions.Shrink.NumberOfShards)
}
if current.Actions.SetPriority != nil {
fmt.Printf("%s priority: %d\n",
indent, *current.Actions.SetPriority.Priority)
}
if current.Actions.Delete != nil && current.Actions.Delete.DeleteSearchableSnapshot != nil {
fmt.Printf("%s delete searchable snapshots: %t\n",
indent, *current.Actions.Delete.DeleteSearchableSnapshot)
}
}
if phase == "delete" {
fmt.Printf("%s delete immediately\n", indent)
}
indent += " "
}
return nil
}
func getIlmCurrentPhase(policy types.IlmPolicy, currentPhase string) *types.Phase {
switch currentPhase {
case "hot":
return policy.Phases.Hot
case "warm":
return policy.Phases.Warm
case "cold":
return policy.Phases.Cold
case "frozen":
return policy.Phases.Frozen
default:
return policy.Phases.Delete
}
}
func ilmPhaseString(phase *types.Phase, short bool) string { func ilmPhaseString(phase *types.Phase, short bool) string {
if phase == nil { if phase == nil {
return "" return ""
@@ -220,6 +319,24 @@ func IlmExplain(conf *cfg.Config, index string) error {
} }
func IlmCreate(conf *cfg.Config, policyname string) error { func IlmCreate(conf *cfg.Config, policyname string) error {
var policy *types.IlmPolicy = nil
res, err := conf.DefaultCluster.ES().Ilm.GetLifecycle().
Policy(policyname).
Do(context.Background())
if err == nil {
if conf.Debug {
repr.Println(res)
}
ilm, exists := res[policyname]
if !exists {
return errors.New("no ilm policy retrieved")
}
policy = &ilm.Policy
}
ilm := conf.DefaultCluster.ES().Ilm.PutLifecycle(policyname) ilm := conf.DefaultCluster.ES().Ilm.PutLifecycle(policyname)
cfg := conf.Ilm cfg := conf.Ilm
@@ -231,6 +348,16 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
rollover := &types.RolloverAction{} rollover := &types.RolloverAction{}
haveroll := false haveroll := false
if policy != nil {
// update
actions = policy.Phases.Hot.Actions
if policy.Phases.Hot.Actions.Rollover != nil {
rollover = policy.Phases.Hot.Actions.Rollover
haveroll = true
}
}
if cfg.HotMinAge != "" { if cfg.HotMinAge != "" {
hot.MinAge = cfg.HotMinAge hot.MinAge = cfg.HotMinAge
} }
@@ -255,12 +382,22 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
hot.Actions = actions.IlmActionsCaster() hot.Actions = actions.IlmActionsCaster()
phases.PhasesCaster().Hot = &hot phases.PhasesCaster().Hot = &hot
} else {
if policy != nil {
// update
phases.PhasesCaster().Hot = policy.Phases.Hot
}
} }
if cfg.HaveWarm() { if cfg.HaveWarm() {
warm := types.Phase{} warm := types.Phase{}
var actions types.IlmActionsVariant = esdsl.NewIlmActions() var actions types.IlmActionsVariant = esdsl.NewIlmActions()
if policy != nil {
// update
actions = policy.Phases.Warm.Actions
}
if cfg.WarmForceMerge != 0 { if cfg.WarmForceMerge != 0 {
actions.IlmActionsCaster().Forcemerge = &types.ForceMergeAction{MaxNumSegments: cfg.WarmForceMerge} actions.IlmActionsCaster().Forcemerge = &types.ForceMergeAction{MaxNumSegments: cfg.WarmForceMerge}
} }
@@ -279,12 +416,22 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
warm.Actions = actions.IlmActionsCaster() warm.Actions = actions.IlmActionsCaster()
phases.PhasesCaster().Warm = &warm phases.PhasesCaster().Warm = &warm
} else {
if policy != nil && policy.Phases.Warm != nil {
// update
phases.PhasesCaster().Warm = policy.Phases.Warm
}
} }
if cfg.HaveCold() { if cfg.HaveCold() {
cold := types.Phase{} cold := types.Phase{}
var actions types.IlmActionsVariant = esdsl.NewIlmActions() var actions types.IlmActionsVariant = esdsl.NewIlmActions()
if policy != nil {
// update
actions = policy.Phases.Cold.Actions
}
if cfg.ColdForceMerge != 0 { if cfg.ColdForceMerge != 0 {
actions.IlmActionsCaster().Forcemerge = &types.ForceMergeAction{MaxNumSegments: cfg.ColdForceMerge} actions.IlmActionsCaster().Forcemerge = &types.ForceMergeAction{MaxNumSegments: cfg.ColdForceMerge}
} }
@@ -304,12 +451,22 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
cold.Actions = actions.IlmActionsCaster() cold.Actions = actions.IlmActionsCaster()
phases.PhasesCaster().Cold = &cold phases.PhasesCaster().Cold = &cold
} else {
if policy != nil && policy.Phases.Cold != nil {
// update
phases.PhasesCaster().Cold = policy.Phases.Cold
}
} }
if cfg.HaveFrozen() { if cfg.HaveFrozen() {
froze := types.Phase{} froze := types.Phase{}
var actions types.IlmActionsVariant = esdsl.NewIlmActions() var actions types.IlmActionsVariant = esdsl.NewIlmActions()
if policy != nil {
// update
actions = policy.Phases.Frozen.Actions
}
if cfg.FrozenMinAge != "" { if cfg.FrozenMinAge != "" {
froze.MinAge = cfg.FrozenMinAge froze.MinAge = cfg.FrozenMinAge
} }
@@ -321,6 +478,11 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
froze.Actions = actions.IlmActionsCaster() froze.Actions = actions.IlmActionsCaster()
phases.PhasesCaster().Frozen = &froze phases.PhasesCaster().Frozen = &froze
} else {
if policy != nil && policy.Phases.Frozen != nil {
// update
phases.PhasesCaster().Frozen = policy.Phases.Frozen
}
} }
if cfg.HaveDelete() { if cfg.HaveDelete() {
@@ -328,6 +490,11 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
delete := types.DeleteAction{} delete := types.DeleteAction{}
var actions types.IlmActionsVariant = esdsl.NewIlmActions() var actions types.IlmActionsVariant = esdsl.NewIlmActions()
if policy != nil {
// update
actions = policy.Phases.Delete.Actions
}
if cfg.DeleteMinAge != "" { if cfg.DeleteMinAge != "" {
del.MinAge = cfg.DeleteMinAge del.MinAge = cfg.DeleteMinAge
} }
@@ -339,16 +506,21 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
actions.IlmActionsCaster().Delete = &delete actions.IlmActionsCaster().Delete = &delete
del.Actions = actions.IlmActionsCaster() del.Actions = actions.IlmActionsCaster()
phases.PhasesCaster().Delete = &del phases.PhasesCaster().Delete = &del
} else {
if policy != nil && policy.Phases.Delete != nil {
// update
phases.PhasesCaster().Delete = policy.Phases.Delete
}
} }
put := &putlifecycle.Request{} put := &putlifecycle.Request{}
policy := &types.IlmPolicy{} newpolicy := &types.IlmPolicy{}
policy.IlmPolicyCaster().Phases = *phases.PhasesCaster() newpolicy.IlmPolicyCaster().Phases = *phases.PhasesCaster()
put.Policy = policy put.Policy = newpolicy
ilm.Request(put) ilm.Request(put)
_, err := ilm.Do(context.Background()) _, err = ilm.Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed create ilm policy: %s", esErrorString(err)) return fmt.Errorf("failed create ilm policy: %s", esErrorString(err))
} }

View File

@@ -190,9 +190,9 @@ func getIlmPhaseData(conf *cfg.Config) ([]PhaseData, error) {
wg := &sync.WaitGroup{} wg := &sync.WaitGroup{}
wg.Add(3) wg.Add(3)
go getApiData(conf.DefaultCluster.ES(), wg, responses, "indicesbytes") go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "indicesbytes")
go getApiData(conf.DefaultCluster.ES(), wg, responses, "explain") go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "explain")
go getApiData(conf.DefaultCluster.ES(), wg, responses, "policies") go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "policies")
wg.Wait() wg.Wait()

View File

@@ -20,9 +20,12 @@ import (
"context" "context"
"fmt" "fmt"
"log/slog" "log/slog"
"strings"
"time"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer" "codeberg.org/scip/esctl/pkg/printer"
"github.com/dustin/go-humanize"
) )
func NodeList(conf *cfg.Config) error { func NodeList(conf *cfg.Config) error {
@@ -56,3 +59,145 @@ func NodeList(conf *cfg.Config) error {
return nil return nil
} }
func NodeNames(conf *cfg.Config) ([]string, error) {
nodes, err := conf.DefaultCluster.ES().Cat.Nodes().
Do(context.Background())
if err != nil {
return nil, fmt.Errorf("failed to get nodes: %s", esErrorString(err))
}
slog.Debug("ES result", "nodes", nodes)
nodelist := make([]string, len(nodes))
for idx, node := range nodes {
nodelist[idx] = *node.Name
}
return nodelist, err
}
// FIXME: adding the settings metric leads to json unmarshall error:
// https://github.com/elastic/go-elasticsearch/issues/1524
func NodeShow(conf *cfg.Config, nodename string) error {
res, err := conf.DefaultCluster.ES().Nodes.Info().
NodeId(nodename).
Metric("os, jvm, thread_pool, remote_cluster_server").
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get node info: %s", esErrorString(err))
}
slog.Debug("ES result", "node", res)
stats, err := conf.DefaultCluster.ES().Nodes.Stats().
NodeId(nodename).
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get node stats: %s", esErrorString(err))
}
slog.Debug("ES result", "stat", stats)
for id, info := range res.Nodes {
stat := stats.Nodes[id]
table := printer.NewTable(conf, 2, 0)
table.Addheaders(nodename+" property", "value")
roles := make([]string, len(info.Roles))
for idx, role := range info.Roles {
roles[idx] = role.Name
}
k8snode := info.Attributes["k8s_node_name"]
table.Entries = [][]string{
{"Id", id},
{"Name", nodename},
{"Kubernetes node", k8snode},
{"Ip address", info.Ip},
{"Node rank", *stat.AdaptiveSelection[id].Rank},
{"JVM", info.Jvm.VmName + " " + info.Jvm.Version},
{"JVM Started", time.UnixMilli(info.Jvm.StartTimeInMillis).String()},
{"OS", info.Os.PrettyName + " " + info.Os.Version},
{"Node roles", strings.Join(roles, ",")},
{"Node version", info.Version},
{"HTTP clients", fmt.Sprintf("%d", *stat.Http.CurrentOpen)},
{"CPUs", fmt.Sprintf("%d", *info.Os.AllocatedProcessors)},
{"Load 15m/5m/1m", fmt.Sprintf("%.2f/%.2f/%.2f",
stat.Os.Cpu.LoadAverage["15m"],
stat.Os.Cpu.LoadAverage["5m"],
stat.Os.Cpu.LoadAverage["1m"],
)},
{"Open FD's", fmt.Sprintf("%d", *stat.Process.OpenFileDescriptors)},
{"Response time avg", fmt.Sprintf("%dns", *stat.AdaptiveSelection[id].AvgResponseTimeNs)},
{"Memory usage (used/avail)",
humanize.Bytes(uint64(*stat.Os.Mem.UsedInBytes)) + " / " + humanize.Bytes(uint64(*stat.Os.Mem.TotalInBytes))},
}
if len(stat.Fs.Data) > 0 {
fs := stat.Fs.Data[0]
table.Entries = append(table.Entries, [][]string{
{"Storage usage (used/avail)",
humanize.Bytes(uint64(*fs.AvailableInBytes)) + " / " + humanize.Bytes(uint64(*fs.TotalInBytes))},
{"Storage mount", *fs.Mount},
}...)
}
if err := table.Print(); err != nil {
return err
}
}
return nil
}
func NodeClients(conf *cfg.Config, nodename string) error {
stats, err := conf.DefaultCluster.ES().Nodes.Stats().
NodeId(nodename).
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get node stats: %s", esErrorString(err))
}
slog.Debug("ES result", "stat", stats)
table := printer.NewTable(conf, 5, 0)
table.Addheaders("agent", "id", "when", "from host", "url")
for _, stat := range stats.Nodes {
for _, client := range stat.Http.Clients {
if client.ClosedTimeMillis == nil {
agent := ""
if client.Agent != nil {
agent = *client.Agent
}
uri := ""
if client.LastUri != nil {
uri = *client.LastUri
if !conf.All {
parts := strings.Split(uri, "?")
uri = parts[0]
}
}
table.AddRow(
agent,
fmt.Sprintf("%d", *client.Id),
time.UnixMilli(*client.LastRequestTimeMillis).String(),
*client.RemoteAddress,
uri,
)
}
}
break
}
table.Sort()
return table.Print()
}

View File

@@ -18,10 +18,10 @@ package es
import ( import (
"context" "context"
"errors"
"fmt" "fmt"
"sync" "sync"
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/elastic/go-elasticsearch/v9" "github.com/elastic/go-elasticsearch/v9"
"github.com/elastic/go-elasticsearch/v9/typedapi/cat/indices" "github.com/elastic/go-elasticsearch/v9/typedapi/cat/indices"
"github.com/elastic/go-elasticsearch/v9/typedapi/cat/tasks" "github.com/elastic/go-elasticsearch/v9/typedapi/cat/tasks"
@@ -43,12 +43,14 @@ const (
ResponseTasks ResponseTasks
ResponseExplain ResponseExplain
ResponseLifecycle ResponseLifecycle
ResponseHealthReport
) )
type apiResponse struct { type apiResponse struct {
error error error error
info *info.Response info *info.Response
health *health.Response health *health.Response
healthreport *HealthReport
ccr *stats.Response ccr *stats.Response
stats *clusterstats.Response stats *clusterstats.Response
indices *indices.Response indices *indices.Response
@@ -59,12 +61,16 @@ type apiResponse struct {
which int which int
} }
func getApiData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, func getApiData(
reschan chan apiResponse, which string) { conf *cfg.Config,
es *elasticsearch.TypedClient,
wg *sync.WaitGroup,
reschan chan apiResponse,
which string) {
defer wg.Done() defer wg.Done()
ar := apiResponse{} ar := apiResponse{}
arerr := errors.New("") var arerr error
switch which { switch which {
case "health": case "health":
@@ -75,6 +81,13 @@ func getApiData(es *elasticsearch.TypedClient, wg *sync.WaitGroup,
ar.which = ResponseHealth ar.which = ResponseHealth
arerr = err arerr = err
case "healthreport":
report, err := getHealthReport(conf)
ar.healthreport = report
ar.which = ResponseHealthReport
arerr = err
case "info": case "info":
res, err := es.Info(). res, err := es.Info().
Do(context.Background()) Do(context.Background())
@@ -83,7 +96,7 @@ func getApiData(es *elasticsearch.TypedClient, wg *sync.WaitGroup,
ar.which = ResponseInfo ar.which = ResponseInfo
arerr = err arerr = err
case "ccrstats": case "ccr":
res, err := es.Ccr.Stats(). res, err := es.Ccr.Stats().
Do(context.Background()) Do(context.Background())

View File

@@ -26,7 +26,6 @@ import (
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer" "codeberg.org/scip/esctl/pkg/printer"
"github.com/alecthomas/repr"
"github.com/elastic/go-elasticsearch/v9/typedapi/core/search" "github.com/elastic/go-elasticsearch/v9/typedapi/core/search"
"github.com/elastic/go-elasticsearch/v9/typedapi/esdsl" "github.com/elastic/go-elasticsearch/v9/typedapi/esdsl"
"github.com/elastic/go-elasticsearch/v9/typedapi/indices/validatequery" "github.com/elastic/go-elasticsearch/v9/typedapi/indices/validatequery"
@@ -149,25 +148,6 @@ func validateSearch(conf *cfg.Config, queries []string) error {
return nil return nil
} }
func Debug(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES().Search().
Index(conf.Index).
Size(0).
Aggregations(map[string]types.Aggregations{
"min_ts": *esdsl.NewMinAggregation().Field("@timestamp").AggregationsCaster(),
"max_ts": *esdsl.NewMaxAggregation().Field("@timestamp").AggregationsCaster(),
}).
Do(context.Background())
if err != nil {
return err
}
repr.Println(res)
return nil
}
func searchOnce(conf *cfg.Config, search *search.Search) error { func searchOnce(conf *cfg.Config, search *search.Search) error {
res, err := search. res, err := search.
From(conf.From). From(conf.From).

View File

@@ -28,6 +28,7 @@ import (
"github.com/olekukonko/tablewriter" "github.com/olekukonko/tablewriter"
"github.com/olekukonko/tablewriter/renderer" "github.com/olekukonko/tablewriter/renderer"
"github.com/olekukonko/tablewriter/tw" "github.com/olekukonko/tablewriter/tw"
"github.com/seeruk/go-wordwrap"
"gopkg.in/yaml.v3" "gopkg.in/yaml.v3"
) )
@@ -38,10 +39,11 @@ type Table struct {
lenHeaders []int lenHeaders []int
alignInts bool alignInts bool
maxwidth int
} }
func NewTable(conf *cfg.Config, columns, rows int) *Table { func NewTable(conf *cfg.Config, columns, rows int) *Table {
table := Table{Mode: conf.Output} table := Table{Mode: conf.Output, maxwidth: cfg.GetTermWidth()}
table.Headers = make([]string, columns) table.Headers = make([]string, columns)
table.Entries = make([][]string, rows) table.Entries = make([][]string, rows)
@@ -118,17 +120,31 @@ func (data *Table) PrintTSV() error {
} }
for _, entries := range data.Entries { for _, entries := range data.Entries {
currentWidth := 0
for idx, entry := range entries { for idx, entry := range entries {
length := visibleLen(entry) length := visibleLen(entry)
if data.lenHeaders[idx] < length { if data.lenHeaders[idx] < length {
if length > currentWidth+data.maxwidth {
data.lenHeaders[idx] = data.maxwidth - currentWidth
} else {
data.lenHeaders[idx] = length data.lenHeaders[idx] = length
} }
} }
currentWidth += data.lenHeaders[idx]
}
} }
// output // output headers
for idx, header := range data.Headers { for idx, header := range data.Headers {
if idx+1 != len(data.Headers) {
fmt.Print(header, strings.Repeat(" ", data.lenHeaders[idx]-visibleLen(header))) fmt.Print(header, strings.Repeat(" ", data.lenHeaders[idx]-visibleLen(header)))
} else {
// no padding for last header
fmt.Print(header)
}
if idx < len(data.Headers)-1 { if idx < len(data.Headers)-1 {
fmt.Print(" ") fmt.Print(" ")
} }
@@ -136,14 +152,37 @@ func (data *Table) PrintTSV() error {
fmt.Println() fmt.Println()
for _, entries := range data.Entries { for _, entries := range data.Entries {
currentWidth := 0
for idx, entry := range entries { for idx, entry := range entries {
length := visibleLen(entry) length := visibleLen(entry)
if length+currentWidth > data.maxwidth {
// text is too wide to be put into one line, wrap it
wrapper := wordwrap.Wrapper(data.maxwidth-currentWidth, false)
wrapped := wrapper(entry)
// and indent it
for idx, line := range strings.Split(wrapped, "\n") {
if idx == 0 {
entry = line
} else {
entry += "\n " + strings.Repeat(" ", currentWidth) + line
}
}
}
currentWidth += data.lenHeaders[idx]
if isInt(entry) && data.alignInts { if isInt(entry) && data.alignInts {
// align right // align right
fmt.Print(strings.Repeat(" ", data.lenHeaders[idx]-length), entry) fmt.Print(strings.Repeat(" ", data.lenHeaders[idx]-length), entry)
} else { } else if length < data.lenHeaders[idx] && idx+1 != len(entries) {
// pad right, if required
fmt.Print(entry, strings.Repeat(" ", data.lenHeaders[idx]-length)) fmt.Print(entry, strings.Repeat(" ", data.lenHeaders[idx]-length))
} else {
// no padding for last entry
fmt.Print(entry)
} }
if idx < len(data.Headers)-1 { if idx < len(data.Headers)-1 {