mirror of
https://codeberg.org/scip/esctl.git
synced 2026-08-25 01:24:18 +02:00
Compare commits
3 Commits
55a1005857
...
feature/es
| Author | SHA1 | Date | |
|---|---|---|---|
| 7f1018f0e1 | |||
| 3486bc0370 | |||
| 374ff99916 |
25
README.md
25
README.md
@@ -29,7 +29,6 @@ Features:
|
|||||||
logical condition (OR, AND), use PIT, limit datetime (ES date math
|
logical condition (OR, AND), use PIT, limit datetime (ES date math
|
||||||
can be used), etc. It is however not yet possible to create
|
can be used), etc. It is however not yet possible to create
|
||||||
recursive searches like: `(cond1 AND cond2) OR (cond3 OR cond4)`.
|
recursive searches like: `(cond1 AND cond2) OR (cond3 OR cond4)`.
|
||||||
- Search using ES|QL language: `esctl searchql`.
|
|
||||||
- Cross cluster replication (ccr): view, pause, resume, delete
|
- Cross cluster replication (ccr): view, pause, resume, delete
|
||||||
replication. You can also manage follower configuration.
|
replication. You can also manage follower configuration.
|
||||||
- Index management: manage aliases, create, modify, delete indices,
|
- Index management: manage aliases, create, modify, delete indices,
|
||||||
@@ -353,30 +352,6 @@ $ esctl search -i foo* -F title=zeitbuchung message=pause | jq
|
|||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
You can also search using [ES|QL](https://www.elastic.co/docs/reference/query-languages/esql/esql-getting-started):
|
|
||||||
|
|
||||||
```console
|
|
||||||
$ esctl searchql "from hyperdrive | sort @timestamp | limit 5"
|
|
||||||
@TIMESTAMP MESSAGE TAG
|
|
||||||
2026-06-24T08:24:56.000Z arosu loop
|
|
||||||
2026-06-24T08:24:58.000Z hami loop
|
|
||||||
2026-06-24T08:24:59.000Z ishininu loop
|
|
||||||
2026-06-24T08:25:00.000Z uyomoruron loop
|
|
||||||
2026-06-24T08:25:02.000Z ishimime loop
|
|
||||||
```
|
|
||||||
|
|
||||||
There are several output modes (json, yaml, csv), to get esql output as CSV:
|
|
||||||
|
|
||||||
```console
|
|
||||||
$ esctl searchql "from hyperdrive | sort @timestamp | limit 5" -o csv
|
|
||||||
@timestamp,message,tag
|
|
||||||
2026-06-24T08:24:56.000Z,arosu,loop
|
|
||||||
2026-06-24T08:24:58.000Z,hami,loop
|
|
||||||
2026-06-24T08:24:59.000Z,ishininu,loop
|
|
||||||
2026-06-24T08:25:00.000Z,uyomoruron,loop
|
|
||||||
2026-06-24T08:25:02.000Z,ishimime,loop
|
|
||||||
```
|
|
||||||
|
|
||||||
To check which field mappings are available for an index:
|
To check which field mappings are available for an index:
|
||||||
```console
|
```console
|
||||||
$ esctl index show foo2
|
$ esctl index show foo2
|
||||||
|
|||||||
@@ -34,7 +34,6 @@ const (
|
|||||||
Cilm
|
Cilm
|
||||||
Cnode
|
Cnode
|
||||||
Cclustersettings
|
Cclustersettings
|
||||||
Cindextemplate
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func complete(conf *cfg.Config, cmd *cli.Command, what int) {
|
func complete(conf *cfg.Config, cmd *cli.Command, what int) {
|
||||||
@@ -66,8 +65,6 @@ func complete(conf *cfg.Config, cmd *cli.Command, what int) {
|
|||||||
list, err = es.NodeNames(conf)
|
list, err = es.NodeNames(conf)
|
||||||
case Cclustersettings:
|
case Cclustersettings:
|
||||||
list, err = es.ClusterSettingsNames(conf)
|
list, err = es.ClusterSettingsNames(conf)
|
||||||
case Cindextemplate:
|
|
||||||
list, err = es.IndexTemplateNames(conf)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -49,23 +49,8 @@ func IndexTemplateList(conf *cfg.Config) *cli.Command {
|
|||||||
Aliases: []string{"ls"},
|
Aliases: []string{"ls"},
|
||||||
Usage: "list index templates",
|
Usage: "list index templates",
|
||||||
|
|
||||||
Flags: []cli.Flag{
|
|
||||||
&cli.StringSliceFlag{
|
|
||||||
Name: "filter-index-pattern",
|
|
||||||
Usage: "show only index templates which use patterns, which match the filter",
|
|
||||||
Destination: &conf.Filter,
|
|
||||||
Aliases: []string{"F"},
|
|
||||||
},
|
|
||||||
&cli.BoolFlag{
|
|
||||||
Name: "hidden",
|
|
||||||
Usage: "show hidden index templates as well",
|
|
||||||
Destination: &conf.Hidden,
|
|
||||||
Aliases: []string{"H"},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
|
|
||||||
Action: func(ctx context.Context, cmd *cli.Command) error {
|
Action: func(ctx context.Context, cmd *cli.Command) error {
|
||||||
return es.IndexTemplateList(conf, cmd.Args().Get(0))
|
return es.IndexTemplateList(conf)
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -86,7 +71,7 @@ func IndexTemplateShow(conf *cfg.Config) *cli.Command {
|
|||||||
},
|
},
|
||||||
|
|
||||||
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
|
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
|
||||||
complete(conf, cmd, Cindextemplate)
|
complete(conf, cmd, Cindex)
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -230,7 +215,7 @@ func IndexTemplateDelete(conf *cfg.Config) *cli.Command {
|
|||||||
},
|
},
|
||||||
|
|
||||||
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
|
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
|
||||||
complete(conf, cmd, Cindextemplate)
|
complete(conf, cmd, Cindex)
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
13
cmd/root.go
13
cmd/root.go
@@ -21,7 +21,6 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
golog "log"
|
golog "log"
|
||||||
"os"
|
"os"
|
||||||
"runtime/debug"
|
|
||||||
"runtime/pprof"
|
"runtime/pprof"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
@@ -119,12 +118,6 @@ func Main() int {
|
|||||||
Usage: "enable HTTP debugging",
|
Usage: "enable HTTP debugging",
|
||||||
Destination: &conf.DebugHTTP,
|
Destination: &conf.DebugHTTP,
|
||||||
},
|
},
|
||||||
&cli.BoolFlag{
|
|
||||||
Name: "debug-goroutines",
|
|
||||||
Value: false,
|
|
||||||
Usage: "enable goroutine debugging",
|
|
||||||
Destination: &conf.DebugGoRoutines,
|
|
||||||
},
|
|
||||||
&cli.BoolFlag{
|
&cli.BoolFlag{
|
||||||
Name: "align-ints",
|
Name: "align-ints",
|
||||||
Aliases: []string{"I"},
|
Aliases: []string{"I"},
|
||||||
@@ -260,12 +253,6 @@ func Version(conf *cfg.Config) *cli.Command {
|
|||||||
Usage: "show esctl version information",
|
Usage: "show esctl version information",
|
||||||
|
|
||||||
Action: func(ctx context.Context, cmd *cli.Command) error {
|
Action: func(ctx context.Context, cmd *cli.Command) error {
|
||||||
info, _ := debug.ReadBuildInfo()
|
|
||||||
|
|
||||||
if strings.Contains(info.Main.Version, "+dirty") {
|
|
||||||
cfg.COMMIT += "+dirty"
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err := fmt.Printf(versionFmt,
|
_, err := fmt.Printf(versionFmt,
|
||||||
cfg.Version, cfg.BUILD, cfg.BRANCH, cfg.COMMIT, cfg.GOVERSION, cfg.APIVERSION)
|
cfg.Version, cfg.BUILD, cfg.BRANCH, cfg.COMMIT, cfg.GOVERSION, cfg.APIVERSION)
|
||||||
|
|
||||||
|
|||||||
@@ -194,18 +194,18 @@ func (cluster *Cluster) getDefaultOptions() []elasticsearch.Option {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (cluster *Cluster) getTransport() elastictransport.Option {
|
func (cluster *Cluster) getTransport() elastictransport.Option {
|
||||||
transport := new(http.Transport{
|
transport := &http.Transport{
|
||||||
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
||||||
})
|
}
|
||||||
|
|
||||||
if cluster.DebugHTTP {
|
if cluster.DebugHTTP {
|
||||||
return elastictransport.WithTransport(
|
return elastictransport.WithTransport(
|
||||||
new(DebugTransport{Transport: transport}),
|
&DebugTransport{Transport: transport},
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
return elastictransport.WithTransport(
|
return elastictransport.WithTransport(
|
||||||
new(CompatibilityTransport{Transport: transport}),
|
&CompatibilityTransport{Transport: transport},
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -28,7 +28,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
Version string = `v0.0.27`
|
Version string = `v0.0.26`
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -90,7 +90,6 @@ type Config struct {
|
|||||||
Force bool // ccr follower renew: -f
|
Force bool // ccr follower renew: -f
|
||||||
|
|
||||||
DebugHTTP bool // root: --debug-http
|
DebugHTTP bool // root: --debug-http
|
||||||
DebugGoRoutines bool // root: --debug-goroutines
|
|
||||||
Separator string // role diff: -s
|
Separator string // role diff: -s
|
||||||
NotDeployed bool // role diff: -n
|
NotDeployed bool // role diff: -n
|
||||||
Undefined bool // role diff: -u
|
Undefined bool // role diff: -u
|
||||||
@@ -114,7 +113,7 @@ type Config struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func NewConfig() *Config {
|
func NewConfig() *Config {
|
||||||
return new(Config{Clusters: map[string]*Cluster{}})
|
return &Config{Clusters: map[string]*Cluster{}}
|
||||||
}
|
}
|
||||||
|
|
||||||
func getDefaultPath() string {
|
func getDefaultPath() string {
|
||||||
@@ -217,7 +216,7 @@ func (conf *Config) LoadConfig() error {
|
|||||||
return fmt.Errorf("failed to read config file: %w", err)
|
return fmt.Errorf("failed to read config file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
newconf := new(Config{})
|
newconf := &Config{}
|
||||||
|
|
||||||
err = yaml.Unmarshal(data, newconf)
|
err = yaml.Unmarshal(data, newconf)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -28,7 +28,6 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"log"
|
"log"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"maps"
|
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
@@ -159,7 +158,7 @@ func ApiRepl(conf *cfg.Config) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func pageJsonOutput(conf *cfg.Config, raw []byte) {
|
func pageJsonOutput(conf *cfg.Config, raw []byte) {
|
||||||
tmpconf := new(cfg.Config{HaveJQ: conf.HaveJQ})
|
tmpconf := &cfg.Config{HaveJQ: conf.HaveJQ}
|
||||||
|
|
||||||
if conf.Pager != "" {
|
if conf.Pager != "" {
|
||||||
tmpconf.HaveJQ = false
|
tmpconf.HaveJQ = false
|
||||||
@@ -207,16 +206,16 @@ func CallAPI(conf *cfg.Config, verb, path, data string) ([]byte, error) {
|
|||||||
verb = strings.ToUpper(verb)
|
verb = strings.ToUpper(verb)
|
||||||
|
|
||||||
// we're using port-forwards anyway
|
// we're using port-forwards anyway
|
||||||
noVerifyTransport := new(http.Transport{
|
noVerifyTransport := &http.Transport{
|
||||||
TLSClientConfig: new(tls.Config{InsecureSkipVerify: true}),
|
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
||||||
})
|
}
|
||||||
|
|
||||||
client := new(http.Client{Transport: noVerifyTransport})
|
client := &http.Client{Transport: noVerifyTransport}
|
||||||
|
|
||||||
if conf.DebugHTTP {
|
if conf.DebugHTTP {
|
||||||
client = new(http.Client{
|
client = &http.Client{
|
||||||
Transport: new(cfg.DebugTransport{
|
Transport: &cfg.DebugTransport{
|
||||||
Transport: noVerifyTransport})})
|
Transport: noVerifyTransport}}
|
||||||
}
|
}
|
||||||
|
|
||||||
req, err := http.NewRequest(verb, conf.DefaultCluster.Uri+path, bytes.NewBuffer([]byte(data)))
|
req, err := http.NewRequest(verb, conf.DefaultCluster.Uri+path, bytes.NewBuffer([]byte(data)))
|
||||||
@@ -367,7 +366,16 @@ func ApiList(conf *cfg.Config, pattern string) error {
|
|||||||
func ApiPathNames() []string {
|
func ApiPathNames() []string {
|
||||||
assets.LoadAssetOpenApi()
|
assets.LoadAssetOpenApi()
|
||||||
|
|
||||||
return slices.Collect(maps.Keys(assets.OpenAPI.Spec().Paths.Paths))
|
paths := make([]string, len(assets.OpenAPI.Spec().Paths.Paths))
|
||||||
|
|
||||||
|
idx := 0
|
||||||
|
|
||||||
|
for path := range assets.OpenAPI.Spec().Paths.Paths {
|
||||||
|
paths[idx] = path
|
||||||
|
idx++
|
||||||
|
}
|
||||||
|
|
||||||
|
return paths
|
||||||
}
|
}
|
||||||
|
|
||||||
func ApiShow(conf *cfg.Config, showpath, verb string) error {
|
func ApiShow(conf *cfg.Config, showpath, verb string) error {
|
||||||
@@ -528,7 +536,7 @@ func getApiExample(op *Op) string {
|
|||||||
// otherwise showpath+verb have to match precisely.
|
// otherwise showpath+verb have to match precisely.
|
||||||
func matchOperation(showpath, verb string) (*Op, error) {
|
func matchOperation(showpath, verb string) (*Op, error) {
|
||||||
ops := []*Op{}
|
ops := []*Op{}
|
||||||
op := new(Op{})
|
op := &Op{}
|
||||||
|
|
||||||
var found bool
|
var found bool
|
||||||
|
|
||||||
|
|||||||
@@ -17,10 +17,6 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
|||||||
package es
|
package es
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
|
||||||
// "encoding/json/jsontext"
|
|
||||||
// "encoding/json/v2"
|
|
||||||
|
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
@@ -53,22 +49,11 @@ func getHealthReport(conf *cfg.Config) (*HealthReport, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
report := new(HealthReport{})
|
report := HealthReport{}
|
||||||
|
|
||||||
// FIXME: use this once jsonv2 is no more experimental it already
|
|
||||||
// builds and works like intended, but golangci-lint doesn't
|
|
||||||
// recognize it with: go: unknown GOEXPERIMENT jsonv2
|
|
||||||
//
|
|
||||||
// if err := json.UnmarshalDecode(
|
|
||||||
// jsontext.NewDecoder(
|
|
||||||
// bytes.NewBuffer(raw)),
|
|
||||||
// &report); err != nil {
|
|
||||||
// return nil, fmt.Errorf("failed to unmarshal healthreport response: %w", err)
|
|
||||||
// }
|
|
||||||
|
|
||||||
if err := json.Unmarshal(raw, &report); err != nil {
|
if err := json.Unmarshal(raw, &report); err != nil {
|
||||||
return nil, fmt.Errorf("failed to unmarshal healthreport response: %w", err)
|
return nil, fmt.Errorf("failed to unmarshal healthreport response: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return report, nil
|
return &report, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -101,39 +101,23 @@ func getClusterStatus(conf *cfg.Config) (*apiResponse, error) {
|
|||||||
es := conf.DefaultCluster.ES()
|
es := conf.DefaultCluster.ES()
|
||||||
|
|
||||||
responses := make(chan apiResponse, gocount)
|
responses := make(chan apiResponse, gocount)
|
||||||
wg := new(sync.WaitGroup{})
|
wg := &sync.WaitGroup{}
|
||||||
|
|
||||||
wg.Go(func() {
|
wg.Add(gocount)
|
||||||
getApiData(conf, es, responses, "health")
|
go getApiData(conf, es, wg, responses, "health")
|
||||||
})
|
go getApiData(conf, es, wg, responses, "healthreport")
|
||||||
|
go getApiData(conf, es, wg, responses, "info")
|
||||||
wg.Go(func() {
|
go getApiData(conf, es, wg, responses, "ccr")
|
||||||
getApiData(conf, es, responses, "healthreport")
|
go getApiData(conf, es, wg, responses, "indices")
|
||||||
})
|
go getApiData(conf, es, wg, responses, "tasks")
|
||||||
|
|
||||||
wg.Go(func() {
|
|
||||||
getApiData(conf, es, responses, "info")
|
|
||||||
})
|
|
||||||
|
|
||||||
wg.Go(func() {
|
|
||||||
getApiData(conf, es, responses, "ccr")
|
|
||||||
})
|
|
||||||
|
|
||||||
wg.Go(func() {
|
|
||||||
getApiData(conf, es, responses, "indices")
|
|
||||||
})
|
|
||||||
|
|
||||||
wg.Go(func() {
|
|
||||||
getApiData(conf, es, responses, "tasks")
|
|
||||||
})
|
|
||||||
|
|
||||||
if conf.Verbose {
|
if conf.Verbose {
|
||||||
getApiData(conf, es, responses, "stats")
|
go getApiData(conf, es, wg, responses, "stats")
|
||||||
}
|
}
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|
||||||
all := new(apiResponse{})
|
all := apiResponse{}
|
||||||
|
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
@@ -160,7 +144,7 @@ func getClusterStatus(conf *cfg.Config) (*apiResponse, error) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return all, err
|
return &all, err
|
||||||
}
|
}
|
||||||
|
|
||||||
func ClusterStatus(conf *cfg.Config) error {
|
func ClusterStatus(conf *cfg.Config) error {
|
||||||
|
|||||||
@@ -29,12 +29,12 @@ func ClusterRerouteMove(conf *cfg.Config, index string) error {
|
|||||||
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
||||||
|
|
||||||
commands := esdsl.NewCommand()
|
commands := esdsl.NewCommand()
|
||||||
moveCommand := new(types.CommandMoveAction{
|
moveCommand := &types.CommandMoveAction{
|
||||||
Shard: conf.Shards,
|
Shard: conf.Shards,
|
||||||
FromNode: conf.FromNode,
|
FromNode: conf.FromNode,
|
||||||
ToNode: conf.ToNode,
|
ToNode: conf.ToNode,
|
||||||
Index: index,
|
Index: index,
|
||||||
})
|
}
|
||||||
|
|
||||||
commands.CommandCaster().Move = moveCommand
|
commands.CommandCaster().Move = moveCommand
|
||||||
|
|
||||||
@@ -52,11 +52,11 @@ func ClusterRerouteAllocateReplica(conf *cfg.Config, index string) error {
|
|||||||
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
||||||
|
|
||||||
commands := esdsl.NewCommand()
|
commands := esdsl.NewCommand()
|
||||||
allocCommand := new(types.CommandAllocateReplicaAction{
|
allocCommand := &types.CommandAllocateReplicaAction{
|
||||||
Shard: conf.Shards,
|
Shard: conf.Shards,
|
||||||
Node: conf.ToNode,
|
Node: conf.ToNode,
|
||||||
Index: index,
|
Index: index,
|
||||||
})
|
}
|
||||||
|
|
||||||
commands.CommandCaster().AllocateReplica = allocCommand
|
commands.CommandCaster().AllocateReplica = allocCommand
|
||||||
|
|
||||||
@@ -74,12 +74,12 @@ func ClusterRerouteCancel(conf *cfg.Config, index string) error {
|
|||||||
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
||||||
|
|
||||||
commands := esdsl.NewCommand()
|
commands := esdsl.NewCommand()
|
||||||
cancelCommand := new(types.CommandCancelAction{
|
cancelCommand := &types.CommandCancelAction{
|
||||||
Shard: conf.Shards,
|
Shard: conf.Shards,
|
||||||
Node: conf.ToNode,
|
Node: conf.ToNode,
|
||||||
Index: index,
|
Index: index,
|
||||||
AllowPrimary: &conf.AllowPrimary,
|
AllowPrimary: &conf.AllowPrimary,
|
||||||
})
|
}
|
||||||
|
|
||||||
commands.CommandCaster().Cancel = cancelCommand
|
commands.CommandCaster().Cancel = cancelCommand
|
||||||
|
|
||||||
@@ -97,12 +97,12 @@ func ClusterRerouteAllocatePrimary(conf *cfg.Config, index string, stale bool) e
|
|||||||
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
move := conf.DefaultCluster.ES().Cluster.Reroute()
|
||||||
|
|
||||||
commands := esdsl.NewCommand()
|
commands := esdsl.NewCommand()
|
||||||
allocCommand := new(types.CommandAllocatePrimaryAction{
|
allocCommand := &types.CommandAllocatePrimaryAction{
|
||||||
Shard: conf.Shards,
|
Shard: conf.Shards,
|
||||||
Node: conf.ToNode,
|
Node: conf.ToNode,
|
||||||
Index: index,
|
Index: index,
|
||||||
AcceptDataLoss: conf.AcceptDataLoss,
|
AcceptDataLoss: conf.AcceptDataLoss,
|
||||||
})
|
}
|
||||||
|
|
||||||
if stale {
|
if stale {
|
||||||
commands.CommandCaster().AllocateStalePrimary = allocCommand
|
commands.CommandCaster().AllocateStalePrimary = allocCommand
|
||||||
|
|||||||
@@ -101,7 +101,7 @@ func DocDelete(conf *cfg.Config, queries []string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
req := new(deletebyquery.Request{})
|
req := &deletebyquery.Request{}
|
||||||
|
|
||||||
if len(queries) == 0 && conf.All {
|
if len(queries) == 0 && conf.All {
|
||||||
req.Query = esdsl.NewMatchAllQuery().QueryCaster()
|
req.Query = esdsl.NewMatchAllQuery().QueryCaster()
|
||||||
|
|||||||
@@ -22,8 +22,6 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"maps"
|
|
||||||
"slices"
|
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"codeberg.org/scip/esctl/pkg/cfg"
|
||||||
@@ -64,7 +62,13 @@ func IlmNames(conf *cfg.Config) ([]string, error) {
|
|||||||
return nil, fmt.Errorf("failed to get ilm policies: %w", esErrorString(err))
|
return nil, fmt.Errorf("failed to get ilm policies: %w", esErrorString(err))
|
||||||
}
|
}
|
||||||
|
|
||||||
names := slices.Collect(maps.Keys(res))
|
names := make([]string, len(res))
|
||||||
|
idx := 0
|
||||||
|
|
||||||
|
for name := range res {
|
||||||
|
names[idx] = name
|
||||||
|
idx++
|
||||||
|
}
|
||||||
|
|
||||||
return names, nil
|
return names, nil
|
||||||
}
|
}
|
||||||
@@ -347,7 +351,7 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
|||||||
|
|
||||||
var actions types.IlmActionsVariant = esdsl.NewIlmActions()
|
var actions types.IlmActionsVariant = esdsl.NewIlmActions()
|
||||||
|
|
||||||
rollover := new(types.RolloverAction{})
|
rollover := &types.RolloverAction{}
|
||||||
haveroll := false
|
haveroll := false
|
||||||
|
|
||||||
if policy != nil {
|
if policy != nil {
|
||||||
@@ -509,8 +513,8 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
|
|||||||
phases.PhasesCaster().Delete = policy.Phases.Delete
|
phases.PhasesCaster().Delete = policy.Phases.Delete
|
||||||
}
|
}
|
||||||
|
|
||||||
put := new(putlifecycle.Request{})
|
put := &putlifecycle.Request{}
|
||||||
newpolicy := new(types.IlmPolicy{})
|
newpolicy := &types.IlmPolicy{}
|
||||||
newpolicy.IlmPolicyCaster().Phases = *phases.PhasesCaster()
|
newpolicy.IlmPolicyCaster().Phases = *phases.PhasesCaster()
|
||||||
put.Policy = newpolicy
|
put.Policy = newpolicy
|
||||||
|
|
||||||
|
|||||||
@@ -185,19 +185,12 @@ func virtualAge(phase *PhaseData) time.Duration {
|
|||||||
// Retrieve all index, ilm-explain and ilm-policies in parallel
|
// Retrieve all index, ilm-explain and ilm-policies in parallel
|
||||||
func getIlmPhaseData(conf *cfg.Config) ([]PhaseData, error) {
|
func getIlmPhaseData(conf *cfg.Config) ([]PhaseData, error) {
|
||||||
responses := make(chan apiResponse, 3)
|
responses := make(chan apiResponse, 3)
|
||||||
wg := new(sync.WaitGroup{})
|
wg := &sync.WaitGroup{}
|
||||||
|
wg.Add(3)
|
||||||
|
|
||||||
wg.Go(func() {
|
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "indicesbytes")
|
||||||
getApiData(conf, conf.DefaultCluster.ES(), responses, "indicesbytes")
|
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "explain")
|
||||||
})
|
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "policies")
|
||||||
|
|
||||||
wg.Go(func() {
|
|
||||||
getApiData(conf, conf.DefaultCluster.ES(), responses, "explain")
|
|
||||||
})
|
|
||||||
|
|
||||||
wg.Go(func() {
|
|
||||||
getApiData(conf, conf.DefaultCluster.ES(), responses, "policies")
|
|
||||||
})
|
|
||||||
|
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|
||||||
@@ -338,7 +331,7 @@ func findNextPhase(policy types.IlmPolicy, currentPhase string) *NextPhase {
|
|||||||
// phase list to determine which comes next
|
// phase list to determine which comes next
|
||||||
phases, start := registerPhases(policy, currentPhase)
|
phases, start := registerPhases(policy, currentPhase)
|
||||||
|
|
||||||
nextPhase := new(NextPhase{})
|
nextPhase := &NextPhase{}
|
||||||
|
|
||||||
// finally determine which phase comes next
|
// finally determine which phase comes next
|
||||||
// exception: hot, where we look for rollover rules
|
// exception: hot, where we look for rollover rules
|
||||||
|
|||||||
@@ -33,23 +33,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// used for completion
|
// used for completion
|
||||||
func IndexTemplateNames(conf *cfg.Config) ([]string, error) {
|
func IndexTemplateList(conf *cfg.Config) error {
|
||||||
res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
|
|
||||||
Do(context.Background())
|
|
||||||
if err != nil {
|
|
||||||
return nil, fmt.Errorf("failed to get index templates: %w", esErrorString(err))
|
|
||||||
}
|
|
||||||
|
|
||||||
names := make([]string, len(res.IndexTemplates))
|
|
||||||
|
|
||||||
for idx, tpl := range res.IndexTemplates {
|
|
||||||
names[idx] = tpl.Name
|
|
||||||
}
|
|
||||||
|
|
||||||
return names, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func IndexTemplateList(conf *cfg.Config, filter string) error {
|
|
||||||
res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
|
res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
|
||||||
Do(context.Background())
|
Do(context.Background())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -58,42 +42,16 @@ func IndexTemplateList(conf *cfg.Config, filter string) error {
|
|||||||
|
|
||||||
slog.Debug("res", "index templates", res)
|
slog.Debug("res", "index templates", res)
|
||||||
|
|
||||||
table := printer.NewTable(conf, 5, 0)
|
table := printer.NewTable(conf, 5, len(res.IndexTemplates))
|
||||||
table.Addheaders("name", "description", "index patterns", "priority")
|
table.Addheaders("name", "description", "priority")
|
||||||
|
|
||||||
for _, tpl := range res.IndexTemplates {
|
for idx, tpl := range res.IndexTemplates {
|
||||||
if !conf.Hidden && strings.HasPrefix(tpl.Name, ".") {
|
desc, err := json.Marshal(tpl.IndexTemplate.Meta_["description"])
|
||||||
continue
|
if err != nil {
|
||||||
|
return fmt.Errorf("failed to unmarshal meta json data: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if filter != "" && !strings.Contains(tpl.Name, filter) {
|
table.Entries[idx] = []any{tpl.Name, desc, tpl.IndexTemplate.Priority}
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(conf.Filter) > 0 {
|
|
||||||
skip := true
|
|
||||||
|
|
||||||
for _, filter := range conf.Filter {
|
|
||||||
for _, pattern := range tpl.IndexTemplate.IndexPatterns {
|
|
||||||
if strings.Contains(pattern, filter) {
|
|
||||||
skip = false
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if skip {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
desc := strings.TrimPrefix(strings.TrimSuffix(string(tpl.IndexTemplate.Meta_["description"]), `"`), `"`)
|
|
||||||
|
|
||||||
var prio int64
|
|
||||||
if tpl.IndexTemplate.Priority != nil {
|
|
||||||
prio = *tpl.IndexTemplate.Priority
|
|
||||||
}
|
|
||||||
|
|
||||||
table.AddRow(tpl.Name, desc, tpl.IndexTemplate.IndexPatterns, prio)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
table.Sort()
|
table.Sort()
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ package es
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"sync"
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"codeberg.org/scip/esctl/pkg/cfg"
|
||||||
"github.com/elastic/go-elasticsearch/v9"
|
"github.com/elastic/go-elasticsearch/v9"
|
||||||
@@ -60,7 +61,14 @@ type apiResponse struct {
|
|||||||
which int
|
which int
|
||||||
}
|
}
|
||||||
|
|
||||||
func getApiData(conf *cfg.Config, es *elasticsearch.TypedClient, reschan chan apiResponse, which string) {
|
func getApiData(
|
||||||
|
conf *cfg.Config,
|
||||||
|
es *elasticsearch.TypedClient,
|
||||||
|
wg *sync.WaitGroup,
|
||||||
|
reschan chan apiResponse,
|
||||||
|
which string) {
|
||||||
|
defer wg.Done()
|
||||||
|
|
||||||
apiRes := apiResponse{}
|
apiRes := apiResponse{}
|
||||||
|
|
||||||
var arerr error
|
var arerr error
|
||||||
|
|||||||
@@ -20,8 +20,6 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"maps"
|
|
||||||
"slices"
|
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"codeberg.org/scip/esctl/pkg/cfg"
|
||||||
"codeberg.org/scip/esctl/pkg/printer"
|
"codeberg.org/scip/esctl/pkg/printer"
|
||||||
@@ -35,7 +33,15 @@ func RoleNames(conf *cfg.Config) ([]string, error) {
|
|||||||
return nil, fmt.Errorf("failed to get roles: %w", esErrorString(err))
|
return nil, fmt.Errorf("failed to get roles: %w", esErrorString(err))
|
||||||
}
|
}
|
||||||
|
|
||||||
return slices.Collect(maps.Keys(res)), nil
|
roles := make([]string, len(res))
|
||||||
|
|
||||||
|
idx := 0
|
||||||
|
for name := range res {
|
||||||
|
roles[idx] = name
|
||||||
|
idx++
|
||||||
|
}
|
||||||
|
|
||||||
|
return roles, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func RoleList(conf *cfg.Config) error {
|
func RoleList(conf *cfg.Config) error {
|
||||||
|
|||||||
@@ -127,7 +127,7 @@ func getCsvRecord(conf *cfg.Config, csvfile, rolename string) (*Record, error) {
|
|||||||
}()
|
}()
|
||||||
|
|
||||||
scanner := bufio.NewScanner(fd)
|
scanner := bufio.NewScanner(fd)
|
||||||
record := new(Record{role: rolename})
|
record := Record{role: rolename}
|
||||||
|
|
||||||
for scanner.Scan() {
|
for scanner.Scan() {
|
||||||
line := strings.TrimSpace(scanner.Text())
|
line := strings.TrimSpace(scanner.Text())
|
||||||
@@ -150,7 +150,7 @@ func getCsvRecord(conf *cfg.Config, csvfile, rolename string) (*Record, error) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return record, nil
|
return &record, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func diffRoles(conf *cfg.Config, records map[string]Record, res getrole.Response) []Register {
|
func diffRoles(conf *cfg.Config, records map[string]Record, res getrole.Response) []Register {
|
||||||
|
|||||||
@@ -56,7 +56,7 @@ func Search(conf *cfg.Config, queries []string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
req := new(search.Request{Query: queryCaster})
|
req := &search.Request{Query: queryCaster}
|
||||||
|
|
||||||
searchEs.Request(req)
|
searchEs.Request(req)
|
||||||
|
|
||||||
@@ -128,7 +128,7 @@ func validateSearch(conf *cfg.Config, queries []string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
req := new(validatequery.Request{Query: queryCaster})
|
req := &validatequery.Request{Query: queryCaster}
|
||||||
|
|
||||||
validate.Request(req)
|
validate.Request(req)
|
||||||
|
|
||||||
|
|||||||
@@ -80,7 +80,7 @@ func NewFilter(query string) (*filter, error) {
|
|||||||
return nil, errors.New("search queries must be in the form field<sep>pattern where <sep> must be one of: = or !=")
|
return nil, errors.New("search queries must be in the form field<sep>pattern where <sep> must be one of: = or !=")
|
||||||
}
|
}
|
||||||
|
|
||||||
flt := new(filter{term: part[0], filter: part[1], criteria: criteria})
|
flt := &filter{term: part[0], filter: part[1], criteria: criteria}
|
||||||
|
|
||||||
if strings.Contains(part[0], ",") {
|
if strings.Contains(part[0], ",") {
|
||||||
// a MultiMatchQuery, match across multiple fields at once
|
// a MultiMatchQuery, match across multiple fields at once
|
||||||
|
|||||||
@@ -66,13 +66,13 @@ func SnapshotList(conf *cfg.Config) error {
|
|||||||
snapshots := []*Snapshot{} // original snapshot names
|
snapshots := []*Snapshot{} // original snapshot names
|
||||||
|
|
||||||
for _, snapshot := range sres {
|
for _, snapshot := range sres {
|
||||||
snap := new(Snapshot{
|
snap := &Snapshot{
|
||||||
Name: *snapshot.Id,
|
Name: *snapshot.Id,
|
||||||
Status: *snapshot.Status,
|
Status: *snapshot.Status,
|
||||||
Start: fmt.Sprintf("%s", snapshot.StartTime),
|
Start: fmt.Sprintf("%s", snapshot.StartTime),
|
||||||
Forindex: indexFromSnapshot(*snapshot.Id),
|
Forindex: indexFromSnapshot(*snapshot.Id),
|
||||||
Orphaned: "no",
|
Orphaned: "no",
|
||||||
})
|
}
|
||||||
|
|
||||||
_, exists := indicies[snap.Forindex]
|
_, exists := indicies[snap.Forindex]
|
||||||
if !exists {
|
if !exists {
|
||||||
|
|||||||
@@ -29,13 +29,13 @@ import (
|
|||||||
const LevelNotice = slog.Level(2)
|
const LevelNotice = slog.Level(2)
|
||||||
|
|
||||||
func Init(conf *cfg.Config) {
|
func Init(conf *cfg.Config) {
|
||||||
logLevel := new(slog.LevelVar{})
|
logLevel := &slog.LevelVar{}
|
||||||
|
|
||||||
opts := new(yadu.Options{
|
opts := &yadu.Options{
|
||||||
Level: logLevel,
|
Level: logLevel,
|
||||||
AddSource: true,
|
AddSource: true,
|
||||||
NoColor: !isatty.IsTerminal(os.Stdout.Fd()),
|
NoColor: !isatty.IsTerminal(os.Stdout.Fd()),
|
||||||
})
|
}
|
||||||
|
|
||||||
buildInfo, _ := debug.ReadBuildInfo()
|
buildInfo, _ := debug.ReadBuildInfo()
|
||||||
|
|
||||||
|
|||||||
@@ -27,7 +27,7 @@ func (b *ByteSize) String() string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func Bytes(size int64) *ByteSize {
|
func Bytes(size int64) *ByteSize {
|
||||||
return new(ByteSize{size: uint64(size)})
|
return &ByteSize{size: uint64(size)}
|
||||||
}
|
}
|
||||||
|
|
||||||
func ByteString(size int64) string {
|
func ByteString(size int64) string {
|
||||||
|
|||||||
@@ -1,55 +0,0 @@
|
|||||||
/*
|
|
||||||
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 printer
|
|
||||||
|
|
||||||
import (
|
|
||||||
"fmt"
|
|
||||||
"runtime/metrics"
|
|
||||||
"strconv"
|
|
||||||
)
|
|
||||||
|
|
||||||
func printGoRoutineMetrics() {
|
|
||||||
fmt.Println("\nGoroutine metrics:")
|
|
||||||
printMetric("/sched/goroutines-created:goroutines", "Created")
|
|
||||||
printMetric("/sched/goroutines:goroutines", "Live")
|
|
||||||
printMetric("/sched/goroutines/not-in-go:goroutines", "Syscall/CGO")
|
|
||||||
printMetric("/sched/goroutines/runnable:goroutines", "Runnable")
|
|
||||||
printMetric("/sched/goroutines/running:goroutines", "Running")
|
|
||||||
printMetric("/sched/goroutines/waiting:goroutines", "Waiting")
|
|
||||||
|
|
||||||
fmt.Println("Thread metrics:")
|
|
||||||
printMetric("/sched/gomaxprocs:threads", "Max")
|
|
||||||
printMetric("/sched/threads/total:threads", "Live")
|
|
||||||
}
|
|
||||||
|
|
||||||
func printMetric(name string, descr string) {
|
|
||||||
sample := []metrics.Sample{{Name: name}}
|
|
||||||
metrics.Read(sample)
|
|
||||||
|
|
||||||
var val string
|
|
||||||
|
|
||||||
switch sample[0].Value.Kind() {
|
|
||||||
case metrics.KindFloat64, metrics.KindFloat64Histogram:
|
|
||||||
val = fmt.Sprintf("%.2f", sample[0].Value.Float64())
|
|
||||||
case metrics.KindUint64:
|
|
||||||
val = strconv.FormatUint(sample[0].Value.Uint64(), 10)
|
|
||||||
case metrics.KindBad:
|
|
||||||
val = "n/a"
|
|
||||||
}
|
|
||||||
|
|
||||||
fmt.Printf(" %s: %v\n", descr, val)
|
|
||||||
}
|
|
||||||
@@ -40,15 +40,10 @@ type Table struct {
|
|||||||
lenHeaders []int
|
lenHeaders []int
|
||||||
alignInts bool
|
alignInts bool
|
||||||
maxwidth int
|
maxwidth int
|
||||||
debugGoRoutines bool
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewTable(conf *cfg.Config, columns, rows int) *Table {
|
func NewTable(conf *cfg.Config, columns, rows int) *Table {
|
||||||
table := new(Table{
|
table := Table{Mode: conf.Output, maxwidth: cfg.GetTermWidth()}
|
||||||
Mode: conf.Output,
|
|
||||||
maxwidth: cfg.GetTermWidth(),
|
|
||||||
debugGoRoutines: conf.DebugGoRoutines,
|
|
||||||
})
|
|
||||||
|
|
||||||
table.Headers = make([]string, columns)
|
table.Headers = make([]string, columns)
|
||||||
table.RawHeaders = make([]string, columns)
|
table.RawHeaders = make([]string, columns)
|
||||||
@@ -56,18 +51,14 @@ func NewTable(conf *cfg.Config, columns, rows int) *Table {
|
|||||||
table.lenHeaders = make([]int, columns)
|
table.lenHeaders = make([]int, columns)
|
||||||
table.alignInts = conf.AlignInts
|
table.alignInts = conf.AlignInts
|
||||||
|
|
||||||
return table
|
return &table
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewTableEmpty(conf *cfg.Config) *Table {
|
func NewTableEmpty(conf *cfg.Config) *Table {
|
||||||
table := new(Table{
|
table := Table{Mode: conf.Output, maxwidth: cfg.GetTermWidth()}
|
||||||
Mode: conf.Output,
|
|
||||||
maxwidth: cfg.GetTermWidth(),
|
|
||||||
debugGoRoutines: conf.DebugGoRoutines,
|
|
||||||
})
|
|
||||||
table.alignInts = conf.AlignInts
|
table.alignInts = conf.AlignInts
|
||||||
|
|
||||||
return table
|
return &table
|
||||||
}
|
}
|
||||||
|
|
||||||
func (table *Table) WithHeaders(headers ...string) *Table {
|
func (table *Table) WithHeaders(headers ...string) *Table {
|
||||||
@@ -84,24 +75,16 @@ func (table *Table) WithHeaders(headers ...string) *Table {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (table *Table) Print() error {
|
func (table *Table) Print() error {
|
||||||
var err error
|
|
||||||
|
|
||||||
switch table.Mode {
|
switch table.Mode {
|
||||||
case "json":
|
case "json":
|
||||||
err = table.PrintJSON()
|
return table.PrintJSON()
|
||||||
case "yaml":
|
case "yaml":
|
||||||
err = table.PrintYAML()
|
return table.PrintYAML()
|
||||||
case "csv":
|
case "csv":
|
||||||
err = table.PrintCSV()
|
return table.PrintCSV()
|
||||||
default:
|
default:
|
||||||
err = table.PrintTSV()
|
return table.PrintTSV()
|
||||||
}
|
}
|
||||||
|
|
||||||
if table.debugGoRoutines {
|
|
||||||
printGoRoutineMetrics()
|
|
||||||
}
|
|
||||||
|
|
||||||
return err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -166,11 +149,9 @@ func (table *Table) PrintTSV() error {
|
|||||||
wrapped := wrapper(entry)
|
wrapped := wrapper(entry)
|
||||||
|
|
||||||
// and indent it
|
// and indent it
|
||||||
first := true
|
for idx, line := range strings.Split(wrapped, "\n") {
|
||||||
for line := range strings.Lines(wrapped) {
|
if idx == 0 {
|
||||||
if first {
|
|
||||||
entry = line
|
entry = line
|
||||||
first = false
|
|
||||||
} else {
|
} else {
|
||||||
entry += "\n " + strings.Repeat(" ", currentWidth) + line
|
entry += "\n " + strings.Repeat(" ", currentWidth) + line
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user