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
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:
## 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)
[![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
@@ -26,13 +30,19 @@ Features:
- Cross cluster replication (ccr): view, pause, resume, delete
replication. You can also manage follower configuration.
- 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.
- Shard management: only list shards 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`
subcommand, which is for internal use. It can be used to verify if
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
elasticsearch API. You can run API calls on the current selected
cluster w/o the hassle to specify the whole url, credentials etc. It
@@ -51,87 +61,92 @@ Features:
Command tree:
```console
api
list
repl
show
ccr
follower
add
delete
pause
renew
resume
show
unfollow
info
pause
resume
status
cluster
list
settings
list
set
status
datastream
create
delete
list
rollover
show
debug
doc
add
delete
show
help
help-jsonpath
ilm
create
list
retry
show
status
index
alias
create
delete
list
rollover
allocation
close
create
delete
fields
ilm
list
modify
show
template
create
delete
list
modify
show
node
list
show
role
diff
list
show
search
shard
list
show
snapshot
list
show
task
cancel
list
version
api - api access and documentation
list - list index of API calls
show - show an API doc
repl - interactive API repl
ccr - manage cross cluster replication
status - cross cluster replication status (yaml config with 2 clusters required)
pause - pause shard allocation
resume - resume shard allocation
follower - manage ccr follower indices
show - show ccr follower index details
add - add ccr follower index
delete - delete ccr follower index
unfollow - unfollow ccr follower index
pause - pause ccr index to follow
resume - resume ccr index to follow
renew - renew ccr follower index
info - show ccr remote info
cluster - manage cluster[s]
status - show cluster status
switch - set current elasticsearch cluster
list - list configured clusters
settings - cluster settings management
list - show cluster settings
set - set|update cluster settings
datastream - manage data streams
list - list indicies
show - show details about an data stream
create - create a new data stream
delete - delete a data stream
rollover - roll over a data stream
doc - manage documents
add - add JSON document index
show - show a JSON document
delete - delete JSON document[s] from index[es]
ilm - manage index lifecycle
retry - retry applying an ILM profile to an index
status - get the current index lifecycle management status
list - list index lifecycle policies
show - show details about an index lifecycle policy
create - create a index lifecycle policy
index - manage indicies
list - list indicies
show - show details about an index
create - create a new index
delete - delete an index
close - close an index
allocation - explain index allocation
modify - modify an index
fields - show info about field capabilities
ilm - show ilm status
alias - manage index aliases
create - create an index alias
list - list index aliases
delete - delete an index alias
rollover - roll over an index alias
template - manage index templates
list - list index templates
show - show details about an index template
create - create a new index template
modify - modify a new index template
delete - delete an index template
node - manage nodes
list - list nodes
show - show details about a node
role - manage roles
list - list roles
show - show details about a role
diff - show differences between roles and CSV baseline
search - search within an index
shard - manage shards
list - list shards
show - show details about a shard
snapshot - manage snapshots
list - list snapshots
show - show details about a snapshot
task - manage tasks
list - list tasks
cancel - cancel running task
version - show esctl version information
debug - developer only
help-jsonpath - show jsonpath help
completion - Output shell completion script for bash, zsh, fish, or Powershell
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:
@@ -144,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
@@ -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
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

View File

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

View File

@@ -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 <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"
"os"
"runtime/pprof"
"strings"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
@@ -41,11 +42,11 @@ func Finish(err error) int {
func Main() int {
conf := cfg.NewConfig()
tree := false
cmd := &cli.Command{
Name: "esctl",
Usage: "manage elasticsearch from cli",
//Version: cfg.Version,
EnableShellCompletion: true,
Flags: []cli.Flag{
@@ -63,6 +64,20 @@ func Main() int {
Usage: "enable HTTP debugging",
Destination: &conf.DebugHTTP,
},
&cli.BoolFlag{
Name: "show-command-tree",
Value: false,
Usage: "generate a command tree",
Destination: &tree,
Hidden: true,
},
&cli.BoolFlag{
Name: "align-ints",
Aliases: []string{"I"},
Value: false,
Usage: "right align integers in tabular output",
Destination: &conf.AlignInts,
},
&cli.StringFlag{
Name: "config",
Aliases: []string{"c"},
@@ -94,25 +109,33 @@ func Main() int {
},
Commands: []*cli.Command{
Search(conf),
Cluster(conf),
Api(conf),
Ccr(conf),
Index(conf),
Ilm(conf),
Cluster(conf),
Datastream(conf),
Doc(conf),
Ilm(conf),
Index(conf),
Node(conf),
Roles(conf),
Search(conf),
Shard(conf),
Snapshot(conf),
Node(conf),
Doc(conf),
Task(conf),
Version(conf),
Debug(conf),
Roles(conf),
Task(conf),
HelpJsonPath(conf),
Api(conf),
},
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {
if tree {
if err := Tree(cmd); err != nil {
return nil, err
}
os.Exit(0)
}
if err := conf.Init(); err != nil {
if len(os.Args) > 1 {
return nil, err
@@ -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/tidwall/gjson v1.19.0
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
)

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/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.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/go.mod h1:RbqR21r5mrJuqunuUZ/Dhy/avygyECGrLceyNeo4LiM=
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
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
@@ -57,6 +45,7 @@ type Config struct {
DefaultCluster *Cluster
HaveJQ bool // determined at runtime by ourselfes
ProfileFile string // for internal use (golang profiling)
AlignInts bool // -I
Index string // index: -i
Failed, Partials bool // index: flags
@@ -119,8 +108,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 +134,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
}
} else {
if len(conf.Clusters) == 1 {
for name, cluster := range conf.Clusters {
conf.DefaultCluster = cluster
conf.CurrentCluster = name
}
} else {
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 {
// 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 {
if cluster.Default {
conf.DefaultCluster = cluster
conf.CurrentCluster = name
break
}
}
}
@@ -250,86 +249,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 {
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)

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{}
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))

View File

@@ -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))

View File

@@ -20,7 +20,6 @@ import (
"context"
"fmt"
"log/slog"
"slices"
"strings"
"sync"
@@ -59,35 +58,36 @@ 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)

View File

@@ -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 {

View File

@@ -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

View File

@@ -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 {

View File

@@ -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))

View File

@@ -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()

View File

@@ -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))

View File

@@ -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)

View File

@@ -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))

View File

@@ -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))
}

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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")

View File

@@ -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).

View File

@@ -46,7 +46,7 @@ type filter struct {
}
// 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.
//
@@ -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())

View File

@@ -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))

View File

@@ -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))
}

View File

@@ -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 {

View File

@@ -37,6 +37,7 @@ type Table struct {
Entries [][]string
lenHeaders []int
alignInts bool
}
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.Entries = make([][]string, rows)
table.lenHeaders = make([]int, columns)
table.alignInts = conf.AlignInts
return &table
}
@@ -137,7 +139,7 @@ func (data *Table) PrintTSV() error {
for idx, entry := range entries {
length := visibleLen(entry)
if isInt(entry) {
if isInt(entry) && data.alignInts {
// align right
fmt.Print(strings.Repeat(" ", data.lenHeaders[idx]-length), entry)
} else {