mirror of
https://codeberg.org/scip/esctl.git
synced 2026-08-24 17:04:18 +02:00
Compare commits
9 Commits
0.0.22
...
feature/au
| Author | SHA1 | Date | |
|---|---|---|---|
| 489074da01 | |||
|
|
b573d63c54 | ||
|
|
e3b96b3e3b | ||
|
|
6436bcfb22 | ||
|
|
f5c8a23589 | ||
|
|
f9eded5f92 | ||
|
|
b7a87051c0 | ||
|
|
81987b75f2 | ||
| 34bf4fbed9 |
223
README.md
223
README.md
@@ -113,6 +113,7 @@ Configure `esctl` with environment variables:
|
||||
- `ES_URI`: elasticsearch uri
|
||||
- `ES_USER`: username
|
||||
- `ES_PASS`: password
|
||||
- `ES_TOKEN`: API token, instead of user+password
|
||||
|
||||
Or create a config file such as this:
|
||||
|
||||
@@ -121,11 +122,10 @@ clusters:
|
||||
foobar:
|
||||
uri: https://es.foo.bar:9200/
|
||||
user: elastic
|
||||
pass: 123456
|
||||
pass: ******
|
||||
other:
|
||||
uri: https://myes.foo:9200/
|
||||
user: elastic
|
||||
pass: asdasdasd
|
||||
token: ******
|
||||
```
|
||||
|
||||
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
|
||||
changing heavily every now and then. New commands are being added
|
||||
constantly as well.
|
||||
@@ -440,95 +474,100 @@ constantly as well.
|
||||
### Command tree:
|
||||
|
||||
```console
|
||||
api - api access and documentation
|
||||
list - list index of API calls
|
||||
show - show an API doc
|
||||
repl - interactive API repl
|
||||
ccr - manage cross cluster replication
|
||||
status - cross cluster replication status (yaml config with 2 clusters required)
|
||||
pause - pause shard allocation
|
||||
resume - resume shard allocation
|
||||
follower - manage ccr follower indices
|
||||
show - show ccr follower index details
|
||||
add - add ccr follower index
|
||||
delete - delete ccr follower index
|
||||
unfollow - unfollow ccr follower index
|
||||
pause - pause ccr index to follow
|
||||
resume - resume ccr index to follow
|
||||
renew - renew ccr follower index
|
||||
info - show ccr remote info
|
||||
cluster - manage cluster[s]
|
||||
status - show cluster status
|
||||
switch - set current elasticsearch cluster
|
||||
list - list configured clusters
|
||||
settings - cluster settings management
|
||||
list - show cluster settings
|
||||
set - set|update cluster settings
|
||||
datastream - manage data streams
|
||||
list - list indicies
|
||||
show - show details about an data stream
|
||||
create - create a new data stream
|
||||
delete - delete a data stream
|
||||
rollover - roll over a data stream
|
||||
doc - manage documents
|
||||
add - add JSON document index
|
||||
show - show a JSON document
|
||||
delete - delete JSON document[s] from index[es]
|
||||
ilm - manage index lifecycle
|
||||
retry - retry applying an ILM profile to an index
|
||||
status - get the current index lifecycle management status
|
||||
list - list index lifecycle policies
|
||||
show - show details about an index lifecycle policy
|
||||
create - create a index lifecycle policy
|
||||
forecast - calculate index phase movements
|
||||
list - list index rollover config
|
||||
show - show rollover forecast over all indices
|
||||
index - manage indicies
|
||||
list - list indicies
|
||||
show - show details about an index
|
||||
create - create a new index
|
||||
modify - modify anindex
|
||||
delete - delete an index
|
||||
close - close an index
|
||||
allocation - explain index allocation
|
||||
fields - show info about field capabilities
|
||||
ilm - show ilm status
|
||||
alias - manage index aliases
|
||||
create - create an index alias
|
||||
list - list index aliases
|
||||
delete - delete an index alias
|
||||
rollover - roll over an index alias
|
||||
template - manage index templates
|
||||
list - list index templates
|
||||
show - show details about an index template
|
||||
create - create a new index template
|
||||
modify - modify a new index template
|
||||
delete - delete an index template
|
||||
node - manage nodes
|
||||
list - list nodes
|
||||
show - show details about a node
|
||||
role - manage roles
|
||||
list - list roles
|
||||
show - show details about a role
|
||||
diff - show differences between roles and CSV baseline
|
||||
search - search within an index
|
||||
shard - manage shards
|
||||
list - list shards
|
||||
show - show details about a shard
|
||||
snapshot - manage snapshots
|
||||
list - list snapshots
|
||||
show - show details about a snapshot
|
||||
task - manage tasks
|
||||
list - list tasks
|
||||
cancel - cancel running task
|
||||
version - show esctl version information
|
||||
debug - developer only
|
||||
help-jsonpath - show jsonpath help
|
||||
completion - Output shell completion script for bash, zsh, fish, or Powershell
|
||||
pwsh - Output pwsh completion script
|
||||
bash - Output bash completion script
|
||||
zsh - Output zsh completion script
|
||||
fish - Output fish completion script
|
||||
api - api access and documentation
|
||||
list - list index of API calls
|
||||
show - show an API doc
|
||||
repl - interactive API repl
|
||||
ccr - manage cross cluster replication
|
||||
status - cross cluster replication status (yaml config with 2 clusters required)
|
||||
pause - pause shard allocation
|
||||
resume - resume shard allocation
|
||||
follower - manage ccr follower indices
|
||||
show - show ccr follower index details
|
||||
add - add ccr follower index
|
||||
delete - delete ccr follower index
|
||||
unfollow - unfollow ccr follower index
|
||||
pause - pause ccr index to follow
|
||||
resume - resume ccr index to follow
|
||||
renew - renew ccr follower index
|
||||
info - show ccr remote info
|
||||
cluster - manage cluster[s]
|
||||
status - show cluster status
|
||||
switch - set current elasticsearch cluster
|
||||
list - list configured clusters
|
||||
settings - cluster settings management
|
||||
list - show 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
|
||||
list - list indicies
|
||||
show - show details about an data stream
|
||||
create - create a new data stream
|
||||
delete - delete a data stream
|
||||
rollover - roll over a data stream
|
||||
doc - manage documents
|
||||
add - add JSON document index
|
||||
show - show a JSON document
|
||||
delete - delete JSON document[s] from index[es]
|
||||
ilm - manage index lifecycle
|
||||
retry - retry applying an ILM profile to an index
|
||||
status - get the current index lifecycle management status
|
||||
list - list index lifecycle policies
|
||||
show - show details about an index lifecycle policy
|
||||
create - create a new lifecycle policy
|
||||
update - update an lifecycle policy
|
||||
forecast - calculate index phase movements
|
||||
list - list index rollover config
|
||||
show - show rollover forecast over all indices
|
||||
explain - explain ilm condition of an index
|
||||
index - manage indicies
|
||||
list - list indicies
|
||||
show - show details about an index
|
||||
create - create a new index
|
||||
update - update an index
|
||||
delete - delete an index
|
||||
close - close an index
|
||||
allocation - explain index allocation
|
||||
fields - show info about field capabilities
|
||||
ilm - show ilm status
|
||||
alias - manage index aliases
|
||||
create - create an index alias
|
||||
list - list index aliases
|
||||
delete - delete an index alias
|
||||
rollover - roll over an index alias
|
||||
template - manage index templates
|
||||
list - list index templates
|
||||
show - show details about an index template
|
||||
create - create a new index template
|
||||
update - update a new index template
|
||||
delete - delete an index template
|
||||
node - manage nodes
|
||||
list - list nodes
|
||||
show - show details about a node
|
||||
clients - show node http clients
|
||||
role - manage roles
|
||||
list - list roles
|
||||
show - show details about a role
|
||||
diff - show differences between roles and CSV baseline
|
||||
search - search within an index
|
||||
shard - manage shards
|
||||
list - list shards
|
||||
show - show details about a shard
|
||||
snapshot - manage snapshots
|
||||
list - list snapshots
|
||||
show - show details about a snapshot
|
||||
task - manage tasks
|
||||
list - list tasks
|
||||
cancel - cancel running task
|
||||
version - show esctl version information
|
||||
debug - developer only
|
||||
help-jsonpath - show jsonpath help
|
||||
help-command-overview - show overview of all available commands
|
||||
```
|
||||
|
||||
# Development
|
||||
|
||||
@@ -37,6 +37,7 @@ func Cluster(conf *cfg.Config) *cli.Command {
|
||||
ClusterSwitch(conf),
|
||||
ClusterList(conf),
|
||||
ClusterSettings(conf),
|
||||
ClusterReroute(conf),
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -60,12 +61,6 @@ func ClusterStatus(conf *cfg.Config) *cli.Command {
|
||||
Aliases: []string{"s"},
|
||||
|
||||
Flags: []cli.Flag{
|
||||
&cli.BoolFlag{
|
||||
Name: "all",
|
||||
Usage: "show status of all clusters",
|
||||
Destination: &conf.All,
|
||||
Aliases: []string{"a"},
|
||||
},
|
||||
&cli.BoolFlag{
|
||||
Name: "verbose",
|
||||
Usage: "include verbose statistics",
|
||||
|
||||
225
cmd/cluster_reroute.go
Normal file
225
cmd/cluster_reroute.go
Normal 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)
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -32,6 +32,7 @@ const (
|
||||
Ccluster
|
||||
Capi
|
||||
Cilm
|
||||
Cnode
|
||||
)
|
||||
|
||||
func complete(cmd *cli.Command, what int) {
|
||||
@@ -64,6 +65,8 @@ func complete(cmd *cli.Command, what int) {
|
||||
list = es.ApiPathNames()
|
||||
case Cilm:
|
||||
list, err = es.IlmNames(conf)
|
||||
case Cnode:
|
||||
list, err = es.NodeNames(conf)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
|
||||
32
cmd/ilm.go
32
cmd/ilm.go
@@ -36,7 +36,8 @@ func Ilm(conf *cfg.Config) *cli.Command {
|
||||
IlmStatus(conf),
|
||||
IlmList(conf),
|
||||
IlmShow(conf),
|
||||
IlmCreate(conf),
|
||||
IlmCreate(conf, false),
|
||||
IlmCreate(conf, true),
|
||||
IlmForecast(conf),
|
||||
IlmExplain(conf),
|
||||
},
|
||||
@@ -103,6 +104,15 @@ func IlmShow(conf *cfg.Config) *cli.Command {
|
||||
Usage: "show details about an index lifecycle 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 {
|
||||
policy := cmd.Args().Get(0)
|
||||
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{
|
||||
Name: "create",
|
||||
Aliases: []string{"+"},
|
||||
Usage: "create a index lifecycle policy",
|
||||
UsageText: "create [options] <policy>",
|
||||
Name: name,
|
||||
Aliases: []string{alias},
|
||||
Usage: usage,
|
||||
UsageText: name + " [options] <policy>",
|
||||
|
||||
Flags: []cli.Flag{
|
||||
&cli.StringFlag{
|
||||
|
||||
@@ -154,8 +154,8 @@ func IndexCreate(conf *cfg.Config, modify bool) *cli.Command {
|
||||
usage := "create a new index"
|
||||
|
||||
if modify {
|
||||
name = "modify"
|
||||
usage = "modify anindex"
|
||||
name = "update"
|
||||
usage = "update an index"
|
||||
}
|
||||
|
||||
return &cli.Command{
|
||||
|
||||
@@ -83,8 +83,8 @@ func IndexTemplateCreate(conf *cfg.Config, modify bool) *cli.Command {
|
||||
required := true
|
||||
|
||||
if modify {
|
||||
name = "modify"
|
||||
alias = "mod"
|
||||
name = "update"
|
||||
alias = "upd"
|
||||
required = false
|
||||
}
|
||||
|
||||
|
||||
45
cmd/node.go
45
cmd/node.go
@@ -18,6 +18,7 @@ package cmd
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"codeberg.org/scip/esctl/pkg/cfg"
|
||||
"codeberg.org/scip/esctl/pkg/es"
|
||||
@@ -34,6 +35,7 @@ func Node(conf *cfg.Config) *cli.Command {
|
||||
Commands: []*cli.Command{
|
||||
NodeList(conf),
|
||||
NodeShow(conf),
|
||||
NodeClients(conf),
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -57,10 +59,47 @@ func NodeShow(conf *cfg.Config) *cli.Command {
|
||||
Usage: "show details about a 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 {
|
||||
// FIXME: implement es.NodeShow()
|
||||
// return es.NodeShow(conf, cmd.Args().Get(0))
|
||||
return nil
|
||||
node := cmd.Args().Get(0)
|
||||
if node == "" {
|
||||
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)
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
102
cmd/root.go
102
cmd/root.go
@@ -42,7 +42,6 @@ func Finish(err error) int {
|
||||
|
||||
func Main() int {
|
||||
conf := cfg.NewConfig()
|
||||
tree := false
|
||||
|
||||
cmd := &cli.Command{
|
||||
Name: "esctl",
|
||||
@@ -64,13 +63,6 @@ func Main() int {
|
||||
Usage: "enable HTTP debugging",
|
||||
Destination: &conf.DebugHTTP,
|
||||
},
|
||||
&cli.BoolFlag{
|
||||
Name: "show-command-tree",
|
||||
Value: false,
|
||||
Usage: "generate a command tree",
|
||||
Destination: &tree,
|
||||
Hidden: true,
|
||||
},
|
||||
&cli.BoolFlag{
|
||||
Name: "align-ints",
|
||||
Aliases: []string{"I"},
|
||||
@@ -125,17 +117,10 @@ func Main() int {
|
||||
Version(conf),
|
||||
Debug(conf),
|
||||
HelpJsonPath(conf),
|
||||
HelpUsage(conf),
|
||||
},
|
||||
|
||||
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 len(os.Args) > 1 {
|
||||
return nil, err
|
||||
@@ -261,22 +246,81 @@ func Debug(conf *cfg.Config) *cli.Command {
|
||||
}
|
||||
}
|
||||
|
||||
func Tree(cmd *cli.Command) error {
|
||||
max := 20
|
||||
func HelpUsage(conf *cfg.Config) *cli.Command {
|
||||
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 {
|
||||
path := cmd.Path()
|
||||
command := path[len(path)-1]
|
||||
Flags: []cli.Flag{
|
||||
&cli.BoolFlag{
|
||||
Name: "hidden",
|
||||
Usage: "include hidden commands",
|
||||
Destination: &conf.Hidden,
|
||||
Aliases: []string{"H"},
|
||||
},
|
||||
},
|
||||
|
||||
if len(path) == 1 || command == "help" {
|
||||
return nil
|
||||
}
|
||||
Action: func(ctx context.Context, cmd *cli.Command) error {
|
||||
maxCommandWidth := 0
|
||||
|
||||
indent := strings.Repeat(" ", len(path[1:])-1)
|
||||
space := strings.Repeat(" ", max-(len(command)+len(indent)))
|
||||
// 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])
|
||||
|
||||
fmt.Printf("%s%s %s - %s\n", indent, command, space, cmd.Usage)
|
||||
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()
|
||||
command := path[len(path)-1]
|
||||
|
||||
if len(path) == 1 || command == "help" {
|
||||
return nil
|
||||
}
|
||||
|
||||
indent := strings.Repeat(" ", len(path[1:])-1)
|
||||
space := strings.Repeat(" ", maxCommandWidth-(len(command)+len(indent)))
|
||||
|
||||
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
1
go.mod
@@ -85,6 +85,7 @@ require (
|
||||
github.com/olekukonko/errors v1.2.0 // indirect
|
||||
github.com/olekukonko/ll v0.1.8 // 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/pretty v1.2.0 // indirect
|
||||
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect
|
||||
|
||||
2
go.sum
2
go.sum
@@ -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/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/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/go.mod h1:0CfEIISq7TuYL3j771MWULgwwjU+GofnZX9QAmXWZgo=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
|
||||
174
pkg/cfg/automate.go
Normal file
174
pkg/cfg/automate.go
Normal 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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -17,27 +17,33 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
package cfg
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/elastic/elastic-transport-go/v8/elastictransport"
|
||||
"github.com/elastic/go-elasticsearch/v9"
|
||||
"golang.org/x/term"
|
||||
"gopkg.in/yaml.v3"
|
||||
)
|
||||
|
||||
// used in general config struct
|
||||
type Cluster struct {
|
||||
Uri, User, Pass string
|
||||
client *elasticsearch.TypedClient
|
||||
Default bool
|
||||
Name, Uri, User, Pass, Token string
|
||||
client *elasticsearch.TypedClient
|
||||
Default, DebugHTTP bool
|
||||
}
|
||||
|
||||
// used just for writing back to the config file
|
||||
type ClusterConfig struct {
|
||||
Uri, User, Pass string
|
||||
Default bool
|
||||
Uri, User, Pass, Token string
|
||||
Default bool
|
||||
}
|
||||
|
||||
// to write the config, we avoid all other config settings
|
||||
@@ -45,17 +51,96 @@ type WriteConfig struct {
|
||||
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 {
|
||||
if cluster.client == nil {
|
||||
fmt.Println("no current cluster, use 'esctl cluster switch <name>' to set one")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
if err := cluster.CheckAuth(); err != nil {
|
||||
fmt.Printf("Error: %s", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
return cluster.client
|
||||
}
|
||||
|
||||
func (cluster *Cluster) SetClient(client *elasticsearch.TypedClient) {
|
||||
cluster.client = client
|
||||
// add authentication to es client, if not yet done
|
||||
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)
|
||||
@@ -72,6 +157,7 @@ func (conf *Config) SwitchCluster(name string) error {
|
||||
Uri: cluster.Uri,
|
||||
User: cluster.User,
|
||||
Pass: cluster.Pass,
|
||||
Token: cluster.Token,
|
||||
Default: false,
|
||||
}
|
||||
|
||||
@@ -97,29 +183,75 @@ func (conf *Config) SwitchCluster(name string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (conf *Config) SetupES() error {
|
||||
// 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")
|
||||
// We do NOT use go-elasticsearch to check for cluster reachability,
|
||||
// because at this stage, auth may not have been configured. So
|
||||
// instead we just connect to the cluster using plan net/http, ignore
|
||||
// HTTP response status and return true if we could just reach ith
|
||||
func (cluster *Cluster) IsReachable() (bool, error) {
|
||||
ctx, cancel := context.WithTimeout(
|
||||
context.Background(),
|
||||
time.Duration(500)*time.Millisecond)
|
||||
defer cancel()
|
||||
|
||||
for _, cluster := range conf.Clusters {
|
||||
es, err := elasticsearch.NewTyped(
|
||||
elasticsearch.WithAddresses(cluster.Uri),
|
||||
elasticsearch.WithBasicAuth(cluster.User, cluster.Pass),
|
||||
elasticsearch.WithTransportOptions(
|
||||
conf.getTransport(),
|
||||
elastictransport.WithHeader(headers),
|
||||
),
|
||||
)
|
||||
req, err := http.NewRequestWithContext(
|
||||
ctx,
|
||||
"GET",
|
||||
cluster.Uri,
|
||||
nil,
|
||||
)
|
||||
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to setup elasticsearch connection: %w", err)
|
||||
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 {
|
||||
return fmt.Errorf("failed to setup elasticsearch connection: %w", err)
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
cluster.SetClient(es)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -28,7 +28,7 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
Version string = `v0.0.22`
|
||||
Version string = `v0.0.25`
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -103,6 +103,9 @@ type Config struct {
|
||||
Tag string // api ls: -t
|
||||
|
||||
Ilm Ilm // ilm create
|
||||
|
||||
FromNode, ToNode string // cluster reroute move: -f + -t
|
||||
AllowPrimary, AcceptDataLoss bool // cluster reroute cancel: -p,-a
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
@@ -150,6 +153,9 @@ func (conf *Config) Init() error {
|
||||
conf.DefaultCluster.Default = true
|
||||
}
|
||||
} else {
|
||||
// load auto config, if any
|
||||
auto := NewAutomator()
|
||||
|
||||
// we need to determine ourselfes
|
||||
if len(conf.Clusters) == 1 {
|
||||
// ok, just one cluster configured, use this, of course
|
||||
@@ -158,6 +164,14 @@ func (conf *Config) Init() error {
|
||||
conf.CurrentCluster = name
|
||||
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 {
|
||||
// multiple ones exists, look if one is set as default
|
||||
for name, cluster := range conf.Clusters {
|
||||
@@ -210,18 +224,17 @@ func (conf *Config) PrintDebug() {
|
||||
|
||||
func (conf *Config) LoadEnv() error {
|
||||
cluster := Cluster{
|
||||
Uri: os.Getenv("ES_URI"),
|
||||
User: os.Getenv("ES_USER"),
|
||||
Pass: os.Getenv("ES_PASS"),
|
||||
Uri: os.Getenv("ES_URI"),
|
||||
User: os.Getenv("ES_USER"),
|
||||
Pass: os.Getenv("ES_PASS"),
|
||||
Token: os.Getenv("ES_TOKEN"),
|
||||
}
|
||||
|
||||
switch {
|
||||
case cluster.Uri == "":
|
||||
return errors.New("ES_URI unset")
|
||||
case cluster.User == "":
|
||||
return errors.New("ES_USER unset")
|
||||
case cluster.Pass == "":
|
||||
return errors.New("ES_PASS unset")
|
||||
case cluster.User == "" || cluster.Token == "":
|
||||
return errors.New("ES_USER and ES_TOKEN unset")
|
||||
}
|
||||
|
||||
conf.Clusters["default"] = &cluster
|
||||
|
||||
@@ -42,6 +42,7 @@ type Ilm struct {
|
||||
DeleteSearchableSnapshots bool
|
||||
|
||||
MinAge, FromPhase, Within string // forecast
|
||||
Tree bool
|
||||
}
|
||||
|
||||
func (cfg *Ilm) HaveHot() bool {
|
||||
|
||||
38
pkg/cfg/term.go
Normal file
38
pkg/cfg/term.go
Normal 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
|
||||
}
|
||||
@@ -18,13 +18,10 @@ package cfg
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/tls"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"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
|
||||
@@ -62,17 +59,3 @@ func (t *DebugTransport) RoundTrip(req *http.Request) (*http.Response, error) {
|
||||
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -41,12 +41,10 @@ import (
|
||||
"github.com/charmbracelet/lipgloss"
|
||||
"github.com/chzyer/readline"
|
||||
"github.com/go-openapi/spec"
|
||||
"golang.org/x/term"
|
||||
)
|
||||
|
||||
const (
|
||||
DefaultMargin = 4
|
||||
intro = `Input format: verb path [data]"
|
||||
intro = `Input format: verb path [data]"
|
||||
|
||||
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("accept", "application/json")
|
||||
req.Header.Add("Authorization", "Basic "+encodeAuth(conf.DefaultCluster.User, conf.DefaultCluster.Pass))
|
||||
|
||||
// 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))
|
||||
}
|
||||
|
||||
// actually execute the request
|
||||
resp, err := client.Do(req)
|
||||
@@ -318,19 +326,19 @@ func ApiShow(conf *cfg.Config, showpath, verb string) error {
|
||||
|
||||
cleanMarkup := regexp.MustCompile(`<[^<>]+>`)
|
||||
|
||||
width := getTermWidth()
|
||||
width := cfg.GetTermWidth()
|
||||
params := getApiParameters(op, showpath, width)
|
||||
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().
|
||||
Bold(true)
|
||||
var paragraph = lipgloss.NewStyle().
|
||||
MarginBottom(1).
|
||||
MarginLeft(DefaultMargin)
|
||||
MarginLeft(cfg.DefaultMargin)
|
||||
var boldparagraph = lipgloss.NewStyle().
|
||||
MarginBottom(1).
|
||||
MarginLeft(DefaultMargin).
|
||||
MarginLeft(cfg.DefaultMargin).
|
||||
Bold(true)
|
||||
var indentparagraph = lipgloss.NewStyle().
|
||||
MarginBottom(1).
|
||||
@@ -413,7 +421,7 @@ func getApiParameters(op *Op, path string, width int) Params {
|
||||
}
|
||||
} else {
|
||||
par := Param{
|
||||
Description: string(markdown.Render(param.Description, width, DefaultMargin)),
|
||||
Description: string(markdown.Render(param.Description, width, cfg.DefaultMargin)),
|
||||
Param: param.Name,
|
||||
}
|
||||
|
||||
@@ -517,14 +525,3 @@ func findOperation(item spec.PathItemProps) []*Op {
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
59
pkg/es/cluser_health_report.go
Normal file
59
pkg/es/cluser_health_report.go
Normal 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
|
||||
}
|
||||
@@ -17,7 +17,6 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
package es
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strings"
|
||||
@@ -26,44 +25,40 @@ import (
|
||||
"codeberg.org/scip/esctl/pkg/cfg"
|
||||
"codeberg.org/scip/esctl/pkg/printer"
|
||||
"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"
|
||||
"github.com/elastic/go-elasticsearch/v9/typedapi/core/info"
|
||||
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
|
||||
)
|
||||
|
||||
type ClusterIndices map[string]map[string]*types.IndicesRecord
|
||||
|
||||
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
|
||||
|
||||
for name, cluster := range conf.Clusters {
|
||||
reachable := "no"
|
||||
current := "no"
|
||||
errmsg := ""
|
||||
|
||||
_, err := cluster.ES().Cluster.Health().
|
||||
Do(context.Background())
|
||||
online, err := cluster.IsReachable()
|
||||
|
||||
if err == nil {
|
||||
if online {
|
||||
reachable = printer.Colorize(conf, "green", "reachable")
|
||||
}
|
||||
|
||||
if cluster.Default {
|
||||
current = printer.Colorize(conf, "green", "yes")
|
||||
|
||||
if err != nil {
|
||||
if !online {
|
||||
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++
|
||||
}
|
||||
|
||||
@@ -78,135 +73,146 @@ func ClusterList(conf *cfg.Config) error {
|
||||
|
||||
// We're using goroutines here to parallelize API requests, since we
|
||||
// have to do 3 of'em for each cluster. This speeds things up.
|
||||
func ClusterStatus(conf *cfg.Config) error {
|
||||
clusters := []string{}
|
||||
gocount := 5
|
||||
func getClusterStatus(conf *cfg.Config) (*apiResponse, error) {
|
||||
gocount := 6
|
||||
if conf.Verbose {
|
||||
gocount++
|
||||
}
|
||||
|
||||
if conf.All {
|
||||
for key := range conf.Clusters {
|
||||
clusters = append(clusters, key)
|
||||
}
|
||||
} else {
|
||||
clusters = []string{"default"}
|
||||
es := conf.DefaultCluster.ES()
|
||||
|
||||
responses := make(chan apiResponse, gocount)
|
||||
wg := &sync.WaitGroup{}
|
||||
|
||||
wg.Add(gocount)
|
||||
go getApiData(conf, es, wg, responses, "health")
|
||||
go getApiData(conf, es, wg, responses, "healthreport")
|
||||
go getApiData(conf, es, wg, responses, "info")
|
||||
go getApiData(conf, es, wg, responses, "ccr")
|
||||
go getApiData(conf, es, wg, responses, "indices")
|
||||
go getApiData(conf, es, wg, responses, "tasks")
|
||||
|
||||
if conf.Verbose {
|
||||
go getApiData(conf, es, wg, responses, "stats")
|
||||
}
|
||||
|
||||
for _, cluster := range clusters {
|
||||
es := conf.DefaultCluster.ES()
|
||||
if cluster != "default" {
|
||||
es = conf.Clusters[cluster].ES()
|
||||
wg.Wait()
|
||||
|
||||
all := apiResponse{}
|
||||
|
||||
for i := 0; i < gocount; i++ {
|
||||
r := <-responses
|
||||
|
||||
if r.error != nil {
|
||||
return nil, r.error
|
||||
}
|
||||
|
||||
responses := make(chan apiResponse, gocount)
|
||||
wg := &sync.WaitGroup{}
|
||||
|
||||
wg.Add(gocount)
|
||||
go getApiData(es, wg, responses, "health")
|
||||
go getApiData(es, wg, responses, "info")
|
||||
go getApiData(es, wg, responses, "ccrstats")
|
||||
go getApiData(es, wg, responses, "indices")
|
||||
go getApiData(es, wg, responses, "tasks")
|
||||
|
||||
if conf.Verbose {
|
||||
go getApiData(es, wg, responses, "stats")
|
||||
switch r.which {
|
||||
case ResponseHealth:
|
||||
all.health = r.health
|
||||
case ResponseCcr:
|
||||
all.ccr = r.ccr
|
||||
case ResponseInfo:
|
||||
all.info = r.info
|
||||
case ResponseStats:
|
||||
all.stats = r.stats
|
||||
case ResponseIndices:
|
||||
all.indices = r.indices
|
||||
case ResponseTasks:
|
||||
all.tasks = r.tasks
|
||||
case ResponseHealthReport:
|
||||
all.healthreport = r.healthreport
|
||||
}
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
return &all, nil
|
||||
}
|
||||
|
||||
var clusterhealth *health.Response
|
||||
var info *info.Response
|
||||
var ccrstats *stats.Response
|
||||
var clusterstats *clusterstats.Response
|
||||
var indexstats *indices.Response
|
||||
var taskstatus *tasks.Response
|
||||
func ClusterStatus(conf *cfg.Config) error {
|
||||
res, err := getClusterStatus(conf)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for i := 0; i < gocount; i++ {
|
||||
r := <-responses
|
||||
slog.Debug("ES result", "cluster health", res.health)
|
||||
|
||||
if r.error != nil {
|
||||
return r.error
|
||||
}
|
||||
isleader := len(res.ccr.AutoFollowStats.AutoFollowedClusters) == 0
|
||||
|
||||
switch r.which {
|
||||
case ResponseHealth:
|
||||
clusterhealth = r.health
|
||||
case ResponseCcr:
|
||||
ccrstats = r.ccr
|
||||
case ResponseInfo:
|
||||
info = r.info
|
||||
case ResponseStats:
|
||||
clusterstats = r.stats
|
||||
case ResponseIndices:
|
||||
indexstats = r.indices
|
||||
case ResponseTasks:
|
||||
taskstatus = r.tasks
|
||||
ccrfollowing := ""
|
||||
if len(res.ccr.AutoFollowStats.AutoFollowedClusters) > 0 {
|
||||
// is following another cluster
|
||||
ccrfollowing = fmt.Sprintf("%s (%d/%d)",
|
||||
res.ccr.AutoFollowStats.AutoFollowedClusters[0].ClusterName,
|
||||
res.ccr.AutoFollowStats.NumberOfSuccessfulFollowIndices,
|
||||
res.ccr.AutoFollowStats.NumberOfFailedFollowIndices,
|
||||
)
|
||||
}
|
||||
|
||||
// look for red indices, if any
|
||||
redindices := 0
|
||||
for _, index := range *res.indices {
|
||||
if *index.Health == "red" {
|
||||
redindices++
|
||||
}
|
||||
}
|
||||
|
||||
// look for long running tasks
|
||||
longtasks := 0
|
||||
for _, task := range *res.tasks {
|
||||
if strings.Contains(*task.RunningTime, "d") {
|
||||
longtasks++
|
||||
}
|
||||
}
|
||||
|
||||
table := printer.NewTable(conf, 2, 7)
|
||||
table.Addheaders(conf.DefaultCluster.Name, "status")
|
||||
|
||||
table.Entries = [][]string{
|
||||
{"Cluster Name", res.health.ClusterName},
|
||||
{"ES Status", printer.Colorize(conf, res.health.Status.Name, res.health.Status.Name)},
|
||||
{"ES Version", res.info.Version.Int},
|
||||
{"Is Leader", fmt.Sprintf("%t", isleader)},
|
||||
{"Active Shards", fmt.Sprintf("%d", res.health.ActiveShards)},
|
||||
{"Active Primary Shards", fmt.Sprintf("%d", res.health.ActivePrimaryShards)},
|
||||
{"Unassigned Shards", fmt.Sprintf("%d", res.health.UnassignedShards)},
|
||||
{"Unassigned Primary Shards", fmt.Sprintf("%d", res.health.UnassignedPrimaryShards)},
|
||||
{"Pending Tasks", fmt.Sprintf("%d", res.health.NumberOfPendingTasks)},
|
||||
{"Nodes", fmt.Sprintf("%d", res.health.NumberOfNodes)},
|
||||
{"Red Indices", fmt.Sprintf("%d", redindices)},
|
||||
{"Long Running Tasks", fmt.Sprintf("%d", longtasks)},
|
||||
}
|
||||
|
||||
if !isleader {
|
||||
table.Entries = append(table.Entries, [][]string{
|
||||
{"AutoFollow (success/failed indices)", ccrfollowing},
|
||||
{"Followed Indices", fmt.Sprintf("%d", len(res.ccr.FollowStats.Indices))},
|
||||
}...)
|
||||
}
|
||||
|
||||
if conf.Verbose {
|
||||
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, ",")})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
slog.Debug("ES result", "cluster health", clusterhealth)
|
||||
|
||||
isleader := len(ccrstats.AutoFollowStats.AutoFollowedClusters) == 0
|
||||
|
||||
ccrfollowing := ""
|
||||
if len(ccrstats.AutoFollowStats.AutoFollowedClusters) > 0 {
|
||||
// is following another cluster
|
||||
ccrfollowing = fmt.Sprintf("%s (%d/%d)",
|
||||
ccrstats.AutoFollowStats.AutoFollowedClusters[0].ClusterName,
|
||||
ccrstats.AutoFollowStats.NumberOfSuccessfulFollowIndices,
|
||||
ccrstats.AutoFollowStats.NumberOfFailedFollowIndices,
|
||||
)
|
||||
}
|
||||
|
||||
// look for red indices, if any
|
||||
redindices := 0
|
||||
for _, index := range *indexstats {
|
||||
if *index.Health == "red" {
|
||||
redindices++
|
||||
}
|
||||
}
|
||||
|
||||
// look for long running tasks
|
||||
longtasks := 0
|
||||
for _, task := range *taskstatus {
|
||||
if strings.Contains(*task.RunningTime, "d") {
|
||||
longtasks++
|
||||
}
|
||||
}
|
||||
|
||||
table := printer.NewTable(conf, 2, 7)
|
||||
table.Addheaders(cluster, "status")
|
||||
|
||||
table.Entries = [][]string{
|
||||
{"Cluster Name", clusterhealth.ClusterName},
|
||||
{"ES Status", printer.Colorize(conf, clusterhealth.Status.Name, clusterhealth.Status.Name)},
|
||||
{"ES Version", info.Version.Int},
|
||||
{"Is Leader", fmt.Sprintf("%t", isleader)},
|
||||
{"Active Shards", fmt.Sprintf("%d", clusterhealth.ActiveShards)},
|
||||
{"Active Primary Shards", fmt.Sprintf("%d", clusterhealth.ActivePrimaryShards)},
|
||||
{"Unassigned Shards", fmt.Sprintf("%d", clusterhealth.UnassignedShards)},
|
||||
{"Unassigned Primary Shards", fmt.Sprintf("%d", clusterhealth.UnassignedPrimaryShards)},
|
||||
{"Pending Tasks", fmt.Sprintf("%d", clusterhealth.NumberOfPendingTasks)},
|
||||
{"Nodes", fmt.Sprintf("%d", clusterhealth.NumberOfNodes)},
|
||||
{"Red Indices", fmt.Sprintf("%d", redindices)},
|
||||
{"Long Running Tasks", fmt.Sprintf("%d", longtasks)},
|
||||
}
|
||||
|
||||
if !isleader {
|
||||
table.Entries = append(table.Entries, [][]string{
|
||||
{"AutoFollow (success/failed indices)", ccrfollowing},
|
||||
{"Followed Indices", fmt.Sprintf("%d", len(ccrstats.FollowStats.Indices))},
|
||||
}...)
|
||||
}
|
||||
|
||||
if conf.Verbose {
|
||||
table = gatherClusterStats(conf, clusterstats, table)
|
||||
}
|
||||
|
||||
if err := table.Print(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := table.Print(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
125
pkg/es/cluster_reroute.go
Normal file
125
pkg/es/cluster_reroute.go
Normal 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
48
pkg/es/debug.go
Normal 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
|
||||
}
|
||||
180
pkg/es/ilm.go
180
pkg/es/ilm.go
@@ -121,6 +121,10 @@ func IlmShow(conf *cfg.Config, policy string) error {
|
||||
return errors.New("no ilm policy retrieved")
|
||||
}
|
||||
|
||||
if conf.Ilm.Tree {
|
||||
return IlmShowTree(conf, ilm.Policy)
|
||||
}
|
||||
|
||||
table := printer.NewTable(conf, 2, 5)
|
||||
table.Addheaders("ilm policy setting", "value")
|
||||
|
||||
@@ -135,6 +139,101 @@ func IlmShow(conf *cfg.Config, policy string) error {
|
||||
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 {
|
||||
if phase == nil {
|
||||
return ""
|
||||
@@ -220,6 +319,24 @@ func IlmExplain(conf *cfg.Config, index 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)
|
||||
cfg := conf.Ilm
|
||||
|
||||
@@ -231,6 +348,16 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
||||
rollover := &types.RolloverAction{}
|
||||
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 != "" {
|
||||
hot.MinAge = cfg.HotMinAge
|
||||
}
|
||||
@@ -255,12 +382,22 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
||||
|
||||
hot.Actions = actions.IlmActionsCaster()
|
||||
phases.PhasesCaster().Hot = &hot
|
||||
} else {
|
||||
if policy != nil {
|
||||
// update
|
||||
phases.PhasesCaster().Hot = policy.Phases.Hot
|
||||
}
|
||||
}
|
||||
|
||||
if cfg.HaveWarm() {
|
||||
warm := types.Phase{}
|
||||
var actions types.IlmActionsVariant = esdsl.NewIlmActions()
|
||||
|
||||
if policy != nil {
|
||||
// update
|
||||
actions = policy.Phases.Warm.Actions
|
||||
}
|
||||
|
||||
if cfg.WarmForceMerge != 0 {
|
||||
actions.IlmActionsCaster().Forcemerge = &types.ForceMergeAction{MaxNumSegments: cfg.WarmForceMerge}
|
||||
}
|
||||
@@ -279,12 +416,22 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
||||
|
||||
warm.Actions = actions.IlmActionsCaster()
|
||||
phases.PhasesCaster().Warm = &warm
|
||||
} else {
|
||||
if policy != nil && policy.Phases.Warm != nil {
|
||||
// update
|
||||
phases.PhasesCaster().Warm = policy.Phases.Warm
|
||||
}
|
||||
}
|
||||
|
||||
if cfg.HaveCold() {
|
||||
cold := types.Phase{}
|
||||
var actions types.IlmActionsVariant = esdsl.NewIlmActions()
|
||||
|
||||
if policy != nil {
|
||||
// update
|
||||
actions = policy.Phases.Cold.Actions
|
||||
}
|
||||
|
||||
if cfg.ColdForceMerge != 0 {
|
||||
actions.IlmActionsCaster().Forcemerge = &types.ForceMergeAction{MaxNumSegments: cfg.ColdForceMerge}
|
||||
}
|
||||
@@ -304,12 +451,22 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
||||
|
||||
cold.Actions = actions.IlmActionsCaster()
|
||||
phases.PhasesCaster().Cold = &cold
|
||||
} else {
|
||||
if policy != nil && policy.Phases.Cold != nil {
|
||||
// update
|
||||
phases.PhasesCaster().Cold = policy.Phases.Cold
|
||||
}
|
||||
}
|
||||
|
||||
if cfg.HaveFrozen() {
|
||||
froze := types.Phase{}
|
||||
var actions types.IlmActionsVariant = esdsl.NewIlmActions()
|
||||
|
||||
if policy != nil {
|
||||
// update
|
||||
actions = policy.Phases.Frozen.Actions
|
||||
}
|
||||
|
||||
if cfg.FrozenMinAge != "" {
|
||||
froze.MinAge = cfg.FrozenMinAge
|
||||
}
|
||||
@@ -321,6 +478,11 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
||||
|
||||
froze.Actions = actions.IlmActionsCaster()
|
||||
phases.PhasesCaster().Frozen = &froze
|
||||
} else {
|
||||
if policy != nil && policy.Phases.Frozen != nil {
|
||||
// update
|
||||
phases.PhasesCaster().Frozen = policy.Phases.Frozen
|
||||
}
|
||||
}
|
||||
|
||||
if cfg.HaveDelete() {
|
||||
@@ -328,6 +490,11 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
||||
delete := types.DeleteAction{}
|
||||
var actions types.IlmActionsVariant = esdsl.NewIlmActions()
|
||||
|
||||
if policy != nil {
|
||||
// update
|
||||
actions = policy.Phases.Delete.Actions
|
||||
}
|
||||
|
||||
if cfg.DeleteMinAge != "" {
|
||||
del.MinAge = cfg.DeleteMinAge
|
||||
}
|
||||
@@ -339,16 +506,21 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
||||
actions.IlmActionsCaster().Delete = &delete
|
||||
del.Actions = actions.IlmActionsCaster()
|
||||
phases.PhasesCaster().Delete = &del
|
||||
} else {
|
||||
if policy != nil && policy.Phases.Delete != nil {
|
||||
// update
|
||||
phases.PhasesCaster().Delete = policy.Phases.Delete
|
||||
}
|
||||
}
|
||||
|
||||
put := &putlifecycle.Request{}
|
||||
policy := &types.IlmPolicy{}
|
||||
policy.IlmPolicyCaster().Phases = *phases.PhasesCaster()
|
||||
put.Policy = policy
|
||||
newpolicy := &types.IlmPolicy{}
|
||||
newpolicy.IlmPolicyCaster().Phases = *phases.PhasesCaster()
|
||||
put.Policy = newpolicy
|
||||
|
||||
ilm.Request(put)
|
||||
|
||||
_, err := ilm.Do(context.Background())
|
||||
_, err = ilm.Do(context.Background())
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed create ilm policy: %s", esErrorString(err))
|
||||
}
|
||||
|
||||
@@ -190,9 +190,9 @@ func getIlmPhaseData(conf *cfg.Config) ([]PhaseData, error) {
|
||||
wg := &sync.WaitGroup{}
|
||||
wg.Add(3)
|
||||
|
||||
go getApiData(conf.DefaultCluster.ES(), wg, responses, "indicesbytes")
|
||||
go getApiData(conf.DefaultCluster.ES(), wg, responses, "explain")
|
||||
go getApiData(conf.DefaultCluster.ES(), wg, responses, "policies")
|
||||
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "indicesbytes")
|
||||
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "explain")
|
||||
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "policies")
|
||||
|
||||
wg.Wait()
|
||||
|
||||
|
||||
145
pkg/es/node.go
145
pkg/es/node.go
@@ -20,9 +20,12 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"codeberg.org/scip/esctl/pkg/cfg"
|
||||
"codeberg.org/scip/esctl/pkg/printer"
|
||||
"github.com/dustin/go-humanize"
|
||||
)
|
||||
|
||||
func NodeList(conf *cfg.Config) error {
|
||||
@@ -56,3 +59,145 @@ func NodeList(conf *cfg.Config) error {
|
||||
|
||||
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()
|
||||
}
|
||||
|
||||
@@ -18,10 +18,10 @@ package es
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"codeberg.org/scip/esctl/pkg/cfg"
|
||||
"github.com/elastic/go-elasticsearch/v9"
|
||||
"github.com/elastic/go-elasticsearch/v9/typedapi/cat/indices"
|
||||
"github.com/elastic/go-elasticsearch/v9/typedapi/cat/tasks"
|
||||
@@ -43,12 +43,14 @@ const (
|
||||
ResponseTasks
|
||||
ResponseExplain
|
||||
ResponseLifecycle
|
||||
ResponseHealthReport
|
||||
)
|
||||
|
||||
type apiResponse struct {
|
||||
error error
|
||||
info *info.Response
|
||||
health *health.Response
|
||||
healthreport *HealthReport
|
||||
ccr *stats.Response
|
||||
stats *clusterstats.Response
|
||||
indices *indices.Response
|
||||
@@ -59,12 +61,16 @@ type apiResponse struct {
|
||||
which int
|
||||
}
|
||||
|
||||
func getApiData(es *elasticsearch.TypedClient, wg *sync.WaitGroup,
|
||||
reschan chan apiResponse, which string) {
|
||||
func getApiData(
|
||||
conf *cfg.Config,
|
||||
es *elasticsearch.TypedClient,
|
||||
wg *sync.WaitGroup,
|
||||
reschan chan apiResponse,
|
||||
which string) {
|
||||
defer wg.Done()
|
||||
|
||||
ar := apiResponse{}
|
||||
arerr := errors.New("")
|
||||
var arerr error
|
||||
|
||||
switch which {
|
||||
case "health":
|
||||
@@ -75,6 +81,13 @@ func getApiData(es *elasticsearch.TypedClient, wg *sync.WaitGroup,
|
||||
ar.which = ResponseHealth
|
||||
arerr = err
|
||||
|
||||
case "healthreport":
|
||||
report, err := getHealthReport(conf)
|
||||
|
||||
ar.healthreport = report
|
||||
ar.which = ResponseHealthReport
|
||||
arerr = err
|
||||
|
||||
case "info":
|
||||
res, err := es.Info().
|
||||
Do(context.Background())
|
||||
@@ -83,7 +96,7 @@ func getApiData(es *elasticsearch.TypedClient, wg *sync.WaitGroup,
|
||||
ar.which = ResponseInfo
|
||||
arerr = err
|
||||
|
||||
case "ccrstats":
|
||||
case "ccr":
|
||||
res, err := es.Ccr.Stats().
|
||||
Do(context.Background())
|
||||
|
||||
|
||||
@@ -26,7 +26,6 @@ import (
|
||||
|
||||
"codeberg.org/scip/esctl/pkg/cfg"
|
||||
"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/esdsl"
|
||||
"github.com/elastic/go-elasticsearch/v9/typedapi/indices/validatequery"
|
||||
@@ -149,25 +148,6 @@ func validateSearch(conf *cfg.Config, queries []string) error {
|
||||
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 {
|
||||
res, err := search.
|
||||
From(conf.From).
|
||||
|
||||
@@ -28,6 +28,7 @@ import (
|
||||
"github.com/olekukonko/tablewriter"
|
||||
"github.com/olekukonko/tablewriter/renderer"
|
||||
"github.com/olekukonko/tablewriter/tw"
|
||||
"github.com/seeruk/go-wordwrap"
|
||||
"gopkg.in/yaml.v3"
|
||||
)
|
||||
|
||||
@@ -38,10 +39,11 @@ type Table struct {
|
||||
|
||||
lenHeaders []int
|
||||
alignInts bool
|
||||
maxwidth int
|
||||
}
|
||||
|
||||
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.Entries = make([][]string, rows)
|
||||
@@ -118,17 +120,31 @@ func (data *Table) PrintTSV() error {
|
||||
}
|
||||
|
||||
for _, entries := range data.Entries {
|
||||
currentWidth := 0
|
||||
for idx, entry := range entries {
|
||||
length := visibleLen(entry)
|
||||
|
||||
if data.lenHeaders[idx] < length {
|
||||
data.lenHeaders[idx] = length
|
||||
if length > currentWidth+data.maxwidth {
|
||||
data.lenHeaders[idx] = data.maxwidth - currentWidth
|
||||
} else {
|
||||
data.lenHeaders[idx] = length
|
||||
}
|
||||
}
|
||||
|
||||
currentWidth += data.lenHeaders[idx]
|
||||
}
|
||||
}
|
||||
|
||||
// output
|
||||
// output headers
|
||||
for idx, header := range data.Headers {
|
||||
fmt.Print(header, strings.Repeat(" ", data.lenHeaders[idx]-visibleLen(header)))
|
||||
if idx+1 != len(data.Headers) {
|
||||
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 {
|
||||
fmt.Print(" ")
|
||||
}
|
||||
@@ -136,14 +152,37 @@ func (data *Table) PrintTSV() error {
|
||||
fmt.Println()
|
||||
|
||||
for _, entries := range data.Entries {
|
||||
currentWidth := 0
|
||||
|
||||
for idx, entry := range entries {
|
||||
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 {
|
||||
// align right
|
||||
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))
|
||||
} else {
|
||||
// no padding for last entry
|
||||
fmt.Print(entry)
|
||||
}
|
||||
|
||||
if idx < len(data.Headers)-1 {
|
||||
|
||||
Reference in New Issue
Block a user