Compare commits

..

8 Commits

32 changed files with 555 additions and 445 deletions

View File

@@ -12,6 +12,13 @@ The programming language used for this project will always be
# Contributing # Contributing
Everyone can contribute to this project except you are using "AI"
(LLM) generated content (code, issue text, etc), even if only a part
of it is. This is a human made project. It is being made to have fun
at work and learning coding and APIs. If you want fast results without
the hassle to invest time then please refrain from
contributing.
You can contribute to this project in various ways: You can contribute to this project in various ways:
## Open an issue ## Open an issue

211
README.md
View File

@@ -1,4 +1,8 @@
[![status-badge](https://ci.codeberg.org/api/badges/16999/status.svg)](https://ci.codeberg.org/repos/16999) [![status-badge](https://ci.codeberg.org/api/badges/16999/status.svg)](https://ci.codeberg.org/repos/16999)
[![humanmande](https://img.shields.io/badge/human-made-green)](CONTRIBUTING.md)
[![License](https://img.shields.io/badge/license-GPL-blue.svg)](https://codeberg.org/scip/esctl/blob/master/LICENSE)
[![Go Report Card](https://goreportcard.com/badge/codeberg.org/scip/esctl)](https://goreportcard.com/report/codeberg.org/scip/esctl)
[![Latest Release](https://flat.badgen.net/codeberg/release/scip/esctl)](https://codeberg.org/scip/esctl/releases)
# esctl # esctl
@@ -26,13 +30,19 @@ Features:
- 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,
display field mappings etc. display field mappings etc. Automatic rollover of aliases supported.
- Index template management: create, modify, delete etc
- Index alias management: create, modify, delete etc
- Node management: only list nodes yet. - Node management: only list nodes yet.
- Shard management: only list shards yet. - Shard management: only list shards yet.
- Snapshot management: only list snapshots yet. - Snapshot management: only list snapshots yet.
- ILM management: list, create, delete etc.
- Task management: list and cancel tasks
- Role management: only list roles yet. There's also a `role diff` - Role management: only list roles yet. There's also a `role diff`
subcommand, which is for internal use. It can be used to verify if subcommand, which is for internal use. It can be used to verify if
role defs in a CSV match the deployed roles. role defs in a CSV match the deployed roles.
- API documentation (`api list` and `api show <path>`) with
interactive markdown pager for endpoint documentation.
- Repl: this is an interactive REPL (read eval print loop) towards the - Repl: this is an interactive REPL (read eval print loop) towards the
elasticsearch API. You can run API calls on the current selected elasticsearch API. You can run API calls on the current selected
cluster w/o the hassle to specify the whole url, credentials etc. It cluster w/o the hassle to specify the whole url, credentials etc. It
@@ -51,87 +61,92 @@ Features:
Command tree: Command tree:
```console ```console
api api - api access and documentation
list list - list index of API calls
repl show - show an API doc
show repl - interactive API repl
ccr ccr - manage cross cluster replication
follower status - cross cluster replication status (yaml config with 2 clusters required)
add pause - pause shard allocation
delete resume - resume shard allocation
pause follower - manage ccr follower indices
renew show - show ccr follower index details
resume add - add ccr follower index
show delete - delete ccr follower index
unfollow unfollow - unfollow ccr follower index
info pause - pause ccr index to follow
pause resume - resume ccr index to follow
resume renew - renew ccr follower index
status info - show ccr remote info
cluster cluster - manage cluster[s]
list status - show cluster status
settings switch - set current elasticsearch cluster
list list - list configured clusters
set settings - cluster settings management
status list - show cluster settings
datastream set - set|update cluster settings
create datastream - manage data streams
delete list - list indicies
list show - show details about an data stream
rollover create - create a new data stream
show delete - delete a data stream
debug rollover - roll over a data stream
doc doc - manage documents
add add - add JSON document index
delete show - show a JSON document
show delete - delete JSON document[s] from index[es]
help ilm - manage index lifecycle
help-jsonpath retry - retry applying an ILM profile to an index
ilm status - get the current index lifecycle management status
create list - list index lifecycle policies
list show - show details about an index lifecycle policy
retry create - create a index lifecycle policy
show index - manage indicies
status list - list indicies
index show - show details about an index
alias create - create a new index
create delete - delete an index
delete close - close an index
list allocation - explain index allocation
rollover modify - modify an index
allocation fields - show info about field capabilities
close ilm - show ilm status
create alias - manage index aliases
delete create - create an index alias
fields list - list index aliases
ilm delete - delete an index alias
list rollover - roll over an index alias
modify template - manage index templates
show list - list index templates
template show - show details about an index template
create create - create a new index template
delete modify - modify a new index template
list delete - delete an index template
modify node - manage nodes
show list - list nodes
node show - show details about a node
list role - manage roles
show list - list roles
role show - show details about a role
diff diff - show differences between roles and CSV baseline
list search - search within an index
show shard - manage shards
search list - list shards
shard show - show details about a shard
list snapshot - manage snapshots
show list - list snapshots
snapshot show - show details about a snapshot
list task - manage tasks
show list - list tasks
task cancel - cancel running task
cancel version - show esctl version information
list debug - developer only
version help-jsonpath - show jsonpath help
completion - Output shell completion script for bash, zsh, fish, or Powershell
bash - Output bash completion script
zsh - Output zsh completion script
fish - Output fish completion script
pwsh - Output pwsh completion script
``` ```
Configure `esctl` with environment variables: Configure `esctl` with environment variables:
@@ -144,7 +159,7 @@ Or create a config file such as this:
```yaml ```yaml
clusters: clusters:
default: foobar:
uri: https://es.foo.bar:9200/ uri: https://es.foo.bar:9200/
user: elastic user: elastic
pass: 123456 pass: 123456
@@ -158,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 default config file in `~/.config/esctl/config.yaml`. In this case you
can omit `-c ...`. 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. 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 ## Installation
The tool does not have any dependencies. Just download the binary for The tool does not have any dependencies. Just download the binary for

View File

@@ -28,7 +28,7 @@ import (
func Api(conf *cfg.Config) *cli.Command { func Api(conf *cfg.Config) *cli.Command {
return &cli.Command{ return &cli.Command{
Name: "api", Name: "api",
Usage: "manage apiuments", Usage: "api access and documentation",
Commands: []*cli.Command{ Commands: []*cli.Command{
ApiList(conf), ApiList(conf),

View File

@@ -18,6 +18,7 @@ package cmd
import ( import (
"context" "context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es" "codeberg.org/scip/esctl/pkg/es"
@@ -33,6 +34,7 @@ func Cluster(conf *cfg.Config) *cli.Command {
Commands: []*cli.Command{ Commands: []*cli.Command{
ClusterStatus(conf), ClusterStatus(conf),
ClusterSwitch(conf),
ClusterList(conf), ClusterList(conf),
ClusterSettings(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 <cluster-name>",
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)
},
}
}

View File

@@ -22,6 +22,7 @@ import (
golog "log" golog "log"
"os" "os"
"runtime/pprof" "runtime/pprof"
"strings"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es" "codeberg.org/scip/esctl/pkg/es"
@@ -41,11 +42,11 @@ func Finish(err error) int {
func Main() int { func Main() int {
conf := cfg.NewConfig() conf := cfg.NewConfig()
tree := false
cmd := &cli.Command{ cmd := &cli.Command{
Name: "esctl", Name: "esctl",
Usage: "manage elasticsearch from cli", Usage: "manage elasticsearch from cli",
//Version: cfg.Version,
EnableShellCompletion: true, EnableShellCompletion: true,
Flags: []cli.Flag{ Flags: []cli.Flag{
@@ -63,6 +64,20 @@ func Main() int {
Usage: "enable HTTP debugging", Usage: "enable HTTP debugging",
Destination: &conf.DebugHTTP, Destination: &conf.DebugHTTP,
}, },
&cli.BoolFlag{
Name: "show-command-tree",
Value: false,
Usage: "generate a command tree",
Destination: &tree,
Hidden: true,
},
&cli.BoolFlag{
Name: "align-ints",
Aliases: []string{"I"},
Value: false,
Usage: "right align integers in tabular output",
Destination: &conf.AlignInts,
},
&cli.StringFlag{ &cli.StringFlag{
Name: "config", Name: "config",
Aliases: []string{"c"}, Aliases: []string{"c"},
@@ -94,25 +109,33 @@ func Main() int {
}, },
Commands: []*cli.Command{ Commands: []*cli.Command{
Search(conf), Api(conf),
Cluster(conf),
Ccr(conf), Ccr(conf),
Index(conf), Cluster(conf),
Ilm(conf),
Datastream(conf), Datastream(conf),
Doc(conf),
Ilm(conf),
Index(conf),
Node(conf),
Roles(conf),
Search(conf),
Shard(conf), Shard(conf),
Snapshot(conf), Snapshot(conf),
Node(conf), Task(conf),
Doc(conf),
Version(conf), Version(conf),
Debug(conf), Debug(conf),
Roles(conf),
Task(conf),
HelpJsonPath(conf), HelpJsonPath(conf),
Api(conf),
}, },
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) { Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {
if tree {
if err := Tree(cmd); err != nil {
return nil, err
}
os.Exit(0)
}
if err := conf.Init(); err != nil { if err := conf.Init(); err != nil {
if len(os.Args) > 1 { if len(os.Args) > 1 {
return nil, err return nil, err
@@ -237,3 +260,23 @@ func Debug(conf *cfg.Config) *cli.Command {
}, },
} }
} }
func Tree(cmd *cli.Command) error {
max := 20
return cmd.Walk(func(cmd *cli.Command) error {
path := cmd.Path()
command := path[len(path)-1]
if len(path) == 1 || command == "help" {
return nil
}
indent := strings.Repeat(" ", len(path[1:])-1)
space := strings.Repeat(" ", max-(len(command)+len(indent)))
fmt.Printf("%s%s %s - %s\n", indent, command, space, cmd.Usage)
return nil
})
}

View File

@@ -1,31 +0,0 @@
#!/bin/bash
getcommands() {
local command
command="$*"
if ! test "$command" = "help"; then
./esctl $command -h | sed -e '/^COMMANDS:/,/^$/!d;//d' -e 's/^ //' -e 's/[, ].*//' | sort
fi
}
printcommands() {
local indent basecommand commands command
indent="$1"
basecommand="$2"
commands="$3"
for command in $commands; do
echo "$indent" "$command"
commands=$(getcommands $basecommand $command)
if test -n "$commands"; then
printcommands "${indent} " "$basecommand $command" "$commands"
fi
done
}
commands=$(getcommands "")
printcommands "" "" "$commands"

2
go.mod
View File

@@ -32,7 +32,7 @@ require (
github.com/olekukonko/tablewriter v1.1.4 github.com/olekukonko/tablewriter v1.1.4
github.com/tidwall/gjson v1.19.0 github.com/tidwall/gjson v1.19.0
github.com/tlinden/yadu v0.1.3 github.com/tlinden/yadu v0.1.3
github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de github.com/urfave/cli/v3 v3.10.1-0.20260623012112-f980ca84bf65
gopkg.in/yaml.v3 v3.0.1 gopkg.in/yaml.v3 v3.0.1
) )

2
go.sum
View File

@@ -189,6 +189,8 @@ github.com/tlinden/yadu v0.1.3 h1:5cRCUmj+l5yvlM2irtpFBIJwVV2DPEgYSaWvF19FtcY=
github.com/tlinden/yadu v0.1.3/go.mod h1:l3bRmHKL9zGAR6pnBHY2HRPxBecf7L74BoBgOOpTcUA= github.com/tlinden/yadu v0.1.3/go.mod h1:l3bRmHKL9zGAR6pnBHY2HRPxBecf7L74BoBgOOpTcUA=
github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de h1:ESKPiS7inVoBnv4FmgGNZdjWBI/wmvaragoyD3D9nM4= github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de h1:ESKPiS7inVoBnv4FmgGNZdjWBI/wmvaragoyD3D9nM4=
github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso= github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso=
github.com/urfave/cli/v3 v3.10.1-0.20260623012112-f980ca84bf65 h1:NXitXXO9DupDLEFj8hIyYtsubRxSFTp5JXqohElT0nU=
github.com/urfave/cli/v3 v3.10.1-0.20260623012112-f980ca84bf65/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso=
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no= github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no=
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e/go.mod h1:RbqR21r5mrJuqunuUZ/Dhy/avygyECGrLceyNeo4LiM= github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e/go.mod h1:RbqR21r5mrJuqunuUZ/Dhy/avygyECGrLceyNeo4LiM=
go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA=

126
pkg/cfg/cluster.go Normal file
View File

@@ -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 <http://www.gnu.org/licenses/>.
*/
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 <name>' 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
}

View File

@@ -17,25 +17,18 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
package cfg package cfg
import ( import (
"bytes"
"context"
"crypto/tls"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"log/slog"
"net/http"
"os" "os"
"path/filepath"
"reflect" "reflect"
"github.com/alecthomas/repr" "github.com/alecthomas/repr"
"github.com/elastic/elastic-transport-go/v8/elastictransport"
"github.com/elastic/go-elasticsearch/v9"
"gopkg.in/yaml.v3" "gopkg.in/yaml.v3"
) )
const ( const (
Version string = `v0.0.20` Version string = `v0.0.21`
) )
var ( var (
@@ -43,11 +36,6 @@ var (
APIVERSION, GOVERSION, BUILD, COMMIT, BRANCH string APIVERSION, GOVERSION, BUILD, COMMIT, BRANCH string
) )
type Cluster struct {
Uri, User, Pass string
ES *elasticsearch.TypedClient
}
type Config struct { type Config struct {
ConfigFile string // -c ConfigFile string // -c
CurrentCluster string // -C CurrentCluster string // -C
@@ -57,6 +45,7 @@ type Config struct {
DefaultCluster *Cluster DefaultCluster *Cluster
HaveJQ bool // determined at runtime by ourselfes HaveJQ bool // determined at runtime by ourselfes
ProfileFile string // for internal use (golang profiling) ProfileFile string // for internal use (golang profiling)
AlignInts bool // -I
Index string // index: -i Index string // index: -i
Failed, Partials bool // index: flags Failed, Partials bool // index: flags
@@ -119,8 +108,12 @@ func NewConfig() *Config {
return &Config{Clusters: map[string]*Cluster{}} 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 { func (conf *Config) Init() error {
DefaultConfig := os.Getenv("HOME") + "/.config/esctl/config.yaml" DefaultConfig := getDefaultPath()
switch { switch {
case fileExists(DefaultConfig): case fileExists(DefaultConfig):
@@ -141,29 +134,35 @@ func (conf *Config) Init() error {
} }
if conf.CurrentCluster != "" { if conf.CurrentCluster != "" {
// -C specified, set current cluster explicitly, no matter what the config says
current, exists := conf.Clusters[conf.CurrentCluster] current, exists := conf.Clusters[conf.CurrentCluster]
if !exists { if !exists {
return fmt.Errorf("no cluster with alias %s configured", conf.CurrentCluster) return fmt.Errorf("no cluster with alias %s configured", conf.CurrentCluster)
} else { } else {
conf.DefaultCluster = current conf.DefaultCluster = current
// disable all others
for _, cluster := range conf.Clusters {
cluster.Default = false
}
conf.DefaultCluster.Default = true
} }
} else { } else {
// we need to determine ourselfes
if len(conf.Clusters) == 1 { if len(conf.Clusters) == 1 {
// ok, just one cluster configured, use this, of course
for name, cluster := range conf.Clusters { for name, cluster := range conf.Clusters {
conf.DefaultCluster = cluster conf.DefaultCluster = cluster
conf.CurrentCluster = name conf.CurrentCluster = name
conf.DefaultCluster.Default = true
} }
} else { } else {
// multiple ones exists, look if one is set as default
for name, cluster := range conf.Clusters { for name, cluster := range conf.Clusters {
_, err := cluster.ES.Cluster.Health(). if cluster.Default {
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err == nil {
conf.DefaultCluster = cluster conf.DefaultCluster = cluster
conf.CurrentCluster = name conf.CurrentCluster = name
break
} }
} }
} }
@@ -250,86 +249,21 @@ func (conf *Config) LoadConfig() error {
if len(newconf.Clusters) > 0 { if len(newconf.Clusters) > 0 {
conf.Clusters = newconf.Clusters conf.Clusters = newconf.Clusters
_, exists := conf.Clusters["default"] for _, cluster := range conf.Clusters {
if !exists { if cluster.Default {
// no "default", just use the first we stumble upon
for _, cluster := range conf.Clusters {
conf.DefaultCluster = cluster conf.DefaultCluster = cluster
break break
} }
} else { }
conf.DefaultCluster = newconf.Clusters["default"]
if conf.DefaultCluster == nil {
conf.DefaultCluster = &Cluster{}
} }
} }
return nil 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 { func fileExists(filename string) bool {
info, err := os.Stat(filename) info, err := os.Stat(filename)

78
pkg/cfg/transport.go Normal file
View File

@@ -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 <http://www.gnu.org/licenses/>.
*/
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)
}

View File

@@ -55,12 +55,8 @@ func CcrStatus(conf *cfg.Config, leader, follower string) error {
indices := ClusterIndices{} indices := ClusterIndices{}
for _, alias := range []string{leader, follower} { for _, alias := range []string{leader, follower} {
cat := conf.Clusters[alias].ES.Cat.Indices(). res, err := conf.Clusters[alias].ES().Cat.Indices().
// we need to add custom request headers, required for older ES instances Do(context.Background())
Header("content-type", "application/json").
Header("accept", "application/json")
res, err := cat.Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get indicies on %s: %s", alias, esErrorString(err)) 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 { func CcrRemoteInfo(conf *cfg.Config, index string) error {
res, err := conf.DefaultCluster.ES.Cluster.RemoteInfo(). res, err := conf.DefaultCluster.ES().Cluster.RemoteInfo().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to retrieve follower info: %s", esErrorString(err)) return fmt.Errorf("failed to retrieve follower info: %s", esErrorString(err))

View File

@@ -26,9 +26,7 @@ import (
) )
func getRemoteName(conf *cfg.Config) (string, error) { func getRemoteName(conf *cfg.Config) (string, error) {
res, err := conf.DefaultCluster.ES.Cluster.RemoteInfo(). res, err := conf.DefaultCluster.ES().Cluster.RemoteInfo().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return "", fmt.Errorf("failed to retrieve follower info: %s", esErrorString(err)) 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 { func CcrFollowerResume(conf *cfg.Config, index string) error {
create := conf.DefaultCluster.ES.Ccr.ResumeFollow(index). _, err := conf.DefaultCluster.ES().Ccr.ResumeFollow(index).
Header("content-type", "application/json"). Do(context.Background())
Header("accept", "application/json")
_, err := create.Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to resume ccr following: %s", esErrorString(err)) 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 { func CcrFollowerPause(conf *cfg.Config, index string) error {
create := conf.DefaultCluster.ES.Ccr.PauseFollow(index). _, err := conf.DefaultCluster.ES().Ccr.PauseFollow(index).
Header("content-type", "application/json"). Do(context.Background())
Header("accept", "application/json")
_, err := create.Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to pause ccr following: %s", esErrorString(err)) 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 { func CcrFollowerUnfollow(conf *cfg.Config, index string) error {
create := conf.DefaultCluster.ES.Ccr.ForgetFollower(index). _, err := conf.DefaultCluster.ES().Ccr.ForgetFollower(index).
Header("content-type", "application/json"). Do(context.Background())
Header("accept", "application/json")
_, err := create.Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to unfollow index: %s", esErrorString(err)) return fmt.Errorf("failed to unfollow index: %s", esErrorString(err))
@@ -136,11 +125,9 @@ func CcrFollowerAdd(conf *cfg.Config, index string) error {
return err return err
} }
create := conf.DefaultCluster.ES.Ccr.Follow(index). create := conf.DefaultCluster.ES().Ccr.Follow(index).
LeaderIndex(index). LeaderIndex(index).
RemoteCluster(remote). RemoteCluster(remote)
Header("content-type", "application/json").
Header("accept", "application/json")
if conf.Wait { if conf.Wait {
create.WaitForActiveShards("all") create.WaitForActiveShards("all")
@@ -156,9 +143,7 @@ func CcrFollowerAdd(conf *cfg.Config, index string) error {
} }
func CcrFollowerShow(conf *cfg.Config, index string) error { func CcrFollowerShow(conf *cfg.Config, index string) error {
res, err := conf.DefaultCluster.ES.Ccr.FollowStats(index). res, err := conf.DefaultCluster.ES().Ccr.FollowStats(index).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to retrieve follower index info: %s", esErrorString(err)) return fmt.Errorf("failed to retrieve follower index info: %s", esErrorString(err))

View File

@@ -20,7 +20,6 @@ import (
"context" "context"
"fmt" "fmt"
"log/slog" "log/slog"
"slices"
"strings" "strings"
"sync" "sync"
@@ -59,36 +58,37 @@ type apiResponse struct {
} }
func ClusterList(conf *cfg.Config) error { 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 idx := 0
names := make([]string, len(conf.Clusters)) for name, cluster := range conf.Clusters {
for name := range conf.Clusters { reachable := "no"
names[idx] = name current := "no"
idx++
}
slices.Sort(names) _, err := cluster.ES().Cluster.Health().
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").
Do(context.Background()) Do(context.Background())
if err == nil { 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 { if err := table.Print(); err != nil {
return err return err
} }
@@ -114,9 +114,9 @@ func ClusterStatus(conf *cfg.Config) error {
} }
for _, cluster := range clusters { for _, cluster := range clusters {
es := conf.DefaultCluster.ES es := conf.DefaultCluster.ES()
if cluster != "default" { if cluster != "default" {
es = conf.Clusters[cluster].ES es = conf.Clusters[cluster].ES()
} }
responses := make(chan apiResponse, gocount) responses := make(chan apiResponse, gocount)

View File

@@ -28,9 +28,7 @@ import (
) )
func ClusterSettingsList(conf *cfg.Config) error { func ClusterSettingsList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Cluster.GetSettings(). res, err := conf.DefaultCluster.ES().Cluster.GetSettings().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get cluster settings: %s", esErrorString(err)) 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 { 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() { for _, arg := range args.Slice() {
setting, value := splitArg(arg) setting, value := splitArg(arg)
@@ -99,11 +97,7 @@ func ClusterSettingsSet(conf *cfg.Config, args cli.Args) error {
} }
} }
_, err := put. _, err := put.Do(context.Background())
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to set settings: %s", esErrorString(err)) 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) 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). AddPersistent(setting, message).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {

View File

@@ -37,9 +37,7 @@ const (
) )
func checkClusterIsLeader(conf *cfg.Config, leader string) bool { func checkClusterIsLeader(conf *cfg.Config, leader string) bool {
stats, err := conf.Clusters[leader].ES.Ccr.Stats(). stats, err := conf.Clusters[leader].ES().Ccr.Stats().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
fmt.Printf("failed to get ccr stats from %s: %s", leader, esErrorString(err)) 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{} status := map[string]*health.Response{}
for _, cluster := range []string{leader, follower} { for _, cluster := range []string{leader, follower} {
st, err := conf.Clusters[leader].ES.Cluster.Health(). st, err := conf.Clusters[leader].ES().Cluster.Health().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
fmt.Printf("failed to get health from %s: %s", cluster, esErrorString(err)) 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{} failed := map[string]map[string]string{}
for _, cluster := range []string{leader, follower} { 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). OnlyManaged(true).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
fmt.Printf("failed to get ilm status from %s: %s", cluster, esErrorString(err)) 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] _, followerHasIt := indices[follower][name]
if !followerHasIt { if !followerHasIt {
// fetch index details // fetch index details
res, err := conf.Clusters[leader].ES.Indices.Get(name). res, err := conf.Clusters[leader].ES().Indices.Get(name).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
continue // ignore it then continue // ignore it then
@@ -255,9 +247,7 @@ func findOrphanedIndices(conf *cfg.Config, indices ClusterIndices, leader, follo
_, leaderHasIt := indices[leader][name] _, leaderHasIt := indices[leader][name]
if !leaderHasIt { if !leaderHasIt {
// fetch index details // fetch index details
res, err := conf.Clusters[follower].ES.Indices.Get(name). res, err := conf.Clusters[follower].ES().Indices.Get(name).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
continue // ignore it then continue // ignore it then
@@ -353,8 +343,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c
switch which { switch which {
case "health": case "health":
res, err := es.Cluster.Health(). res, err := es.Cluster.Health().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
ar.health = res ar.health = res
@@ -363,8 +351,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c
case "info": case "info":
res, err := es.Info(). res, err := es.Info().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
ar.info = res ar.info = res
@@ -373,8 +359,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c
case "ccrstats": case "ccrstats":
res, err := es.Ccr.Stats(). res, err := es.Ccr.Stats().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
ar.ccr = res ar.ccr = res
@@ -383,8 +367,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c
case "stats": case "stats":
res, err := es.Cluster.Stats(). res, err := es.Cluster.Stats().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
ar.stats = res ar.stats = res
@@ -393,8 +375,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c
case "indices": case "indices":
res, err := es.Cat.Indices(). res, err := es.Cat.Indices().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
ar.indices = &res ar.indices = &res
@@ -403,8 +383,6 @@ func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan c
case "tasks": case "tasks":
res, err := es.Cat.Tasks(). res, err := es.Cat.Tasks().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
ar.tasks = &res ar.tasks = &res

View File

@@ -30,9 +30,7 @@ import (
) )
func DatastreamNames(conf *cfg.Config) ([]string, error) { func DatastreamNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES.Indices.GetDataStream(). res, err := conf.DefaultCluster.ES().Indices.GetDataStream().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to get data streams: %s", esErrorString(err)) 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 { func DatastreamList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Indices.GetDataStream(). res, err := conf.DefaultCluster.ES().Indices.GetDataStream().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get data streams: %s", esErrorString(err)) 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 { func DatastreamShow(conf *cfg.Config, dsname string) error {
res, err := conf.DefaultCluster.ES.Indices.GetDataStream(). res, err := conf.DefaultCluster.ES().Indices.GetDataStream().
Name(dsname). Name(dsname).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get data stream: %s", esErrorString(err)) 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") 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). Name(dsname).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get data stream stats: %s", esErrorString(err)) 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 { func DatastreamCreate(conf *cfg.Config, dsname string) error {
_, err := conf.DefaultCluster.ES.Indices.CreateDataStream(dsname). _, err := conf.DefaultCluster.ES().Indices.CreateDataStream(dsname).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -226,9 +216,7 @@ func DatastreamCreate(conf *cfg.Config, dsname string) error {
} }
func DatastreamDelete(conf *cfg.Config, dsname string) error { func DatastreamDelete(conf *cfg.Config, dsname string) error {
_, err := conf.DefaultCluster.ES.Indices.DeleteDataStream(dsname). _, err := conf.DefaultCluster.ES().Indices.DeleteDataStream(dsname).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {

View File

@@ -47,10 +47,8 @@ func DocAdd(conf *cfg.Config, jsondoc string) error {
now := fmt.Sprintf("%d", rand.Int64()) 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). Document(data).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to create new doc in index %s: %s", conf.Index, esErrorString(err)) 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 { func DocShow(conf *cfg.Config, id string) error {
res, err := conf.DefaultCluster.ES.Get(conf.Index, id). res, err := conf.DefaultCluster.ES().Get(conf.Index, id).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to retrieve doc in index %s: %s", conf.Index, esErrorString(err)) 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=") { if len(queries) == 1 && strings.Contains(queries[0], "id=") {
id, _ := strings.CutPrefix(queries[0], "id=") id, _ := strings.CutPrefix(queries[0], "id=")
_, err := conf.DefaultCluster.ES.Delete(conf.Index, id). _, err := conf.DefaultCluster.ES().Delete(conf.Index, id).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to delete doc in index %s: %s", conf.Index, esErrorString(err)) 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 req.Query = queryCaster
} }
_, err := conf.DefaultCluster.ES.DeleteByQuery(conf.Index). _, err := conf.DefaultCluster.ES().DeleteByQuery(conf.Index).
Request(req). Request(req).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to delete docs in index %s: %s", conf.Index, esErrorString(err)) return fmt.Errorf("failed to delete docs in index %s: %s", conf.Index, esErrorString(err))

View File

@@ -32,7 +32,7 @@ import (
) )
func IlmRetry(conf *cfg.Config, index string) error { 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()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to retry ilm: %s", esErrorString(err)) 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 { func IlmStatus(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Ilm.GetStatus(). res, err := conf.DefaultCluster.ES().Ilm.GetStatus().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get ilm status: %s", esErrorString(err)) 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) { func IlmNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES.Ilm.GetLifecycle(). res, err := conf.DefaultCluster.ES().Ilm.GetLifecycle().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to get ilm policies: %s", esErrorString(err)) 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 { func IlmList(conf *cfg.Config, pattern string) error {
ilm := conf.DefaultCluster.ES.Ilm.GetLifecycle() ilm := conf.DefaultCluster.ES().Ilm.GetLifecycle()
if pattern != "" { if pattern != "" {
ilm.FilterPath(pattern) ilm.FilterPath(pattern)
@@ -103,7 +103,7 @@ func IlmList(conf *cfg.Config, pattern string) error {
} }
func IlmShow(conf *cfg.Config, policy 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). Policy(policy).
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -165,7 +165,7 @@ func ilmPhaseString(phase *types.Phase, short bool) string {
} }
func IlmExplain(conf *cfg.Config, index string) error { 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()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get ilm state: %s", esErrorString(err)) 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 { 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 cfg := conf.Ilm
var phases types.PhasesVariant = esdsl.NewPhases() var phases types.PhasesVariant = esdsl.NewPhases()

View File

@@ -35,9 +35,7 @@ import (
// used for completion // used for completion
func IndexNames(conf *cfg.Config) ([]string, error) { func IndexNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES.Cat.Indices(). res, err := conf.DefaultCluster.ES().Cat.Indices().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -86,10 +84,7 @@ func filterIndices(conf *cfg.Config, list indices.Response) indices.Response {
} }
func IndexList(conf *cfg.Config) error { func IndexList(conf *cfg.Config) error {
cat := conf.DefaultCluster.ES.Cat.Indices(). 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")
if conf.Failed { if conf.Failed {
cat = cat.Health(healthstatus.Red) cat = cat.Health(healthstatus.Red)
@@ -130,10 +125,7 @@ func IndexList(conf *cfg.Config) error {
} }
func IndexShow(conf *cfg.Config, indexpattern string) error { func IndexShow(conf *cfg.Config, indexpattern string) error {
res, err := conf.DefaultCluster.ES.Indices.Get(indexpattern). 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").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get index: %s", esErrorString(err)) 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 { func IndexCreate(conf *cfg.Config, index string, mappings []string) error {
settings := esdsl.NewIndexSettings() settings := esdsl.NewIndexSettings()
create := conf.DefaultCluster.ES.Indices.Create(index). create := conf.DefaultCluster.ES().Indices.Create(index)
Header("content-type", "application/json").
Header("accept", "application/json")
if conf.Wait { if conf.Wait {
create.WaitForActiveShards("all") 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 { func IndexDelete(conf *cfg.Config, index string) error {
_, err := conf.DefaultCluster.ES.Indices.Delete(index). _, err := conf.DefaultCluster.ES().Indices.Delete(index).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to delete index: %s", esErrorString(err)) 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 { func IndexClose(conf *cfg.Config, index string) error {
create := conf.DefaultCluster.ES.Indices.Close(index). _, err := conf.DefaultCluster.ES().Indices.Close(index).
Header("content-type", "application/json"). Do(context.Background())
Header("accept", "application/json")
_, err := create.Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to close index: %s", esErrorString(err)) 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 { func IndexAllocation(conf *cfg.Config, index string) error {
res, err := conf.DefaultCluster.ES.Cluster.AllocationExplain(). res, err := conf.DefaultCluster.ES().Cluster.AllocationExplain().
Index(index). Index(index).
Primary(conf.Primary). Primary(conf.Primary).
Shard(conf.Shards). Shard(conf.Shards).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get index allocation explain: %s", esErrorString(err)) 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 { func IndexModify(conf *cfg.Config, index string) error {
settings := esdsl.NewIndexSettings().NumberOfReplicas(strconv.Itoa(conf.Replicas)) settings := esdsl.NewIndexSettings().NumberOfReplicas(strconv.Itoa(conf.Replicas))
_, err := conf.DefaultCluster.ES.Indices.PutSettings(). _, err := conf.DefaultCluster.ES().Indices.PutSettings().
Indices(index). Indices(index).
Index(settings). Index(settings).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to modify index settings: %s", esErrorString(err)) 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 { func IndexFields(conf *cfg.Config, index string) error {
res, err := conf.DefaultCluster.ES.FieldCaps(). res, err := conf.DefaultCluster.ES().FieldCaps().
Index(index). Index(index).
Fields("*"). Fields("*").
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to retrieve field capabilties: %s", esErrorString(err)) return fmt.Errorf("failed to retrieve field capabilties: %s", esErrorString(err))

View File

@@ -29,9 +29,7 @@ import (
// FIXME: add filter support, see IndexCreate mapping // FIXME: add filter support, see IndexCreate mapping
func IndexAliasCreate(conf *cfg.Config, index, alias string) error { func IndexAliasCreate(conf *cfg.Config, index, alias string) error {
res, err := conf.DefaultCluster.ES.Indices.PutAlias(index, alias). res, err := conf.DefaultCluster.ES().Indices.PutAlias(index, alias).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
slog.Debug("create alias", "result", res) slog.Debug("create alias", "result", res)
@@ -51,10 +49,8 @@ func IndexAliasList(conf *cfg.Config) error {
filter = *regexp.MustCompile(conf.Filter[0]) filter = *regexp.MustCompile(conf.Filter[0])
} }
res, err := conf.DefaultCluster.ES.Indices.GetAlias(). res, err := conf.DefaultCluster.ES().Indices.GetAlias().
Index("_all"). Index("_all").
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -100,9 +96,7 @@ func IndexAliasList(conf *cfg.Config) error {
} }
func IndexAliasDelete(conf *cfg.Config, index, alias string) error { func IndexAliasDelete(conf *cfg.Config, index, alias string) error {
res, err := conf.DefaultCluster.ES.Indices.DeleteAlias(index, alias). res, err := conf.DefaultCluster.ES().Indices.DeleteAlias(index, alias).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
slog.Debug("delete alias", "result", res) slog.Debug("delete alias", "result", res)

View File

@@ -33,9 +33,7 @@ import (
// used for completion // used for completion
func IndexTemplateList(conf *cfg.Config) error { func IndexTemplateList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Indices.GetIndexTemplate(). res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -65,10 +63,8 @@ func IndexTemplateList(conf *cfg.Config) error {
} }
func IndexTemplateShow(conf *cfg.Config, tplname string) 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). Name(tplname).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -158,9 +154,7 @@ func IndexTemplateCreate(conf *cfg.Config, name string, mappings []string) error
settings := esdsl.NewIndexSettings() settings := esdsl.NewIndexSettings()
maps := esdsl.NewIndexTemplateMapping() maps := esdsl.NewIndexTemplateMapping()
create := conf.DefaultCluster.ES.Indices.PutIndexTemplate(name). create := conf.DefaultCluster.ES().Indices.PutIndexTemplate(name)
Header("content-type", "application/json").
Header("accept", "application/json")
if conf.Shards > 0 { if conf.Shards > 0 {
settings = settings.NumberOfShards(strconv.Itoa(conf.Shards)) settings = settings.NumberOfShards(strconv.Itoa(conf.Shards))
@@ -243,10 +237,8 @@ func IndexTemplateModify(conf *cfg.Config, name string, mappings []string) error
maps := esdsl.NewIndexTemplateMapping() maps := esdsl.NewIndexTemplateMapping()
// load existing index mapping // load existing index mapping
res, err := conf.DefaultCluster.ES.Indices.GetIndexTemplate(). res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
Name(name). Name(name).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -262,9 +254,7 @@ func IndexTemplateModify(conf *cfg.Config, name string, mappings []string) error
tpl := res.IndexTemplates[0] tpl := res.IndexTemplates[0]
// our modify PUT request // our modify PUT request
modify := conf.DefaultCluster.ES.Indices.PutIndexTemplate(name). modify := conf.DefaultCluster.ES().Indices.PutIndexTemplate(name)
Header("content-type", "application/json").
Header("accept", "application/json")
// load existing settings, if any // load existing settings, if any
settings := tpl.IndexTemplate.Template.Settings 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 { func IndexTemplateDelete(conf *cfg.Config, name string) error {
_, err := conf.DefaultCluster.ES.Indices.DeleteIndexTemplate(name). _, err := conf.DefaultCluster.ES().Indices.DeleteIndexTemplate(name).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to delete index template: %s", esErrorString(err)) 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 // find indices matching those, find their associated aliases and
// rollover all we find. Only run when conf.Rollover==true // rollover all we find. Only run when conf.Rollover==true
func rolloverAliasIndexTemplate(conf *cfg.Config, name string) error { func rolloverAliasIndexTemplate(conf *cfg.Config, name string) error {
res, err := conf.DefaultCluster.ES.Indices.GetIndexTemplate(). res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
Name(name). Name(name).
Do(context.Background()) Do(context.Background())
@@ -389,7 +377,7 @@ func rolloverAliasIndexTemplate(conf *cfg.Config, name string) error {
// find all aliases matching the patterns // find all aliases matching the patterns
for _, pattern := range patterns { for _, pattern := range patterns {
res, err := conf.DefaultCluster.ES.Indices.ResolveIndex(pattern). res, err := conf.DefaultCluster.ES().Indices.ResolveIndex(pattern).
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to resolve index pattern: %s", esErrorString(err)) return fmt.Errorf("failed to resolve index pattern: %s", esErrorString(err))

View File

@@ -27,7 +27,7 @@ import (
func NodeList(conf *cfg.Config) error { func NodeList(conf *cfg.Config) error {
// get nodes // 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 { if err != nil {
return fmt.Errorf("failed to get nodes: %s", esErrorString(err)) return fmt.Errorf("failed to get nodes: %s", esErrorString(err))
} }

View File

@@ -28,7 +28,7 @@ import (
) )
func RoleNames(conf *cfg.Config) ([]string, error) { func RoleNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES.Security.GetRole(). res, err := conf.DefaultCluster.ES().Security.GetRole().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to get roles: %s", esErrorString(err)) 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 { func RoleList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Security.GetRole(). res, err := conf.DefaultCluster.ES().Security.GetRole().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get roles: %s", esErrorString(err)) 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 { func RoleShow(conf *cfg.Config, rolename string) error {
res, err := conf.DefaultCluster.ES.Security.GetRole(). res, err := conf.DefaultCluster.ES().Security.GetRole().
Name(rolename). Name(rolename).
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {

View File

@@ -207,7 +207,7 @@ func RoleDiff(conf *cfg.Config, csvfile, role string) error {
return RoleDiffSingle(conf, csvfile, role) return RoleDiffSingle(conf, csvfile, role)
} }
res, err := conf.DefaultCluster.ES.Security.GetRole(). res, err := conf.DefaultCluster.ES().Security.GetRole().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get roles: %s", esErrorString(err)) 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) { 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()) Do(context.Background())
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to get role mappings: %s", esErrorString(err)) 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 { 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). Name(rolename).
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {

View File

@@ -47,7 +47,7 @@ func RolloverConditions(conf *cfg.Config) types.RolloverConditionsVariant {
} }
func RolloverAlias(conf *cfg.Config, alias string) (*rollover.Response, error) { 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 { if conf.Wait {
roll.WaitForActiveShards("all") roll.WaitForActiveShards("all")

View File

@@ -50,7 +50,7 @@ func Search(conf *cfg.Config, queries []string) error {
return validateSearch(conf, queries) 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) queryCaster, err := prepareQuery(conf, queries)
if err != nil { if err != nil {
@@ -121,7 +121,7 @@ func explain(res *types.ExplanationDetail, indent string) {
} }
func validateSearch(conf *cfg.Config, queries []string) error { 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) queryCaster, err := prepareQuery(conf, queries)
if err != nil { if err != nil {
@@ -150,7 +150,7 @@ func validateSearch(conf *cfg.Config, queries []string) error {
} }
func Debug(conf *cfg.Config) error { func Debug(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Search(). res, err := conf.DefaultCluster.ES().Search().
Index(conf.Index). Index(conf.Index).
Size(0). Size(0).
Aggregations(map[string]types.Aggregations{ 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 // 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 { func searchPit(conf *cfg.Config, req *search.Request) error {
ctx := context.Background() 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 { if err != nil {
return fmt.Errorf("failed to open point-in-time request for search: %s", err) return fmt.Errorf("failed to open point-in-time request for search: %s", err)
} }
defer func() { 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 { if err != nil {
log.Fatalf("failed to close PIT: %s", err) log.Fatalf("failed to close PIT: %s", err)
} }
}() }()
search := conf.DefaultCluster.ES.Search(). search := conf.DefaultCluster.ES().Search().
Request(req). Request(req).
Pit(esdsl.NewPointInTimeReference(). Pit(esdsl.NewPointInTimeReference().
Id(pit.Id). Id(pit.Id).

View File

@@ -46,7 +46,7 @@ type filter struct {
} }
// Build a new filter object. We use this to build our elastic query // Build a new filter object. We use this to build our elastic query
// out of it. We support differnt types of queries: // out of it. We support different types of queries:
// //
// - nop filter: no query at all, just return the first N documents. // - nop filter: no query at all, just return the first N documents.
// //
@@ -243,7 +243,7 @@ func addFilters(conf *cfg.Config) ([]types.QueryVariant, error) {
// check if the conf.SortBy field (default: @timestamp) is searchable // check if the conf.SortBy field (default: @timestamp) is searchable
// by using the field capabilities API. // by using the field capabilities API.
func addSort(conf *cfg.Config, search *search.Search) *search.Search { 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). Index(conf.Index).
Fields(conf.SortBy). Fields(conf.SortBy).
Do(context.Background()) Do(context.Background())

View File

@@ -83,9 +83,7 @@ func filterShards(conf *cfg.Config, shardlist shards.Response) shards.Response {
} }
func ShardList(conf *cfg.Config) error { func ShardList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Cat.Shards(). res, err := conf.DefaultCluster.ES().Cat.Shards().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get shards: %s", esErrorString(err)) 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 { func ShardShow(conf *cfg.Config, index string) error {
res, err := conf.DefaultCluster.ES.Cat.Shards().Index(index). res, err := conf.DefaultCluster.ES().Cat.Shards().Index(index).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get shards: %s", esErrorString(err)) return fmt.Errorf("failed to get shards: %s", esErrorString(err))

View File

@@ -43,7 +43,7 @@ type Snapshot struct {
func SnapshotList(conf *cfg.Config) error { func SnapshotList(conf *cfg.Config) error {
// get partial indicies // 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 { if err != nil {
return fmt.Errorf("failed to get indicies: %s", esErrorString(err)) return fmt.Errorf("failed to get indicies: %s", esErrorString(err))
} }
@@ -56,7 +56,7 @@ func SnapshotList(conf *cfg.Config) error {
} }
// get snapshots // 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 { if err != nil {
return fmt.Errorf("failed to get snapshots: %s", esErrorString(err)) 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 { 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 { if err != nil {
return fmt.Errorf("failed to get snapshot: %s", esErrorString(err)) return fmt.Errorf("failed to get snapshot: %s", esErrorString(err))
} }

View File

@@ -28,9 +28,7 @@ import (
) )
func TaskList(conf *cfg.Config) error { func TaskList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Cat.Tasks(). res, err := conf.DefaultCluster.ES().Cat.Tasks().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -65,10 +63,8 @@ func TaskList(conf *cfg.Config) error {
} }
func TaskCancel(conf *cfg.Config, taskid string) error { func TaskCancel(conf *cfg.Config, taskid string) error {
_, err := conf.DefaultCluster.ES.Tasks.Cancel(). _, err := conf.DefaultCluster.ES().Tasks.Cancel().
TaskId(taskid). TaskId(taskid).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {

View File

@@ -37,6 +37,7 @@ type Table struct {
Entries [][]string Entries [][]string
lenHeaders []int lenHeaders []int
alignInts bool
} }
func NewTable(conf *cfg.Config, columns, rows int) *Table { func NewTable(conf *cfg.Config, columns, rows int) *Table {
@@ -45,6 +46,7 @@ func NewTable(conf *cfg.Config, columns, rows int) *Table {
table.Headers = make([]string, columns) table.Headers = make([]string, columns)
table.Entries = make([][]string, rows) table.Entries = make([][]string, rows)
table.lenHeaders = make([]int, columns) table.lenHeaders = make([]int, columns)
table.alignInts = conf.AlignInts
return &table return &table
} }
@@ -137,7 +139,7 @@ func (data *Table) PrintTSV() error {
for idx, entry := range entries { for idx, entry := range entries {
length := visibleLen(entry) length := visibleLen(entry)
if isInt(entry) { 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 {