diff --git a/README.md b/README.md index c284696..ee5e747 100644 --- a/README.md +++ b/README.md @@ -80,6 +80,7 @@ ccr - manage cross cluster replication info - show ccr remote info cluster - manage cluster[s] status - show cluster status + switch - set current elasticsearch cluster list - list configured clusters settings - cluster settings management list - show cluster settings @@ -158,7 +159,7 @@ Or create a config file such as this: ```yaml clusters: - default: + foobar: uri: https://es.foo.bar:9200/ user: elastic pass: 123456 @@ -172,9 +173,37 @@ and specify it with `-c configfile`. You may also put clusters into a default config file in `~/.config/esctl/config.yaml`. In this case you can omit `-c ...`. -If you want to work on a specific cluster, specify its name with the +If you want to work on a specific cluster, you need to make it the +current default one. You can either manually configure it in the +config: + +```yaml +clusters: + foobar: + uri: https://es.foo.bar:9200/ + user: elastic + pass: 123456 + default: true +``` + +or - if the cluster already exists in the config - issue this command: + +```console +esctl cluster switch foobar +``` + +You may also temporary set a cluster as the current default with the global `-C` option. +The following rules apply for default cluster selection: + +1. If there's just one cluster configured (either via environment + variables or config), this one will be selected. +2. If there is one cluster configured with the `default` flag `true`, + use this one. + +Use `esctl cluster ls` to check status and see which one would be used. + ## Installation The tool does not have any dependencies. Just download the binary for diff --git a/cmd/cluster.go b/cmd/cluster.go index d6b4252..5670cda 100644 --- a/cmd/cluster.go +++ b/cmd/cluster.go @@ -18,6 +18,7 @@ package cmd import ( "context" + "errors" "codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/es" @@ -33,6 +34,7 @@ func Cluster(conf *cfg.Config) *cli.Command { Commands: []*cli.Command{ ClusterStatus(conf), + ClusterSwitch(conf), ClusterList(conf), ClusterSettings(conf), }, @@ -77,3 +79,25 @@ func ClusterStatus(conf *cfg.Config) *cli.Command { }, } } + +func ClusterSwitch(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "switch", + Usage: "set current elasticsearch cluster", + UsageText: "switch ", + Aliases: []string{"ctx"}, + + ShellComplete: func(ctx context.Context, cmd *cli.Command) { + complete(cmd, Ccluster) + }, + + Action: func(ctx context.Context, cmd *cli.Command) error { + name := cmd.Args().Get(0) + if name == "" { + return errors.New("no cluster name specified") + } + + return conf.SwitchCluster(name) + }, + } +} diff --git a/cmd/root.go b/cmd/root.go index d7f7b0b..688fec0 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -122,7 +122,10 @@ func Main() int { Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) { if tree { - Tree(cmd) + if err := Tree(cmd); err != nil { + return nil, err + } + os.Exit(0) } @@ -254,7 +257,7 @@ func Debug(conf *cfg.Config) *cli.Command { func Tree(cmd *cli.Command) error { max := 20 - cmd.Walk(func(cmd *cli.Command) error { + return cmd.Walk(func(cmd *cli.Command) error { path := cmd.Path() command := path[len(path)-1] @@ -269,6 +272,4 @@ func Tree(cmd *cli.Command) error { return nil }) - - return nil } diff --git a/pkg/cfg/cluster.go b/pkg/cfg/cluster.go new file mode 100644 index 0000000..ca0ae3a --- /dev/null +++ b/pkg/cfg/cluster.go @@ -0,0 +1,126 @@ +/* +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 . +*/ +package cfg + +import ( + "errors" + "fmt" + "net/http" + "os" + + "github.com/elastic/elastic-transport-go/v8/elastictransport" + "github.com/elastic/go-elasticsearch/v9" + "gopkg.in/yaml.v3" +) + +// used in general config struct +type Cluster struct { + Uri, User, Pass string + client *elasticsearch.TypedClient + Default bool +} + +// used just for writing back to the config file +type ClusterConfig struct { + Uri, User, Pass string + Default bool +} + +// to write the config, we avoid all other config settings +type WriteConfig struct { + Clusters map[string]*ClusterConfig +} + +func (cluster *Cluster) ES() *elasticsearch.TypedClient { + if cluster.client == nil { + fmt.Println("no current cluster, use 'esctl cluster switch ' to set one") + os.Exit(1) + } + + return cluster.client +} + +func (cluster *Cluster) SetClient(client *elasticsearch.TypedClient) { + cluster.client = client +} + +// set Default=true for the given cluster in the config (if exists) +func (conf *Config) SwitchCluster(name string) error { + _, exists := conf.Clusters[name] + if !exists { + return errors.New("no cluster with that name configured") + } + + cfg := WriteConfig{Clusters: map[string]*ClusterConfig{}} + + for clustername, cluster := range conf.Clusters { + cfg.Clusters[clustername] = &ClusterConfig{ + Uri: cluster.Uri, + User: cluster.User, + Pass: cluster.Pass, + Default: false, + } + + if clustername == name { + cfg.Clusters[clustername].Default = true + } + } + + raw, err := yaml.Marshal(cfg) + if err != nil { + return fmt.Errorf("failed to marshal cluster config: %w", err) + } + + outfile := getDefaultPath() + if conf.ConfigFile != "" { + outfile = conf.ConfigFile + } + + if err := os.WriteFile(outfile, raw, 0600); err != nil { + return err + } + + return nil +} + +func (conf *Config) SetupES() error { + // These headers are not needed with ES 9, but with ES 8, we set + // them here so every API call uses it. The only exception being + // the api repl, which does it on its own. + headers := http.Header{} + headers.Add("content-type", "application/json") + headers.Add("Accept", "application/json") + + for _, cluster := range conf.Clusters { + es, err := elasticsearch.NewTyped( + elasticsearch.WithAddresses(cluster.Uri), + elasticsearch.WithBasicAuth(cluster.User, cluster.Pass), + elasticsearch.WithTransportOptions( + conf.getTransport(), + elastictransport.WithHeader(headers), + ), + ) + + if err != nil { + return fmt.Errorf("failed to setup elasticsearch connection: %w", err) + } + + cluster.SetClient(es) + } + + return nil +} diff --git a/pkg/cfg/config.go b/pkg/cfg/config.go index 247a507..9f61157 100644 --- a/pkg/cfg/config.go +++ b/pkg/cfg/config.go @@ -17,25 +17,18 @@ along with this program. If not, see . package cfg import ( - "bytes" - "context" - "crypto/tls" - "encoding/json" "errors" "fmt" - "log/slog" - "net/http" "os" + "path/filepath" "reflect" "github.com/alecthomas/repr" - "github.com/elastic/elastic-transport-go/v8/elastictransport" - "github.com/elastic/go-elasticsearch/v9" "gopkg.in/yaml.v3" ) const ( - Version string = `v0.0.20` + Version string = `v0.0.21` ) var ( @@ -43,11 +36,6 @@ var ( APIVERSION, GOVERSION, BUILD, COMMIT, BRANCH string ) -type Cluster struct { - Uri, User, Pass string - ES *elasticsearch.TypedClient -} - type Config struct { ConfigFile string // -c CurrentCluster string // -C @@ -119,8 +107,12 @@ func NewConfig() *Config { return &Config{Clusters: map[string]*Cluster{}} } +func getDefaultPath() string { + return filepath.Join([]string{os.Getenv("HOME"), ".config", "esctl", "config.yaml"}...) +} + func (conf *Config) Init() error { - DefaultConfig := os.Getenv("HOME") + "/.config/esctl/config.yaml" + DefaultConfig := getDefaultPath() switch { case fileExists(DefaultConfig): @@ -141,29 +133,35 @@ func (conf *Config) Init() error { } if conf.CurrentCluster != "" { + // -C specified, set current cluster explicitly, no matter what the config says current, exists := conf.Clusters[conf.CurrentCluster] if !exists { return fmt.Errorf("no cluster with alias %s configured", conf.CurrentCluster) } else { conf.DefaultCluster = current + + // disable all others + for _, cluster := range conf.Clusters { + cluster.Default = false + } + + conf.DefaultCluster.Default = true } } else { + // we need to determine ourselfes if len(conf.Clusters) == 1 { + // ok, just one cluster configured, use this, of course for name, cluster := range conf.Clusters { conf.DefaultCluster = cluster conf.CurrentCluster = name + conf.DefaultCluster.Default = true } } else { + // multiple ones exists, look if one is set as default for name, cluster := range conf.Clusters { - _, err := cluster.ES.Cluster.Health(). - Header("content-type", "application/json"). - Header("accept", "application/json"). - Do(context.Background()) - - if err == nil { + if cluster.Default { conf.DefaultCluster = cluster conf.CurrentCluster = name - break } } } @@ -250,86 +248,21 @@ func (conf *Config) LoadConfig() error { if len(newconf.Clusters) > 0 { conf.Clusters = newconf.Clusters - _, exists := conf.Clusters["default"] - if !exists { - // no "default", just use the first we stumble upon - for _, cluster := range conf.Clusters { + for _, cluster := range conf.Clusters { + if cluster.Default { conf.DefaultCluster = cluster break } - } else { - conf.DefaultCluster = newconf.Clusters["default"] + } + + if conf.DefaultCluster == nil { + conf.DefaultCluster = &Cluster{} } } return nil } -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) -} - -func (conf *Config) SetupES() error { - for _, cluster := range conf.Clusters { - es, err := elasticsearch.NewTyped( - elasticsearch.WithAddresses(cluster.Uri), - elasticsearch.WithBasicAuth(cluster.User, cluster.Pass), - elasticsearch.WithTransportOptions(conf.getTransport()), - ) - - if err != nil { - return fmt.Errorf("failed to setup elasticsearch connection: %w", err) - } - - cluster.ES = es - } - - return nil -} - -// used to print uri, path and body of a request made by the go-client -type DebugTransport struct { - Transport http.RoundTripper -} - -func (t *DebugTransport) RoundTrip(req *http.Request) (*http.Response, error) { - content := "" - contentline := "" - - if req.ContentLength > 0 { - buf := new(bytes.Buffer) - body, _ := req.GetBody() - - _, err := buf.ReadFrom(body) - if err != nil { - return nil, err - } - - var pretty bytes.Buffer - err = json.Indent(&pretty, buf.Bytes(), "", "\t") - if err != nil { - return nil, fmt.Errorf("json parse error: %s", err) - } - - content = pretty.String() - contentline = buf.String() - } - - slog.Info("req", "host", req.URL.Host, "uri", req.URL.Path, "body", content, "bodyline", contentline) - - return t.Transport.RoundTrip(req) -} - func fileExists(filename string) bool { info, err := os.Stat(filename) diff --git a/pkg/cfg/transport.go b/pkg/cfg/transport.go new file mode 100644 index 0000000..386edfb --- /dev/null +++ b/pkg/cfg/transport.go @@ -0,0 +1,78 @@ +/* +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 . +*/ +package cfg + +import ( + "bytes" + "crypto/tls" + "encoding/json" + "fmt" + "log/slog" + "net/http" + + "github.com/elastic/elastic-transport-go/v8/elastictransport" +) + +// used to print uri, path and body of a request made by the go-client +type DebugTransport struct { + Transport http.RoundTripper +} + +func (t *DebugTransport) RoundTrip(req *http.Request) (*http.Response, error) { + content := "" + contentline := "" + + if req.ContentLength > 0 { + buf := new(bytes.Buffer) + body, _ := req.GetBody() + + _, err := buf.ReadFrom(body) + if err != nil { + return nil, err + } + + var pretty bytes.Buffer + err = json.Indent(&pretty, buf.Bytes(), "", "\t") + if err != nil { + return nil, fmt.Errorf("json parse error: %s", err) + } + + content = pretty.String() + contentline = buf.String() + } + + slog.Info("req", "host", req.URL.Host, "uri", req.URL.Path, + "body", content, "bodyline", contentline, + "headers", req.Header, + ) + + 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) +} diff --git a/pkg/es/ccr.go b/pkg/es/ccr.go index beb2fbe..943c785 100644 --- a/pkg/es/ccr.go +++ b/pkg/es/ccr.go @@ -55,12 +55,8 @@ func CcrStatus(conf *cfg.Config, leader, follower string) error { indices := ClusterIndices{} for _, alias := range []string{leader, follower} { - cat := conf.Clusters[alias].ES.Cat.Indices(). - // we need to add custom request headers, required for older ES instances - Header("content-type", "application/json"). - Header("accept", "application/json") - - res, err := cat.Do(context.Background()) + res, err := conf.Clusters[alias].ES().Cat.Indices(). + Do(context.Background()) if err != nil { return fmt.Errorf("failed to get indicies on %s: %s", alias, esErrorString(err)) } @@ -88,9 +84,7 @@ func CcrStatus(conf *cfg.Config, leader, follower string) error { } func CcrRemoteInfo(conf *cfg.Config, index string) error { - res, err := conf.DefaultCluster.ES.Cluster.RemoteInfo(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Cluster.RemoteInfo(). Do(context.Background()) if err != nil { return fmt.Errorf("failed to retrieve follower info: %s", esErrorString(err)) diff --git a/pkg/es/ccr_follower.go b/pkg/es/ccr_follower.go index 37444b1..e324142 100644 --- a/pkg/es/ccr_follower.go +++ b/pkg/es/ccr_follower.go @@ -26,9 +26,7 @@ import ( ) func getRemoteName(conf *cfg.Config) (string, error) { - res, err := conf.DefaultCluster.ES.Cluster.RemoteInfo(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Cluster.RemoteInfo(). Do(context.Background()) if err != nil { return "", fmt.Errorf("failed to retrieve follower info: %s", esErrorString(err)) @@ -89,11 +87,8 @@ func CcrFollowerRenew(conf *cfg.Config, index string) error { } func CcrFollowerResume(conf *cfg.Config, index string) error { - create := conf.DefaultCluster.ES.Ccr.ResumeFollow(index). - Header("content-type", "application/json"). - Header("accept", "application/json") - - _, err := create.Do(context.Background()) + _, err := conf.DefaultCluster.ES().Ccr.ResumeFollow(index). + Do(context.Background()) if err != nil { return fmt.Errorf("failed to resume ccr following: %s", esErrorString(err)) @@ -103,11 +98,8 @@ func CcrFollowerResume(conf *cfg.Config, index string) error { } func CcrFollowerPause(conf *cfg.Config, index string) error { - create := conf.DefaultCluster.ES.Ccr.PauseFollow(index). - Header("content-type", "application/json"). - Header("accept", "application/json") - - _, err := create.Do(context.Background()) + _, err := conf.DefaultCluster.ES().Ccr.PauseFollow(index). + Do(context.Background()) if err != nil { return fmt.Errorf("failed to pause ccr following: %s", esErrorString(err)) @@ -117,11 +109,8 @@ func CcrFollowerPause(conf *cfg.Config, index string) error { } func CcrFollowerUnfollow(conf *cfg.Config, index string) error { - create := conf.DefaultCluster.ES.Ccr.ForgetFollower(index). - Header("content-type", "application/json"). - Header("accept", "application/json") - - _, err := create.Do(context.Background()) + _, err := conf.DefaultCluster.ES().Ccr.ForgetFollower(index). + Do(context.Background()) if err != nil { return fmt.Errorf("failed to unfollow index: %s", esErrorString(err)) @@ -136,11 +125,9 @@ func CcrFollowerAdd(conf *cfg.Config, index string) error { return err } - create := conf.DefaultCluster.ES.Ccr.Follow(index). + create := conf.DefaultCluster.ES().Ccr.Follow(index). LeaderIndex(index). - RemoteCluster(remote). - Header("content-type", "application/json"). - Header("accept", "application/json") + RemoteCluster(remote) if conf.Wait { create.WaitForActiveShards("all") @@ -156,9 +143,7 @@ func CcrFollowerAdd(conf *cfg.Config, index string) error { } func CcrFollowerShow(conf *cfg.Config, index string) error { - res, err := conf.DefaultCluster.ES.Ccr.FollowStats(index). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Ccr.FollowStats(index). Do(context.Background()) if err != nil { return fmt.Errorf("failed to retrieve follower index info: %s", esErrorString(err)) diff --git a/pkg/es/cluster.go b/pkg/es/cluster.go index 0738286..a33e1fe 100644 --- a/pkg/es/cluster.go +++ b/pkg/es/cluster.go @@ -20,7 +20,6 @@ import ( "context" "fmt" "log/slog" - "slices" "strings" "sync" @@ -59,36 +58,37 @@ type apiResponse struct { } func ClusterList(conf *cfg.Config) error { - table := printer.NewTable(conf, 3, len(conf.Clusters)) + table := printer.NewTable(conf, 4, len(conf.Clusters)) - table.Addheaders("cluster", "uri", "default") + table.Addheaders("cluster", "uri", "reachable", "current") idx := 0 - names := make([]string, len(conf.Clusters)) - for name := range conf.Clusters { - names[idx] = name - idx++ - } + for name, cluster := range conf.Clusters { + reachable := "no" + current := "no" - slices.Sort(names) - - for idx, name := range names { - current := name == "default" || name == conf.CurrentCluster - cluster := conf.Clusters[name] - - _, err := cluster.ES.Cluster.Health(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + _, err := cluster.ES().Cluster.Health(). Do(context.Background()) if err == nil { - name = printer.Colorize(conf, "green", name) + reachable = printer.Colorize(conf, "green", "reachable") } - table.Entries[idx] = []string{name, cluster.Uri, fmt.Sprintf("%t", current)} + if cluster.Default { + current = printer.Colorize(conf, "green", "yes") + + if err != nil { + reachable = printer.Colorize(conf, "red", "no") + } + } + + table.Entries[idx] = []string{name, cluster.Uri, reachable, current} + idx++ } + table.Sort() + if err := table.Print(); err != nil { return err } @@ -114,9 +114,9 @@ func ClusterStatus(conf *cfg.Config) error { } for _, cluster := range clusters { - es := conf.DefaultCluster.ES + es := conf.DefaultCluster.ES() if cluster != "default" { - es = conf.Clusters[cluster].ES + es = conf.Clusters[cluster].ES() } responses := make(chan apiResponse, gocount) diff --git a/pkg/es/cluster_settings.go b/pkg/es/cluster_settings.go index 4b0b5cd..4d8878e 100644 --- a/pkg/es/cluster_settings.go +++ b/pkg/es/cluster_settings.go @@ -28,9 +28,7 @@ import ( ) func ClusterSettingsList(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Cluster.GetSettings(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Cluster.GetSettings(). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get cluster settings: %s", esErrorString(err)) @@ -76,7 +74,7 @@ func ClusterSettingsList(conf *cfg.Config) error { } func ClusterSettingsSet(conf *cfg.Config, args cli.Args) error { - put := conf.Clusters[conf.CurrentCluster].ES.Cluster.PutSettings() + put := conf.DefaultCluster.ES().Cluster.PutSettings() for _, arg := range args.Slice() { setting, value := splitArg(arg) @@ -99,11 +97,7 @@ func ClusterSettingsSet(conf *cfg.Config, args cli.Args) error { } } - _, err := put. - Header("content-type", "application/json"). - Header("accept", "application/json"). - Do(context.Background()) - + _, err := put.Do(context.Background()) if err != nil { return fmt.Errorf("failed to set settings: %s", esErrorString(err)) } @@ -117,10 +111,8 @@ func ClusterSettingsSetSingle(conf *cfg.Config, setting, value string) error { return fmt.Errorf("failed to marshall persistent value <%v> to valid JSON: %s", value, err) } - _, err = conf.Clusters[conf.CurrentCluster].ES.Cluster.PutSettings(). + _, err = conf.DefaultCluster.ES().Cluster.PutSettings(). AddPersistent(setting, message). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { diff --git a/pkg/es/cluster_util.go b/pkg/es/cluster_util.go index 3c63bea..d19eb24 100644 --- a/pkg/es/cluster_util.go +++ b/pkg/es/cluster_util.go @@ -37,9 +37,7 @@ const ( ) func checkClusterIsLeader(conf *cfg.Config, leader string) bool { - stats, err := conf.Clusters[leader].ES.Ccr.Stats(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + stats, err := conf.Clusters[leader].ES().Ccr.Stats(). Do(context.Background()) if err != nil { fmt.Printf("failed to get ccr stats from %s: %s", leader, esErrorString(err)) @@ -64,9 +62,7 @@ func checkClusterStatus(conf *cfg.Config, leader, follower string) bool { status := map[string]*health.Response{} for _, cluster := range []string{leader, follower} { - st, err := conf.Clusters[leader].ES.Cluster.Health(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + st, err := conf.Clusters[leader].ES().Cluster.Health(). Do(context.Background()) if err != nil { fmt.Printf("failed to get health from %s: %s", cluster, esErrorString(err)) @@ -123,10 +119,8 @@ func findIlmErrors(conf *cfg.Config, leader, follower string) bool { failed := map[string]map[string]string{} for _, cluster := range []string{leader, follower} { - ilm, err := conf.Clusters[cluster].ES.Ilm.ExplainLifecycle("_all"). + ilm, err := conf.Clusters[cluster].ES().Ilm.ExplainLifecycle("_all"). OnlyManaged(true). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { fmt.Printf("failed to get ilm status from %s: %s", cluster, esErrorString(err)) @@ -192,9 +186,7 @@ func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, f _, followerHasIt := indices[follower][name] if !followerHasIt { // fetch index details - res, err := conf.Clusters[leader].ES.Indices.Get(name). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.Clusters[leader].ES().Indices.Get(name). Do(context.Background()) if err != nil { continue // ignore it then @@ -255,9 +247,7 @@ func findOrphanedIndices(conf *cfg.Config, indices ClusterIndices, leader, follo _, leaderHasIt := indices[leader][name] if !leaderHasIt { // fetch index details - res, err := conf.Clusters[follower].ES.Indices.Get(name). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.Clusters[follower].ES().Indices.Get(name). Do(context.Background()) if err != nil { continue // ignore it then @@ -353,8 +343,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c switch which { case "health": res, err := es.Cluster.Health(). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) ar.health = res @@ -363,8 +351,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c case "info": res, err := es.Info(). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) ar.info = res @@ -373,8 +359,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c case "ccrstats": res, err := es.Ccr.Stats(). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) ar.ccr = res @@ -383,8 +367,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c case "stats": res, err := es.Cluster.Stats(). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) ar.stats = res @@ -393,8 +375,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c case "indices": res, err := es.Cat.Indices(). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) ar.indices = &res @@ -403,8 +383,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c case "tasks": res, err := es.Cat.Tasks(). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) ar.tasks = &res diff --git a/pkg/es/datastream.go b/pkg/es/datastream.go index 0d5677e..a0a1a83 100644 --- a/pkg/es/datastream.go +++ b/pkg/es/datastream.go @@ -30,9 +30,7 @@ import ( ) func DatastreamNames(conf *cfg.Config) ([]string, error) { - res, err := conf.DefaultCluster.ES.Indices.GetDataStream(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Indices.GetDataStream(). Do(context.Background()) if err != nil { return nil, fmt.Errorf("failed to get data streams: %s", esErrorString(err)) @@ -47,9 +45,7 @@ func DatastreamNames(conf *cfg.Config) ([]string, error) { } func DatastreamList(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Indices.GetDataStream(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Indices.GetDataStream(). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get data streams: %s", esErrorString(err)) @@ -129,10 +125,8 @@ func filterDatastreams(conf *cfg.Config, list []types.DataStream) []types.DataSt } func DatastreamShow(conf *cfg.Config, dsname string) error { - res, err := conf.DefaultCluster.ES.Indices.GetDataStream(). + res, err := conf.DefaultCluster.ES().Indices.GetDataStream(). Name(dsname). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get data stream: %s", esErrorString(err)) @@ -144,10 +138,8 @@ func DatastreamShow(conf *cfg.Config, dsname string) error { return errors.New("no data stream found with that name") } - stats, err := conf.DefaultCluster.ES.Indices.DataStreamsStats(). + stats, err := conf.DefaultCluster.ES().Indices.DataStreamsStats(). Name(dsname). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get data stream stats: %s", esErrorString(err)) @@ -213,9 +205,7 @@ func DatastreamShow(conf *cfg.Config, dsname string) error { } func DatastreamCreate(conf *cfg.Config, dsname string) error { - _, err := conf.DefaultCluster.ES.Indices.CreateDataStream(dsname). - Header("content-type", "application/json"). - Header("accept", "application/json"). + _, err := conf.DefaultCluster.ES().Indices.CreateDataStream(dsname). Do(context.Background()) if err != nil { @@ -226,9 +216,7 @@ func DatastreamCreate(conf *cfg.Config, dsname string) error { } func DatastreamDelete(conf *cfg.Config, dsname string) error { - _, err := conf.DefaultCluster.ES.Indices.DeleteDataStream(dsname). - Header("content-type", "application/json"). - Header("accept", "application/json"). + _, err := conf.DefaultCluster.ES().Indices.DeleteDataStream(dsname). Do(context.Background()) if err != nil { diff --git a/pkg/es/doc.go b/pkg/es/doc.go index 6daf4b3..16109c8 100644 --- a/pkg/es/doc.go +++ b/pkg/es/doc.go @@ -47,10 +47,8 @@ func DocAdd(conf *cfg.Config, jsondoc string) error { now := fmt.Sprintf("%d", rand.Int64()) - res, err := conf.DefaultCluster.ES.Create(conf.Index, now). + res, err := conf.DefaultCluster.ES().Create(conf.Index, now). Document(data). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { return fmt.Errorf("failed to create new doc in index %s: %s", conf.Index, esErrorString(err)) @@ -62,9 +60,7 @@ func DocAdd(conf *cfg.Config, jsondoc string) error { } func DocShow(conf *cfg.Config, id string) error { - res, err := conf.DefaultCluster.ES.Get(conf.Index, id). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Get(conf.Index, id). Do(context.Background()) if err != nil { return fmt.Errorf("failed to retrieve doc in index %s: %s", conf.Index, esErrorString(err)) @@ -95,9 +91,7 @@ func DocDelete(conf *cfg.Config, queries []string) error { if len(queries) == 1 && strings.Contains(queries[0], "id=") { id, _ := strings.CutPrefix(queries[0], "id=") - _, err := conf.DefaultCluster.ES.Delete(conf.Index, id). - Header("content-type", "application/json"). - Header("accept", "application/json"). + _, err := conf.DefaultCluster.ES().Delete(conf.Index, id). Do(context.Background()) if err != nil { return fmt.Errorf("failed to delete doc in index %s: %s", conf.Index, esErrorString(err)) @@ -119,10 +113,8 @@ func DocDelete(conf *cfg.Config, queries []string) error { req.Query = queryCaster } - _, err := conf.DefaultCluster.ES.DeleteByQuery(conf.Index). + _, err := conf.DefaultCluster.ES().DeleteByQuery(conf.Index). Request(req). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { return fmt.Errorf("failed to delete docs in index %s: %s", conf.Index, esErrorString(err)) diff --git a/pkg/es/ilm.go b/pkg/es/ilm.go index 0069eee..4383d29 100644 --- a/pkg/es/ilm.go +++ b/pkg/es/ilm.go @@ -32,7 +32,7 @@ import ( ) func IlmRetry(conf *cfg.Config, index string) error { - _, err := conf.DefaultCluster.ES.Ilm.Retry(index). + _, err := conf.DefaultCluster.ES().Ilm.Retry(index). Do(context.Background()) if err != nil { return fmt.Errorf("failed to retry ilm: %s", esErrorString(err)) @@ -42,7 +42,7 @@ func IlmRetry(conf *cfg.Config, index string) error { } func IlmStatus(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Ilm.GetStatus(). + res, err := conf.DefaultCluster.ES().Ilm.GetStatus(). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get ilm status: %s", esErrorString(err)) @@ -54,7 +54,7 @@ func IlmStatus(conf *cfg.Config) error { } func IlmNames(conf *cfg.Config) ([]string, error) { - res, err := conf.DefaultCluster.ES.Ilm.GetLifecycle(). + res, err := conf.DefaultCluster.ES().Ilm.GetLifecycle(). Do(context.Background()) if err != nil { return nil, fmt.Errorf("failed to get ilm policies: %s", esErrorString(err)) @@ -71,7 +71,7 @@ func IlmNames(conf *cfg.Config) ([]string, error) { } func IlmList(conf *cfg.Config, pattern string) error { - ilm := conf.DefaultCluster.ES.Ilm.GetLifecycle() + ilm := conf.DefaultCluster.ES().Ilm.GetLifecycle() if pattern != "" { ilm.FilterPath(pattern) @@ -103,7 +103,7 @@ func IlmList(conf *cfg.Config, pattern string) error { } func IlmShow(conf *cfg.Config, policy string) error { - res, err := conf.DefaultCluster.ES.Ilm.GetLifecycle(). + res, err := conf.DefaultCluster.ES().Ilm.GetLifecycle(). Policy(policy). Do(context.Background()) if err != nil { @@ -165,7 +165,7 @@ func ilmPhaseString(phase *types.Phase, short bool) string { } func IlmExplain(conf *cfg.Config, index string) error { - res, err := conf.DefaultCluster.ES.Ilm.ExplainLifecycle(index). + res, err := conf.DefaultCluster.ES().Ilm.ExplainLifecycle(index). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get ilm state: %s", esErrorString(err)) @@ -198,7 +198,7 @@ func IlmExplain(conf *cfg.Config, index string) error { } func IlmCreate(conf *cfg.Config, policyname string) error { - ilm := conf.DefaultCluster.ES.Ilm.PutLifecycle(policyname) + ilm := conf.DefaultCluster.ES().Ilm.PutLifecycle(policyname) cfg := conf.Ilm var phases types.PhasesVariant = esdsl.NewPhases() diff --git a/pkg/es/index.go b/pkg/es/index.go index a14bc45..f7f1fc0 100644 --- a/pkg/es/index.go +++ b/pkg/es/index.go @@ -35,9 +35,7 @@ import ( // used for completion func IndexNames(conf *cfg.Config) ([]string, error) { - res, err := conf.DefaultCluster.ES.Cat.Indices(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Cat.Indices(). Do(context.Background()) if err != nil { @@ -86,10 +84,7 @@ func filterIndices(conf *cfg.Config, list indices.Response) indices.Response { } func IndexList(conf *cfg.Config) error { - cat := conf.DefaultCluster.ES.Cat.Indices(). - // we need to add custom request headers, required for older ES instances - Header("content-type", "application/json"). - Header("accept", "application/json") + cat := conf.DefaultCluster.ES().Cat.Indices() if conf.Failed { cat = cat.Health(healthstatus.Red) @@ -130,10 +125,7 @@ func IndexList(conf *cfg.Config) error { } func IndexShow(conf *cfg.Config, indexpattern string) error { - res, err := conf.DefaultCluster.ES.Indices.Get(indexpattern). - // we need to add custom request headers, required for older ES instances - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Indices.Get(indexpattern). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get index: %s", esErrorString(err)) @@ -179,9 +171,7 @@ func IndexShow(conf *cfg.Config, indexpattern string) error { func IndexCreate(conf *cfg.Config, index string, mappings []string) error { settings := esdsl.NewIndexSettings() - create := conf.DefaultCluster.ES.Indices.Create(index). - Header("content-type", "application/json"). - Header("accept", "application/json") + create := conf.DefaultCluster.ES().Indices.Create(index) if conf.Wait { create.WaitForActiveShards("all") @@ -229,9 +219,7 @@ func IndexCreate(conf *cfg.Config, index string, mappings []string) error { } func IndexDelete(conf *cfg.Config, index string) error { - _, err := conf.DefaultCluster.ES.Indices.Delete(index). - Header("content-type", "application/json"). - Header("accept", "application/json"). + _, err := conf.DefaultCluster.ES().Indices.Delete(index). Do(context.Background()) if err != nil { return fmt.Errorf("failed to delete index: %s", esErrorString(err)) @@ -241,11 +229,8 @@ func IndexDelete(conf *cfg.Config, index string) error { } func IndexClose(conf *cfg.Config, index string) error { - create := conf.DefaultCluster.ES.Indices.Close(index). - Header("content-type", "application/json"). - Header("accept", "application/json") - - _, err := create.Do(context.Background()) + _, err := conf.DefaultCluster.ES().Indices.Close(index). + Do(context.Background()) if err != nil { return fmt.Errorf("failed to close index: %s", esErrorString(err)) @@ -255,12 +240,10 @@ func IndexClose(conf *cfg.Config, index string) error { } func IndexAllocation(conf *cfg.Config, index string) error { - res, err := conf.DefaultCluster.ES.Cluster.AllocationExplain(). + res, err := conf.DefaultCluster.ES().Cluster.AllocationExplain(). Index(index). Primary(conf.Primary). Shard(conf.Shards). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get index allocation explain: %s", esErrorString(err)) @@ -301,11 +284,9 @@ func IndexAllocation(conf *cfg.Config, index string) error { func IndexModify(conf *cfg.Config, index string) error { settings := esdsl.NewIndexSettings().NumberOfReplicas(strconv.Itoa(conf.Replicas)) - _, err := conf.DefaultCluster.ES.Indices.PutSettings(). + _, err := conf.DefaultCluster.ES().Indices.PutSettings(). Indices(index). Index(settings). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { return fmt.Errorf("failed to modify index settings: %s", esErrorString(err)) @@ -315,11 +296,9 @@ func IndexModify(conf *cfg.Config, index string) error { } func IndexFields(conf *cfg.Config, index string) error { - res, err := conf.DefaultCluster.ES.FieldCaps(). + res, err := conf.DefaultCluster.ES().FieldCaps(). Index(index). Fields("*"). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { return fmt.Errorf("failed to retrieve field capabilties: %s", esErrorString(err)) diff --git a/pkg/es/index_alias.go b/pkg/es/index_alias.go index 0f25085..0f0f49c 100644 --- a/pkg/es/index_alias.go +++ b/pkg/es/index_alias.go @@ -29,9 +29,7 @@ import ( // FIXME: add filter support, see IndexCreate mapping func IndexAliasCreate(conf *cfg.Config, index, alias string) error { - res, err := conf.DefaultCluster.ES.Indices.PutAlias(index, alias). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Indices.PutAlias(index, alias). Do(context.Background()) slog.Debug("create alias", "result", res) @@ -51,10 +49,8 @@ func IndexAliasList(conf *cfg.Config) error { filter = *regexp.MustCompile(conf.Filter[0]) } - res, err := conf.DefaultCluster.ES.Indices.GetAlias(). + res, err := conf.DefaultCluster.ES().Indices.GetAlias(). Index("_all"). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { @@ -100,9 +96,7 @@ func IndexAliasList(conf *cfg.Config) error { } func IndexAliasDelete(conf *cfg.Config, index, alias string) error { - res, err := conf.DefaultCluster.ES.Indices.DeleteAlias(index, alias). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Indices.DeleteAlias(index, alias). Do(context.Background()) slog.Debug("delete alias", "result", res) diff --git a/pkg/es/index_template.go b/pkg/es/index_template.go index d50baca..a2aafcd 100644 --- a/pkg/es/index_template.go +++ b/pkg/es/index_template.go @@ -33,9 +33,7 @@ import ( // used for completion func IndexTemplateList(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Indices.GetIndexTemplate(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate(). Do(context.Background()) if err != nil { @@ -65,10 +63,8 @@ func IndexTemplateList(conf *cfg.Config) error { } func IndexTemplateShow(conf *cfg.Config, tplname string) error { - res, err := conf.DefaultCluster.ES.Indices.GetIndexTemplate(). + res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate(). Name(tplname). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { @@ -158,9 +154,7 @@ func IndexTemplateCreate(conf *cfg.Config, name string, mappings []string) error settings := esdsl.NewIndexSettings() maps := esdsl.NewIndexTemplateMapping() - create := conf.DefaultCluster.ES.Indices.PutIndexTemplate(name). - Header("content-type", "application/json"). - Header("accept", "application/json") + create := conf.DefaultCluster.ES().Indices.PutIndexTemplate(name) if conf.Shards > 0 { settings = settings.NumberOfShards(strconv.Itoa(conf.Shards)) @@ -243,10 +237,8 @@ func IndexTemplateModify(conf *cfg.Config, name string, mappings []string) error maps := esdsl.NewIndexTemplateMapping() // load existing index mapping - res, err := conf.DefaultCluster.ES.Indices.GetIndexTemplate(). + res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate(). Name(name). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil { @@ -262,9 +254,7 @@ func IndexTemplateModify(conf *cfg.Config, name string, mappings []string) error tpl := res.IndexTemplates[0] // our modify PUT request - modify := conf.DefaultCluster.ES.Indices.PutIndexTemplate(name). - Header("content-type", "application/json"). - Header("accept", "application/json") + modify := conf.DefaultCluster.ES().Indices.PutIndexTemplate(name) // load existing settings, if any settings := tpl.IndexTemplate.Template.Settings @@ -355,9 +345,7 @@ func IndexTemplateModify(conf *cfg.Config, name string, mappings []string) error } func IndexTemplateDelete(conf *cfg.Config, name string) error { - _, err := conf.DefaultCluster.ES.Indices.DeleteIndexTemplate(name). - Header("content-type", "application/json"). - Header("accept", "application/json"). + _, err := conf.DefaultCluster.ES().Indices.DeleteIndexTemplate(name). Do(context.Background()) if err != nil { return fmt.Errorf("failed to delete index template: %s", esErrorString(err)) @@ -370,7 +358,7 @@ func IndexTemplateDelete(conf *cfg.Config, name string) error { // find indices matching those, find their associated aliases and // rollover all we find. Only run when conf.Rollover==true func rolloverAliasIndexTemplate(conf *cfg.Config, name string) error { - res, err := conf.DefaultCluster.ES.Indices.GetIndexTemplate(). + res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate(). Name(name). Do(context.Background()) @@ -389,7 +377,7 @@ func rolloverAliasIndexTemplate(conf *cfg.Config, name string) error { // find all aliases matching the patterns for _, pattern := range patterns { - res, err := conf.DefaultCluster.ES.Indices.ResolveIndex(pattern). + res, err := conf.DefaultCluster.ES().Indices.ResolveIndex(pattern). Do(context.Background()) if err != nil { return fmt.Errorf("failed to resolve index pattern: %s", esErrorString(err)) diff --git a/pkg/es/node.go b/pkg/es/node.go index 2676c79..a5b91a8 100644 --- a/pkg/es/node.go +++ b/pkg/es/node.go @@ -27,7 +27,7 @@ import ( func NodeList(conf *cfg.Config) error { // get nodes - nodes, err := conf.DefaultCluster.ES.Cat.Nodes().Do(context.Background()) + nodes, err := conf.DefaultCluster.ES().Cat.Nodes().Do(context.Background()) if err != nil { return fmt.Errorf("failed to get nodes: %s", esErrorString(err)) } diff --git a/pkg/es/role.go b/pkg/es/role.go index de22035..ffff478 100644 --- a/pkg/es/role.go +++ b/pkg/es/role.go @@ -28,7 +28,7 @@ import ( ) func RoleNames(conf *cfg.Config) ([]string, error) { - res, err := conf.DefaultCluster.ES.Security.GetRole(). + res, err := conf.DefaultCluster.ES().Security.GetRole(). Do(context.Background()) if err != nil { return nil, fmt.Errorf("failed to get roles: %s", esErrorString(err)) @@ -45,7 +45,7 @@ func RoleNames(conf *cfg.Config) ([]string, error) { } func RoleList(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Security.GetRole(). + res, err := conf.DefaultCluster.ES().Security.GetRole(). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get roles: %s", esErrorString(err)) @@ -71,7 +71,7 @@ func RoleList(conf *cfg.Config) error { } func RoleShow(conf *cfg.Config, rolename string) error { - res, err := conf.DefaultCluster.ES.Security.GetRole(). + res, err := conf.DefaultCluster.ES().Security.GetRole(). Name(rolename). Do(context.Background()) if err != nil { diff --git a/pkg/es/role_diff.go b/pkg/es/role_diff.go index 9bcd99c..bb5898c 100644 --- a/pkg/es/role_diff.go +++ b/pkg/es/role_diff.go @@ -207,7 +207,7 @@ func RoleDiff(conf *cfg.Config, csvfile, role string) error { return RoleDiffSingle(conf, csvfile, role) } - res, err := conf.DefaultCluster.ES.Security.GetRole(). + res, err := conf.DefaultCluster.ES().Security.GetRole(). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get roles: %s", esErrorString(err)) @@ -242,7 +242,7 @@ func RoleDiff(conf *cfg.Config, csvfile, role string) error { } func getRoleMappingGroups(conf *cfg.Config, rolename string) ([]string, error) { - mappings, err := conf.DefaultCluster.ES.Security.GetRoleMapping(). + mappings, err := conf.DefaultCluster.ES().Security.GetRoleMapping(). Do(context.Background()) if err != nil { return nil, fmt.Errorf("failed to get role mappings: %s", esErrorString(err)) @@ -275,7 +275,7 @@ func compareSlices(name string, a, b []string) { } func RoleDiffSingle(conf *cfg.Config, csvfile, rolename string) error { - res, err := conf.DefaultCluster.ES.Security.GetRole(). + res, err := conf.DefaultCluster.ES().Security.GetRole(). Name(rolename). Do(context.Background()) if err != nil { diff --git a/pkg/es/rollover.go b/pkg/es/rollover.go index 014845e..f7ca971 100644 --- a/pkg/es/rollover.go +++ b/pkg/es/rollover.go @@ -47,7 +47,7 @@ func RolloverConditions(conf *cfg.Config) types.RolloverConditionsVariant { } func RolloverAlias(conf *cfg.Config, alias string) (*rollover.Response, error) { - roll := conf.DefaultCluster.ES.Indices.Rollover(alias) + roll := conf.DefaultCluster.ES().Indices.Rollover(alias) if conf.Wait { roll.WaitForActiveShards("all") diff --git a/pkg/es/search.go b/pkg/es/search.go index e171400..f93c5f1 100644 --- a/pkg/es/search.go +++ b/pkg/es/search.go @@ -50,7 +50,7 @@ func Search(conf *cfg.Config, queries []string) error { return validateSearch(conf, queries) } - searchEs := conf.DefaultCluster.ES.Search().Index(conf.Index) + searchEs := conf.DefaultCluster.ES().Search().Index(conf.Index) queryCaster, err := prepareQuery(conf, queries) if err != nil { @@ -121,7 +121,7 @@ func explain(res *types.ExplanationDetail, indent string) { } func validateSearch(conf *cfg.Config, queries []string) error { - validate := conf.DefaultCluster.ES.Indices.ValidateQuery() + validate := conf.DefaultCluster.ES().Indices.ValidateQuery() queryCaster, err := prepareQuery(conf, queries) if err != nil { @@ -150,7 +150,7 @@ func validateSearch(conf *cfg.Config, queries []string) error { } func Debug(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Search(). + res, err := conf.DefaultCluster.ES().Search(). Index(conf.Index). Size(0). Aggregations(map[string]types.Aggregations{ @@ -187,18 +187,18 @@ func searchOnce(conf *cfg.Config, search *search.Search) error { // https://www.elastic.co/docs/reference/elasticsearch/clients/go/using-the-api/searching#_pit_search_after func searchPit(conf *cfg.Config, req *search.Request) error { ctx := context.Background() - pit, err := conf.DefaultCluster.ES.OpenPointInTime(conf.Index).KeepAlive("1m").Do(ctx) + pit, err := conf.DefaultCluster.ES().OpenPointInTime(conf.Index).KeepAlive("1m").Do(ctx) if err != nil { return fmt.Errorf("failed to open point-in-time request for search: %s", err) } defer func() { - _, err := conf.DefaultCluster.ES.ClosePointInTime().Id(pit.Id).Do(ctx) + _, err := conf.DefaultCluster.ES().ClosePointInTime().Id(pit.Id).Do(ctx) if err != nil { log.Fatalf("failed to close PIT: %s", err) } }() - search := conf.DefaultCluster.ES.Search(). + search := conf.DefaultCluster.ES().Search(). Request(req). Pit(esdsl.NewPointInTimeReference(). Id(pit.Id). diff --git a/pkg/es/search_filter.go b/pkg/es/search_filter.go index 8531a43..8a364fc 100644 --- a/pkg/es/search_filter.go +++ b/pkg/es/search_filter.go @@ -243,7 +243,7 @@ func addFilters(conf *cfg.Config) ([]types.QueryVariant, error) { // check if the conf.SortBy field (default: @timestamp) is searchable // by using the field capabilities API. func addSort(conf *cfg.Config, search *search.Search) *search.Search { - res, err := conf.DefaultCluster.ES.FieldCaps(). + res, err := conf.DefaultCluster.ES().FieldCaps(). Index(conf.Index). Fields(conf.SortBy). Do(context.Background()) diff --git a/pkg/es/shard.go b/pkg/es/shard.go index 2addb71..24395fa 100644 --- a/pkg/es/shard.go +++ b/pkg/es/shard.go @@ -83,9 +83,7 @@ func filterShards(conf *cfg.Config, shardlist shards.Response) shards.Response { } func ShardList(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Cat.Shards(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Cat.Shards(). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get shards: %s", esErrorString(err)) @@ -135,9 +133,7 @@ func printShards(conf *cfg.Config, shardlist shards.Response) error { } func ShardShow(conf *cfg.Config, index string) error { - res, err := conf.DefaultCluster.ES.Cat.Shards().Index(index). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Cat.Shards().Index(index). Do(context.Background()) if err != nil { return fmt.Errorf("failed to get shards: %s", esErrorString(err)) diff --git a/pkg/es/snapshot.go b/pkg/es/snapshot.go index 78c10df..2ef680c 100644 --- a/pkg/es/snapshot.go +++ b/pkg/es/snapshot.go @@ -43,7 +43,7 @@ type Snapshot struct { func SnapshotList(conf *cfg.Config) error { // get partial indicies - ires, err := conf.DefaultCluster.ES.Cat.Indices().Do(context.Background()) + ires, err := conf.DefaultCluster.ES().Cat.Indices().Do(context.Background()) if err != nil { return fmt.Errorf("failed to get indicies: %s", esErrorString(err)) } @@ -56,7 +56,7 @@ func SnapshotList(conf *cfg.Config) error { } // get snapshots - sres, err := conf.DefaultCluster.ES.Cat.Snapshots().Do(context.Background()) + sres, err := conf.DefaultCluster.ES().Cat.Snapshots().Do(context.Background()) if err != nil { return fmt.Errorf("failed to get snapshots: %s", esErrorString(err)) } @@ -106,7 +106,7 @@ func SnapshotList(conf *cfg.Config) error { } func SnapshotShow(conf *cfg.Config, snapshot string) error { - res, err := conf.DefaultCluster.ES.Snapshot.Get("*", snapshot).Do(context.Background()) + res, err := conf.DefaultCluster.ES().Snapshot.Get("*", snapshot).Do(context.Background()) if err != nil { return fmt.Errorf("failed to get snapshot: %s", esErrorString(err)) } diff --git a/pkg/es/task.go b/pkg/es/task.go index 813bcc6..640bdb6 100644 --- a/pkg/es/task.go +++ b/pkg/es/task.go @@ -28,9 +28,7 @@ import ( ) func TaskList(conf *cfg.Config) error { - res, err := conf.DefaultCluster.ES.Cat.Tasks(). - Header("content-type", "application/json"). - Header("accept", "application/json"). + res, err := conf.DefaultCluster.ES().Cat.Tasks(). Do(context.Background()) if err != nil { @@ -65,10 +63,8 @@ func TaskList(conf *cfg.Config) error { } func TaskCancel(conf *cfg.Config, taskid string) error { - _, err := conf.DefaultCluster.ES.Tasks.Cancel(). + _, err := conf.DefaultCluster.ES().Tasks.Cancel(). TaskId(taskid). - Header("content-type", "application/json"). - Header("accept", "application/json"). Do(context.Background()) if err != nil {