mirror of
https://codeberg.org/scip/esctl.git
synced 2026-08-24 18:44:17 +02:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b573d63c54 | ||
|
|
e3b96b3e3b | ||
|
|
6436bcfb22 |
@@ -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
|
||||||
@@ -549,6 +549,7 @@ index - manage indicies
|
|||||||
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
|
||||||
|
|||||||
@@ -61,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",
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
45
cmd/node.go
45
cmd/node.go
@@ -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)
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
1
go.mod
1
go.mod
@@ -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
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/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=
|
||||||
|
|||||||
@@ -17,27 +17,33 @@ 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
|
||||||
}
|
}
|
||||||
|
|
||||||
// to write the config, we avoid all other config settings
|
// to write the config, we avoid all other config settings
|
||||||
@@ -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,23 +183,61 @@ 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) SetupES() error {
|
||||||
|
for name, cluster := range conf.Clusters {
|
||||||
|
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)
|
||||||
|
|||||||
@@ -28,7 +28,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
Version string = `v0.0.23`
|
Version string = `v0.0.24`
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -213,18 +213,17 @@ func (conf *Config) PrintDebug() {
|
|||||||
|
|
||||||
func (conf *Config) LoadEnv() error {
|
func (conf *Config) LoadEnv() error {
|
||||||
cluster := Cluster{
|
cluster := Cluster{
|
||||||
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
|
||||||
|
|||||||
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 (
|
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)
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -41,12 +41,10 @@ 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")
|
||||||
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
|
// 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
|
|
||||||
}
|
|
||||||
|
|||||||
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
|
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,135 +73,146 @@ 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 {
|
es := conf.DefaultCluster.ES()
|
||||||
for key := range conf.Clusters {
|
|
||||||
clusters = append(clusters, key)
|
responses := make(chan apiResponse, gocount)
|
||||||
}
|
wg := &sync.WaitGroup{}
|
||||||
} else {
|
|
||||||
clusters = []string{"default"}
|
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 {
|
wg.Wait()
|
||||||
es := conf.DefaultCluster.ES()
|
|
||||||
if cluster != "default" {
|
all := apiResponse{}
|
||||||
es = conf.Clusters[cluster].ES()
|
|
||||||
|
for i := 0; i < gocount; i++ {
|
||||||
|
r := <-responses
|
||||||
|
|
||||||
|
if r.error != nil {
|
||||||
|
return nil, r.error
|
||||||
}
|
}
|
||||||
|
|
||||||
responses := make(chan apiResponse, gocount)
|
switch r.which {
|
||||||
wg := &sync.WaitGroup{}
|
case ResponseHealth:
|
||||||
|
all.health = r.health
|
||||||
wg.Add(gocount)
|
case ResponseCcr:
|
||||||
go getApiData(es, wg, responses, "health")
|
all.ccr = r.ccr
|
||||||
go getApiData(es, wg, responses, "info")
|
case ResponseInfo:
|
||||||
go getApiData(es, wg, responses, "ccrstats")
|
all.info = r.info
|
||||||
go getApiData(es, wg, responses, "indices")
|
case ResponseStats:
|
||||||
go getApiData(es, wg, responses, "tasks")
|
all.stats = r.stats
|
||||||
|
case ResponseIndices:
|
||||||
if conf.Verbose {
|
all.indices = r.indices
|
||||||
go getApiData(es, wg, responses, "stats")
|
case ResponseTasks:
|
||||||
|
all.tasks = r.tasks
|
||||||
|
case ResponseHealthReport:
|
||||||
|
all.healthreport = r.healthreport
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
wg.Wait()
|
return &all, nil
|
||||||
|
}
|
||||||
|
|
||||||
var clusterhealth *health.Response
|
func ClusterStatus(conf *cfg.Config) error {
|
||||||
var info *info.Response
|
res, err := getClusterStatus(conf)
|
||||||
var ccrstats *stats.Response
|
if err != nil {
|
||||||
var clusterstats *clusterstats.Response
|
return err
|
||||||
var indexstats *indices.Response
|
}
|
||||||
var taskstatus *tasks.Response
|
|
||||||
|
|
||||||
for i := 0; i < gocount; i++ {
|
slog.Debug("ES result", "cluster health", res.health)
|
||||||
r := <-responses
|
|
||||||
|
|
||||||
if r.error != nil {
|
isleader := len(res.ccr.AutoFollowStats.AutoFollowedClusters) == 0
|
||||||
return r.error
|
|
||||||
}
|
|
||||||
|
|
||||||
switch r.which {
|
ccrfollowing := ""
|
||||||
case ResponseHealth:
|
if len(res.ccr.AutoFollowStats.AutoFollowedClusters) > 0 {
|
||||||
clusterhealth = r.health
|
// is following another cluster
|
||||||
case ResponseCcr:
|
ccrfollowing = fmt.Sprintf("%s (%d/%d)",
|
||||||
ccrstats = r.ccr
|
res.ccr.AutoFollowStats.AutoFollowedClusters[0].ClusterName,
|
||||||
case ResponseInfo:
|
res.ccr.AutoFollowStats.NumberOfSuccessfulFollowIndices,
|
||||||
info = r.info
|
res.ccr.AutoFollowStats.NumberOfFailedFollowIndices,
|
||||||
case ResponseStats:
|
)
|
||||||
clusterstats = r.stats
|
}
|
||||||
case ResponseIndices:
|
|
||||||
indexstats = r.indices
|
// look for red indices, if any
|
||||||
case ResponseTasks:
|
redindices := 0
|
||||||
taskstatus = r.tasks
|
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)
|
if err := table.Print(); err != nil {
|
||||||
|
return err
|
||||||
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
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
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
|
||||||
|
}
|
||||||
@@ -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()
|
||||||
|
|
||||||
|
|||||||
145
pkg/es/node.go
145
pkg/es/node.go
@@ -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()
|
||||||
|
}
|
||||||
|
|||||||
@@ -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())
|
||||||
|
|
||||||
|
|||||||
@@ -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).
|
||||||
|
|||||||
@@ -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 {
|
||||||
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 {
|
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 {
|
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 {
|
||||||
|
|||||||
Reference in New Issue
Block a user