Compare commits

..

16 Commits
0.0.4 ... 0.0.8

Author SHA1 Message Date
T. von Dein
dd619ab815 Add CCR stuff, fix search, add docs support, support index maps, refactor (#13) 2026-05-12 18:28:57 +02:00
327bc84bd8 fix cluster list sorting 2026-05-11 14:03:30 +02:00
a785eb642f add settings reference to usage 2026-05-11 14:03:11 +02:00
a209debe09 fix linting issues 2026-05-11 10:46:17 +02:00
d1d97fed1b enhanced config debug print (dump w/o ES structs) 2026-05-11 10:42:58 +02:00
4194a004c0 separate cluster settings code 2026-05-11 10:28:25 +02:00
ceab7698e2 dont debug the whole confi, bump version 2026-05-11 10:23:44 +02:00
d22352da12 fix settings output, recursively determine the json path of settings 2026-05-11 10:22:17 +02:00
3bd26fde7a fix coredump, register default cluster name 2026-05-11 10:21:43 +02:00
T. von Dein
38ed9de0b4 Add CI pipelines, fix linting issues (#12) 2026-05-07 15:01:18 +02:00
T. von Dein
5e82a3706b enhance cluster status (#11) 2026-05-07 12:54:27 +02:00
0122009774 add es version to status output 2026-05-06 13:10:02 +02:00
8109481e23 check at bootstrap which cluster is reachable, fix current cluster 2026-05-06 13:02:28 +02:00
T. von Dein
af75bfa192 add -p -t -D flags to cluster settings ls (#10) 2026-05-06 07:25:27 +02:00
a71be5001a add missing utilities.go 2026-05-05 12:07:31 +02:00
T. von Dein
fcfeaf66a9 Add nodes and cluster settings support (#9) 2026-05-05 12:06:34 +02:00
29 changed files with 1736 additions and 533 deletions

42
.goreleaser.yaml Normal file
View File

@@ -0,0 +1,42 @@
# vim: set ts=2 sw=2 tw=0 fo=cnqoj
version: 2
before:
hooks:
- go mod tidy
gitea_urls:
api: https://codeberg.org/api/v1
download: https://codeberg.org
builds:
- env:
- CGO_ENABLED=0
goos:
- linux
- darwin
changelog:
sort: asc
filters:
exclude:
- "^docs:"
- "^test:"
groups:
- title: Improved
regexp: '^.*?(feat|add|new)(\([[:word:]]+\))??!?:.+$'
order: 0
- title: Fixed
regexp: '^.*?(bug|fix)(\([[:word:]]+\))??!?:.+$'
order: 1
- title: Changed
order: 999
release:
header: "# Release Notes"
footer: >-
---
Full Changelog: [{{ .PreviousTag }}...{{ .Tag }}](https://codeberg.org/scip/epuppy/compare/{{ .PreviousTag }}...{{ .Tag }})

26
.woodpecker/build.yaml Normal file
View File

@@ -0,0 +1,26 @@
matrix:
platform:
- linux/amd64
goversion:
- 1.25
labels:
platform: ${platform}
steps:
build:
when:
event: [push,manual]
image: golang:${goversion}
commands:
- go get
- go build
linter:
when:
event: [push,manual]
image: golang:${goversion}
commands:
- curl -sSfL https://raw.githubusercontent.com/golangci/golangci-lint/HEAD/install.sh | sh -s -- -b $(go env GOPATH)/bin v2.5.0
- golangci-lint --version
- golangci-lint run ./...

15
.woodpecker/release.yaml Normal file
View File

@@ -0,0 +1,15 @@
# build release
labels:
platform: linux/amd64
steps:
goreleaser:
image: goreleaser/goreleaser
when:
event: [tag,manual]
environment:
GITEA_TOKEN:
from_secret: DEPLOY_TOKEN
commands:
- goreleaser release --clean --verbose

18
CODE_OF_CONDUCT.md Normal file
View File

@@ -0,0 +1,18 @@
# CODE_OF_CONDUCT.md (compact model)
Purpose: Maintain a collaborative, harassment-free environment focused
on shipping quality software.
Standards: Be respectful; assume good intent; no harassment; no
discrimination; keep critiques technical.
Scope: All project spaces (issues, PRs, forums, events).
Reporting: Email tom AT vondein DOT org. Acknowledge within 72 hours.
Enforcement: Two maintainers review, one recuses on
conflict. Sanctions range from warning to removal. Summary posted (no
personal details).
Escalation: If you believe maintainers handled a report in bad faith,
escalate to our foundation committee (link) for independent review.

91
CONTRIBUTING.md Normal file
View File

@@ -0,0 +1,91 @@
## Project Goals
The idea behind this project is to build a small commandline tool to
manage an ElasticSearch cluster.
There will be no GUI, no web interface, no public API of some
sort. TUI subcommands maybe added.
The programming language used for this project will always be
[GOLANG](https://go.dev/) with the exception of the documentation
([Perl POD](https://perldoc.perl.org/perlpod)) and the Makefile.
# Contributing
You can contribute to this project in various ways:
## Open an issue
If you encounter a problem or don't understand how the program works
or if you think the documentation is unclear, please don't hesitate to
open an issue.
Please add as much information about the case as possible, such as:
- Your environment (operating system etc)
- program version
- Input data. Please replace sensitive information with mock data!
- Actual program output.
- Expected program output.
- Error message - if any.
Be aware that I am working on this (and some other) project in my
spare time which is scarce. Therefore please don't expect me to
respond to your query within hours or even days. Be patient, but I
WILL respond.
## Pull Requests
Code and documentation help is always much appreciated! Please follow
thes guidelines to successfully contribute:
- Every pull request shall be based on latest `development`
branch. `main` is only used for releases.
- Execute the unit tests before committing: `make test`. There shall
be no errors.
- Strive to be backwards compatible so that users who are already
using the program don't have to change their habits - unless it is
really neccessary.
- Try to add a unit test for your addition.
- Don't ever change existing unit tests!
- Add a meaningful and comprehensive rationale about your contribution:
- Why do you think it might be useful for others?
- What did you actually change or add?
- Is there an open issue which this PR fixes and if so, please link
to that issue.
- [Re-]format your code with `gofmt -s`.
- Avoid unneccesary dependencies, especially for very small functions.
- **If** a new dependency is being added, it must be compatible with
our [license agreement](LICENSE).
- You need to accept that the code or documentation you contribute
will be redistributed under the terms of said license agreement. If
your contribution is considerably large or if you contribute
regularly, then feel free to add your name and if you want your
email address to the *AUTHORS* section of the
manual page.
- Adhere to the above mentioned project goals.
- If you are unsure if your addition or change will be accepted,
better ask before starting coding. Open an issue about your proposal
and let's discuss it! That way we avoid doing unnessesary work on
both sides.
Each pull request will be carefully reviewed and if it is a useful
addition it will be accepted. However, please be prepared that
sometimes a PR will be rejected. The reasons may vary and will be
documented. Perhaps the above guidelines are not matched, or the
addition seems to be not so useful from my perspective, maybe there
are too much changes or there might be changes I don't even
understand.
But whatever happens: your contribution is always welcome!

View File

@@ -1,3 +1,5 @@
[![status-badge](https://ci.codeberg.org/api/badges/16999/status.svg)](https://ci.codeberg.org/repos/16999)
# esctl
Elasticsearch CLI
@@ -12,13 +14,14 @@ USAGE:
esctl [global options] [command [command options]]
VERSION:
v0.0.3
v0.0.4
COMMANDS:
search, / search within an index
index, i manage indicies
snapshot, snap manage snapshots
cluster, c manage cluster[s]
node, snap manage nodes
help, h Shows a list of commands or help for one command
GLOBAL OPTIONS:

View File

@@ -3,3 +3,5 @@
- Fix index names custom completion
- add cluster default <name> which would add a flag to the config, so that no -C is needed subsequently
- add `cluster stats` from `/_cluster/stats` like mem, procs, open files, num indices, shards etc
or add these to `cluster status`, maybe add a `--stats` to include stats there?

103
cmd/ccr.go Normal file
View File

@@ -0,0 +1,103 @@
/*
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 cmd
import (
"context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
"github.com/urfave/cli/v3"
)
func Ccr(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "ccr",
Aliases: []string{"replication", "rep"},
Usage: "manage cross cluster replication",
Commands: []*cli.Command{
CcrStatus(conf),
CcrShardPause(conf),
CcrShardResume(conf),
CcrFollower(conf),
},
}
}
func CcrStatus(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "status",
Aliases: []string{"st"},
Usage: "cross cluster replication status (yaml config with 2 clusters required)",
Flags: []cli.Flag{
&cli.StringFlag{
Name: "exclude",
Usage: "regexp of indicies to exclude",
Destination: &conf.Exclude,
Aliases: []string{"e"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
leader := cmd.Args().Get(0)
follower := cmd.Args().Get(1)
if leader == "" || follower == "" {
return errors.New("no leader and follower aliases specified")
}
_, hasLeader := conf.Clusters[leader]
_, hasFollower := conf.Clusters[follower]
if !hasLeader || !hasFollower {
return errors.New("either leader or follower alias not configured")
}
if err := es.CcrStatus(conf, leader, follower); err != nil {
return err
}
return nil
},
}
}
func CcrShardPause(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "pause",
Usage: "pause shard allocation",
Action: func(ctx context.Context, cmd *cli.Command) error {
return es.CcrShardPause(conf)
},
}
}
func CcrShardResume(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "resume",
Usage: "resume shard allocation",
Action: func(ctx context.Context, cmd *cli.Command) error {
return es.CcrShardResume(conf)
},
}
}

108
cmd/ccr_follower.go Normal file
View File

@@ -0,0 +1,108 @@
/*
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 cmd
import (
"context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
"github.com/urfave/cli/v3"
)
func CcrFollower(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "follower",
Aliases: []string{"f"},
Usage: "manage ccr follower indices",
Commands: []*cli.Command{
CcrFollowerShow(conf),
CcrFollowerAdd(conf),
CcrFollowerDelete(conf),
},
}
}
func CcrFollowerAdd(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "add",
Aliases: []string{"+"},
Usage: "add ccr follower index",
UsageText: "add [options] <index>",
Flags: []cli.Flag{
&cli.BoolFlag{
Name: "wait",
Usage: "wait for active shards",
Destination: &conf.Wait,
Aliases: []string{"w"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
args := cmd.Args()
if args.Len() != 1 {
return errors.New("missing arguments: <index>")
}
return es.CcrFollowerAdd(conf, cmd.Args().Get(0))
},
}
}
func CcrFollowerDelete(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "delete",
Aliases: []string{"-"},
Usage: "delete ccr follower index",
UsageText: "delete <index>",
Action: func(ctx context.Context, cmd *cli.Command) error {
args := cmd.Args()
if args.Len() != 1 {
return errors.New("missing arguments: <index>")
}
// just an alias to index delete
return es.IndexDelete(conf, cmd.Args().Get(0))
},
}
}
func CcrFollowerShow(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "show",
Aliases: []string{"sh"},
Usage: "show ccr follower index details",
UsageText: "show <index>",
Action: func(ctx context.Context, cmd *cli.Command) error {
args := cmd.Args()
if args.Len() != 1 {
return errors.New("missing arguments: <index>")
}
return es.CcrFollowerShow(conf, cmd.Args().Get(0))
},
}
}

View File

@@ -18,7 +18,6 @@ package cmd
import (
"context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
@@ -33,30 +32,26 @@ func Cluster(conf *cfg.Config) *cli.Command {
Usage: "manage cluster[s]",
Commands: []*cli.Command{
Compare(conf),
Status(conf),
List(conf),
ClusterStatus(conf),
ClusterList(conf),
ClusterSettings(conf),
},
}
}
func List(conf *cfg.Config) *cli.Command {
func ClusterList(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "list",
Usage: "list configured clusters",
Aliases: []string{"ls"},
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.List(conf); err != nil {
return err
}
return nil
return es.ClusterList(conf)
},
}
}
func Status(conf *cfg.Config) *cli.Command {
func ClusterStatus(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "status",
Usage: "show cluster status",
@@ -72,50 +67,7 @@ func Status(conf *cfg.Config) *cli.Command {
},
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.Status(conf); err != nil {
return err
}
return nil
},
}
}
func Compare(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "compare",
Aliases: []string{"c"},
Usage: "compare cluster[s] (yaml config with 2 clusters required)",
Flags: []cli.Flag{
&cli.StringFlag{
Name: "exclude",
Usage: "regexp of indicies to exclude",
Destination: &conf.Exclude,
Aliases: []string{"e"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
leader := cmd.Args().Get(0)
follower := cmd.Args().Get(1)
if leader == "" || follower == "" {
return errors.New("no leader and follower aliases specified")
}
_, hasLeader := conf.Clusters[leader]
_, hasFollower := conf.Clusters[follower]
if !hasLeader || !hasFollower {
return errors.New("either leader or follower alias not configured")
}
if err := es.ClusterCompare(conf, leader, follower); err != nil {
return err
}
return nil
return es.ClusterStatus(conf)
},
}
}

119
cmd/cluster_settings.go Normal file
View File

@@ -0,0 +1,119 @@
/*
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 cmd
import (
"context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
"github.com/urfave/cli/v3"
)
const (
SETTINGS = `https://www.elastic.co/docs/reference/elasticsearch/configuration-reference`
)
func ClusterSettings(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "settings",
Usage: "cluster settings management",
Aliases: []string{"config"},
Commands: []*cli.Command{
ClusterSettingsList(conf),
ClusterSettingsSet(conf),
},
}
}
func ClusterSettingsList(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "list",
Usage: "show cluster settings",
Aliases: []string{"ls", "get"},
Flags: []cli.Flag{
&cli.BoolFlag{
Name: "persistent",
Usage: "only show persistent setting[s] (default)",
Destination: &conf.Persistent,
Aliases: []string{"p"},
},
&cli.BoolFlag{
Name: "transient",
Usage: "only show transient setting[s]",
Destination: &conf.Transient,
Aliases: []string{"t"},
},
&cli.BoolFlag{
Name: "default",
Usage: "only show default setting[s]",
Destination: &conf.Default,
Aliases: []string{"D"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.ClusterSettingsList(conf); err != nil {
return err
}
return nil
},
}
}
func ClusterSettingsSet(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "set",
Usage: "set|update cluster settings",
Aliases: []string{"set", "update"},
UsageText: "set [options] setting:value [setting:value ...]\n\nReference: " + SETTINGS,
Flags: []cli.Flag{
&cli.BoolFlag{
Name: "persistent",
Usage: "add persistent setting[s] (default)",
Destination: &conf.Persistent,
Aliases: []string{"p"},
},
&cli.BoolFlag{
Name: "transient",
Usage: "add transient setting[s]",
Destination: &conf.Transient,
Aliases: []string{"t"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
args := cmd.Args()
if args.Len() == 0 {
return errors.New("at least one setting must be specified (format: setting:value)")
}
if err := es.ClusterSettingsSet(conf, cmd.Args()); err != nil {
return err
}
return nil
},
}
}

58
cmd/doc.go Normal file
View File

@@ -0,0 +1,58 @@
/*
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 cmd
import (
"context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
"github.com/urfave/cli/v3"
)
func Doc(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "doc",
Usage: "manage documents",
Commands: []*cli.Command{
DocAdd(conf),
//Delete(conf),
},
}
}
func DocAdd(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "add",
Aliases: []string{"+"},
Usage: "add JSON document index",
UsageText: "add [options] <index> '<json-doc>'",
Action: func(ctx context.Context, cmd *cli.Command) error {
args := cmd.Args()
if args.Len() != 2 {
return errors.New("missing arguments: <index> <json-doc>")
}
return es.DocAdd(conf, cmd.Args().Get(0), cmd.Args().Get(1))
},
}
}

View File

@@ -69,11 +69,7 @@ func IndexList(conf *cfg.Config) *cli.Command {
},
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.IndexList(conf); err != nil {
return err
}
return nil
return es.IndexList(conf)
},
}
}
@@ -85,11 +81,8 @@ func IndexShow(conf *cfg.Config) *cli.Command {
Usage: "show details about an index",
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.IndexShow(conf, cmd.Args().Get(0)); err != nil {
return err
}
return es.IndexShow(conf, cmd.Args().Get(0))
return nil
},
// FIXME: doesn't work at all
@@ -140,11 +133,14 @@ func IndexCreate(conf *cfg.Config) *cli.Command {
},
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.IndexCreate(conf, cmd.Args().Get(0)); err != nil {
return err
args := cmd.Args()
if args.Len() == 0 {
return fmt.Errorf("no index specified")
}
return nil
mappings := args.Slice()[1:]
return es.IndexCreate(conf, args.Get(0), mappings)
},
}
}
@@ -156,11 +152,7 @@ func IndexDelete(conf *cfg.Config) *cli.Command {
Usage: "delete an index",
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.IndexDelete(conf, cmd.Args().Get(0)); err != nil {
return err
}
return nil
return es.IndexDelete(conf, cmd.Args().Get(0))
},
}
}

65
cmd/node.go Normal file
View File

@@ -0,0 +1,65 @@
/*
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 cmd
import (
"context"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
"github.com/urfave/cli/v3"
)
func Node(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "node",
Aliases: []string{"snap"},
Usage: "manage nodes",
Commands: []*cli.Command{
NodeList(conf),
NodeShow(conf),
},
}
}
func NodeList(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "list",
Aliases: []string{"ls"},
Usage: "list nodes",
Action: func(ctx context.Context, cmd *cli.Command) error {
return es.NodeList(conf)
},
}
}
func NodeShow(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "show",
Aliases: []string{"sh"},
Usage: "show details about a node",
UsageText: "show [options] <node>",
Action: func(ctx context.Context, cmd *cli.Command) error {
// return es.NodeShow(conf, cmd.Args().Get(0))
return nil
},
}
}

View File

@@ -76,6 +76,9 @@ func Main() int {
Index(conf),
Snapshot(conf),
Cluster(conf),
Ccr(conf),
Node(conf),
Doc(conf),
},
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {

View File

@@ -18,6 +18,7 @@ package cmd
import (
"context"
"errors"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es"
@@ -30,7 +31,7 @@ func Search(conf *cfg.Config) *cli.Command {
Name: "search",
Aliases: []string{"/"},
Usage: "search within an index",
UsageText: "search [options] <query>",
UsageText: "search [options] <field=pattern> ...",
Flags: []cli.Flag{
&cli.StringFlag{
@@ -63,11 +64,12 @@ func Search(conf *cfg.Config) *cli.Command {
},
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.Search(conf, cmd.Args().Get(0)); err != nil {
return err
args := cmd.Args()
if args.Len() == 0 {
return errors.New("at least one query must be specified (format: field=pattern)")
}
return nil
return es.Search(conf, args.Slice())
},
}
}

View File

@@ -54,11 +54,7 @@ func SnapshotList(conf *cfg.Config) *cli.Command {
},
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.SnapshotList(conf); err != nil {
return err
}
return nil
return es.SnapshotList(conf)
},
}
}
@@ -71,11 +67,7 @@ func SnapshotShow(conf *cfg.Config) *cli.Command {
UsageText: "show [options] <snapshot>",
Action: func(ctx context.Context, cmd *cli.Command) error {
if err := es.SnapshotShow(conf, cmd.Args().Get(0)); err != nil {
return err
}
return nil
return es.SnapshotShow(conf, cmd.Args().Get(0))
},
}
}

View File

@@ -17,6 +17,7 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
package cfg
import (
"context"
"crypto/tls"
"errors"
"fmt"
@@ -30,7 +31,7 @@ import (
)
const (
Version string = `v0.0.4`
Version string = `v0.0.8`
)
type Cluster struct {
@@ -39,19 +40,20 @@ type Cluster struct {
}
type Config struct {
ConfigFile string // -c
CurrentCluster string // -C
Debug bool // -d
Clusters map[string]*Cluster
DefaultCluster *Cluster
Index string // index: -i
Failed, Partials bool // index: flags
Shards, Replicas int // index create: -s -r
Wait bool // index create: -w
From, To, MaxItems int // search: flags
Filter []string // search: -F
Exclude string // cluster compare: -e (regexp)
All bool // cluster status: -a
ConfigFile string // -c
CurrentCluster string // -C
Debug bool // -d
Clusters map[string]*Cluster
DefaultCluster *Cluster
Index string // index: -i
Failed, Partials bool // index: flags
Shards, Replicas int // index create: -s -r
Wait bool // index create: -w
From, To, MaxItems int // search: flags
Filter []string // search: -F
Exclude string // cluster compare: -e (regexp)
All bool // cluster status: -a
Persistent, Transient, Default bool // -p -t -D cluster settings set
}
func NewConfig() *Config {
@@ -75,6 +77,10 @@ func (conf *Config) Init() error {
}
}
if err := conf.SetupES(); err != nil {
return err
}
if conf.CurrentCluster != "" {
current, exists := conf.Clusters[conf.CurrentCluster]
if !exists {
@@ -82,17 +88,48 @@ func (conf *Config) Init() error {
} 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 {
conf.DefaultCluster = cluster
conf.CurrentCluster = name
}
}
}
}
if conf.Debug {
repr.Println(conf)
}
conf.SetupES()
conf.PrintDebug()
return nil
}
func (conf *Config) PrintDebug() {
if !conf.Debug {
return
}
clone := *conf
for name := range clone.Clusters {
clone.Clusters[name] = nil
}
clone.DefaultCluster = nil
fmt.Println("config:")
repr.Println(clone)
}
func (conf *Config) LoadEnv() error {
cluster := Cluster{
Uri: os.Getenv("ES_URI"),

87
pkg/es/ccr.go Normal file
View File

@@ -0,0 +1,87 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"errors"
"fmt"
"log/slog"
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
)
const (
shard_pause = `cluster.routing.allocation.enable`
)
func CcrShardPause(conf *cfg.Config) error {
return ClusterSettingsSetSingle(conf, shard_pause, "none")
}
func CcrShardResume(conf *cfg.Config) error {
return ClusterSettingsSetSingle(conf, shard_pause, "all")
}
func CcrStatus(conf *cfg.Config, leader, follower string) error {
if !checkClusterIsLeader(conf, leader) {
if !checkClusterIsLeader(conf, follower) {
return errors.New("leader/follower attribution is invalid, both clusters are followers")
}
// reverse attribution
f := follower
follower = leader
leader = f
slog.Debug("leader/follower attribution is invalid, reversing", "leader", leader, "follower", follower)
}
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())
if err != nil {
return fmt.Errorf("failed to get indicies on %s: %s", alias, err)
}
indices[alias] = make(map[string]*types.IndicesRecord, len(res))
for _, index := range res {
indices[alias][*index.Index] = &index
}
}
if !checkClusterStatus(conf, leader, follower) {
return errors.New("one of the two clusters is in a failed state")
}
findIlmErrors(conf, leader, follower)
if findIndicesOnlyOnLeader(conf, indices, leader, follower) &&
findOrphanedIndices(conf, indices, leader, follower) &&
findFailedFollowerIndices(conf, indices, follower) {
fmt.Println("everything's hunky-dory.")
}
return nil
}

106
pkg/es/ccr_follower.go Normal file
View File

@@ -0,0 +1,106 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"fmt"
"log/slog"
"codeberg.org/scip/esctl/pkg/cfg"
)
func CcrFollowerAdd(conf *cfg.Config, index string) error {
res, err := conf.DefaultCluster.ES.Cluster.RemoteInfo().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to retrieve follower info: %s", err)
}
remote := ""
for name := range res {
remote = name
break
}
if remote == "" {
return fmt.Errorf("cluster doesn't have a follower: %s", err)
}
create := conf.DefaultCluster.ES.Ccr.Follow(index).
LeaderIndex(index).
RemoteCluster(remote).
Header("content-type", "application/json").
Header("accept", "application/json")
if conf.Wait {
create.WaitForActiveShards("all")
}
_, err = create.Do(context.Background())
if err != nil {
return fmt.Errorf("failed to create follower index: %s", err)
}
return nil
}
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").
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to retrieve follower index info: %s", err)
}
slog.Debug("ES result", "follower stats", res.Indices)
if len(res.Indices) == 0 {
return fmt.Errorf("cluster did not return any follower stats for index %s", index)
}
if len(res.Indices[0].Shards) == 0 {
return fmt.Errorf("cluster did not return any shard stats on follower index %s", index)
}
follower := res.Indices[0].Shards[0]
table := NewTable(2, 9)
table.Addheaders("field", "value")
table.entries = [][]string{
{"name", index},
{"remote_cluster", follower.RemoteCluster},
{"leader_checkpoint", fmt.Sprintf("%d", follower.LeaderGlobalCheckpoint)},
{"follower_checkpoint", fmt.Sprintf("%d", follower.FollowerGlobalCheckpoint)},
{"bytes_read", fmt.Sprintf("%d", follower.BytesRead)},
{"failed_read_requests", fmt.Sprintf("%d", follower.FailedReadRequests)},
{"failed_write_requests", fmt.Sprintf("%d", follower.FailedWriteRequests)},
{"successful_read_requests", fmt.Sprintf("%d", follower.SuccessfulReadRequests)},
{"successful_write_requests", fmt.Sprintf("%d", follower.SuccessfulWriteRequests)},
}
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}

View File

@@ -18,84 +18,79 @@ package es
import (
"context"
"errors"
"fmt"
"log"
"log/slog"
"regexp"
"slices"
"sync"
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/elastic/go-elasticsearch/v9/typedapi/ccr/stats"
"github.com/elastic/go-elasticsearch/v9/typedapi/cluster/health"
"github.com/elastic/go-elasticsearch/v9/typedapi/core/info"
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
)
const (
DefaultExclude = `(part|monitoring|.internal|metrics-endpoint)`
ResponseHealth = iota
ResponseInfo
ResponseCcr
)
type ClusterIndices map[string]map[string]*types.IndicesRecord
func ClusterCompare(conf *cfg.Config, leader, follower string) error {
if !checkClusterFollower(conf, leader) {
return errors.New("leader/follower attribution is invalid, reverse cluster attribution and retry")
}
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())
if err != nil {
return fmt.Errorf("failed to get indicies on %s: %s", alias, err)
}
indices[alias] = make(map[string]*types.IndicesRecord, len(res))
for _, index := range res {
indices[alias][*index.Index] = &index
}
}
if !checkClusterStatus(conf, leader, follower) {
return errors.New("One of the two clusters is in a failed state")
}
findIlmErrors(conf, leader, follower)
if findIndicesOnlyOnLeader(conf, indices, leader, follower) &&
findOrphanedIndices(conf, indices, leader, follower) &&
findFailedFollowerIndices(conf, indices, follower) {
fmt.Println("everything's hunky-dory.")
}
return nil
type apiResponse struct {
error error
info *info.Response
health *health.Response
ccr *stats.Response
which int
}
func List(conf *cfg.Config) error {
table := NewTable(2, len(conf.Clusters))
table.Addheaders("cluster", "uri")
func ClusterList(conf *cfg.Config) error {
table := NewTable(3, len(conf.Clusters))
table.Addheaders("cluster", "uri", "default")
idx := 0
for name, cluster := range conf.Clusters {
table.entries[idx] = []string{name, cluster.Uri}
names := make([]string, len(conf.Clusters))
for name := range conf.Clusters {
names[idx] = name
idx++
}
table.Sort()
table.PrintMarkdown()
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").
Do(context.Background())
if err == nil {
name = Colorize("green", name)
}
table.entries[idx] = []string{name, cluster.Uri, fmt.Sprintf("%t", current)}
}
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}
func Status(conf *cfg.Config) error {
// We're using goroutines here to parallelize API requests, since we
// have to do 3 of'em for each cluster. This speeds things up.
func ClusterStatus(conf *cfg.Config) error {
clusters := []string{}
if conf.All {
for key, _ := range conf.Clusters {
for key := range conf.Clusters {
clusters = append(clusters, key)
}
} else {
@@ -108,367 +103,66 @@ func Status(conf *cfg.Config) error {
es = conf.Clusters[cluster].ES
}
res, err := es.Cluster.Health().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil {
log.Fatalf("Error getting health: %s", err)
responses := make(chan apiResponse, 3)
wg := &sync.WaitGroup{}
wg.Add(3)
go getClusterData(es, wg, responses, "health")
go getClusterData(es, wg, responses, "info")
go getClusterData(es, wg, responses, "ccrstats")
wg.Wait()
var clusterhealth *health.Response
var info *info.Response
var ccrstats *stats.Response
for i := 0; i < 3; i++ {
r := <-responses
if r.error != nil {
return r.error
}
switch r.which {
case ResponseHealth:
clusterhealth = r.health
case ResponseCcr:
ccrstats = r.ccr
case ResponseInfo:
info = r.info
}
}
slog.Debug("ES result", "cluster health", res)
slog.Debug("ES result", "cluster health", clusterhealth)
ccrfollowing := ""
if len(ccrstats.AutoFollowStats.AutoFollowedClusters) > 0 {
// is following another cluster
ccrfollowing = fmt.Sprintf("%s (%d/%d)",
ccrstats.AutoFollowStats.AutoFollowedClusters[0].ClusterName,
ccrstats.AutoFollowStats.NumberOfSuccessfulFollowIndices,
ccrstats.AutoFollowStats.NumberOfFailedFollowIndices,
)
}
table := NewTable(2, 5)
table.Addheaders(cluster, "status")
table.entries = [][]string{
{"Cluster Name", Colorize(*&res.Status.Name, res.ClusterName)},
{"Active Shards", fmt.Sprintf("%d", res.ActiveShards)},
{"Active Primary Shards", fmt.Sprintf("%d", res.ActivePrimaryShards)},
{"Indicies", fmt.Sprintf("%d", len(res.Indices))},
{"Nodes", fmt.Sprintf("%d", res.NumberOfNodes)},
{"Cluster Name", Colorize(clusterhealth.Status.Name, clusterhealth.ClusterName)},
{"ES Version", info.Version.Int},
{"Active Shards", fmt.Sprintf("%d", clusterhealth.ActiveShards)},
{"Active Primary Shards", fmt.Sprintf("%d", clusterhealth.ActivePrimaryShards)},
{"Indicies", fmt.Sprintf("%d", len(clusterhealth.Indices))},
{"Nodes", fmt.Sprintf("%d", clusterhealth.NumberOfNodes)},
{"AutoFollow (success/failed indices)", ccrfollowing},
}
table.PrintMarkdown()
if err := table.PrintMarkdown(); err != nil {
return err
}
}
return nil
}
// look for indicies only on leader
func findIndicesOnlyOnMaster(conf *cfg.Config, indices ClusterIndices, leader, follower string) {
exclude := regexp.MustCompile(DefaultExclude)
if conf.Exclude != "" {
exclude = regexp.MustCompile(conf.Exclude)
}
indexOnlyOnLeader := map[string]*types.IndicesRecord{}
for name, index := range indices[leader] {
if exclude.MatchString(name) {
continue
}
_, 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").
Do(context.Background())
if err != nil {
continue // ignore it then
}
_, defined := res[name]
if !defined {
// json response map didn't contain the index
continue
}
isWritable := false
for _, alias := range res[name].Aliases {
if *alias.IsWriteIndex {
isWritable = true
break
}
}
if isWritable {
// ignore index if associated alias index is writing
continue
}
indexOnlyOnLeader[name] = index
}
}
idx := 0
table := NewTable(3, len(indexOnlyOnLeader))
table.Addheaders("index only on leader", "size", "docscount")
for name, index := range indexOnlyOnLeader {
name := Colorize("red", name)
table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount}
idx++
}
table.Sort()
table.PrintMarkdown()
}
func checkClusterFollower(conf *cfg.Config, leader string) bool {
stats, err := conf.Clusters[leader].ES.Ccr.Stats().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil {
fmt.Printf("failed to get ccr stats from %s: %s", leader, err)
return false
}
if len(stats.AutoFollowStats.AutoFollowedClusters) == 0 {
// is not following anyone
return true
}
if stats.AutoFollowStats.AutoFollowedClusters[0].ClusterName != "" {
fmt.Println("leader/follower attribution is invalid, reverse cluster attribution and retry")
return false
}
return true
}
// checks if both clusters are green
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").
Do(context.Background())
if err != nil {
fmt.Printf("failed to get health from %s: %s", cluster, err)
return false
}
status[cluster] = st
}
table := NewTable(3, 5)
table.Addheaders("setting", "leader:"+leader, "follower:"+follower)
table.entries = [][]string{
{"Cluster Name",
Colorize(*&status[leader].Status.Name, status[leader].ClusterName),
Colorize(*&status[follower].Status.Name, status[follower].ClusterName),
},
{"Active Shards",
fmt.Sprintf("%d", status[leader].ActiveShards),
fmt.Sprintf("%d", status[follower].ActiveShards),
},
{"Active Primary Shards",
fmt.Sprintf("%d", status[leader].ActivePrimaryShards),
fmt.Sprintf("%d", status[follower].ActivePrimaryShards),
},
{"Indicies",
fmt.Sprintf("%d", len(status[leader].Indices)),
fmt.Sprintf("%d", len(status[follower].Indices)),
},
{"Nodes",
fmt.Sprintf("%d", status[leader].NumberOfNodes),
fmt.Sprintf("%d", status[follower].NumberOfNodes),
},
}
table.PrintMarkdown()
if status[leader].Status.Name == "green" && status[follower].Status.Name == "green" {
return true
}
return false
}
// finds indices on both clusters which have ilm errors
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").
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, err)
return false
}
failed[cluster] = map[string]string{}
for name, ilmstate := range ilm.Indices {
count := ilmstate.(*types.LifecycleExplainManaged).FailedStepRetryCount
if count != nil && *count > 0 {
failed[cluster][name] = fmt.Sprintf("%d", *count)
}
}
}
if len(failed[leader]) == 0 && len(failed[follower]) == 0 {
return true
}
for idx, cluster := range []string{leader, follower} {
which := "leader"
if idx > 0 {
which = "follower"
}
if len(failed[cluster]) > 0 {
idx := 0
table := NewTable(2, len(failed[cluster]))
table.Addheaders("ilm errors on "+which, "errors")
for name, count := range failed[cluster] {
table.entries[idx] = []string{name, count}
idx++
}
table.Sort()
table.PrintMarkdown()
}
}
return false
}
// find unsynchronized indicies only present on leader
func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool {
exclude := regexp.MustCompile(DefaultExclude)
if conf.Exclude != "" {
exclude = regexp.MustCompile(conf.Exclude)
}
indexOnlyOnLeader := map[string]*types.IndicesRecord{}
for name, index := range indices[leader] {
if exclude.MatchString(name) {
continue
}
_, 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").
Do(context.Background())
if err != nil {
continue // ignore it then
}
_, defined := res[name]
if !defined {
// json response map didn't contain the index
continue
}
isWritable := false
for _, alias := range res[name].Aliases {
if *alias.IsWriteIndex {
isWritable = true
break
}
}
if isWritable {
// ignore index if associated alias index is writing
continue
}
indexOnlyOnLeader[name] = index
}
}
if len(indexOnlyOnLeader) == 0 {
return true
}
idx := 0
table := NewTable(3, len(indexOnlyOnLeader))
table.Addheaders("index only on leader", "size", "docscount")
for name, index := range indexOnlyOnLeader {
name := Colorize("red", name)
table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount}
idx++
}
table.Sort()
table.PrintMarkdown()
return false
}
// find indices only present on follower
func findOrphanedIndices(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool {
orphaned := map[string]*types.IndicesRecord{}
for name, index := range indices[follower] {
_, 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").
Do(context.Background())
if err != nil {
continue // ignore it then
}
_, defined := res[name]
if !defined {
// json response map didn't contain the index
continue
}
orphaned[name] = index
}
}
if len(orphaned) == 0 {
return true
}
idx := 0
table := NewTable(3, len(orphaned))
table.Addheaders("orphaned index on follower", "size", "docscount")
for name, index := range orphaned {
name := Colorize("red", name)
table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount}
idx++
}
table.Sort()
table.PrintMarkdown()
return false
}
// find red indices on follower
func findFailedFollowerIndices(conf *cfg.Config, indices ClusterIndices, follower string) bool {
red := map[string]*types.IndicesRecord{}
for name, index := range indices[follower] {
if *index.Health == "red" {
red[name] = index
}
}
if len(red) == 0 {
return true
}
idx := 0
table := NewTable(3, len(red))
table.Addheaders("red index on follower", "size", "docscount")
for name, index := range red {
name := Colorize("red", name)
table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount}
idx++
}
table.Sort()
table.PrintMarkdown()
return false
}

160
pkg/es/cluster_settings.go Normal file
View File

@@ -0,0 +1,160 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/urfave/cli/v3"
)
// recursively traverse the raw settings hash and build a flat map
// consisting of the translated path and its value.
//
// e.g.
// logger:
//
// org:
// elasticsearch:
// transport:
// OutboundHandler: "ERROR"
//
// gets:
//
// logger.org.elasticsearch.transport.OutboundHandler: "ERROR"
func getJsonPath(raw map[string]any, topic string) map[string]string {
paths := map[string]string{}
for name, data := range raw {
path := topic + "." + name
switch value := data.(type) {
case string:
paths[path] = value
case map[string]any:
paths = getJsonPath(value, path)
}
}
return paths
}
func ClusterSettingsList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES.Cluster.GetSettings().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get cluster settings: %s", err)
}
table := NewTable(2, 0)
table.Addheaders("setting", "value")
entries := [][]string{}
settingshash := res.Persistent // == map[string]json.RawMessage
switch {
case conf.Transient:
settingshash = res.Transient
case conf.Default:
settingshash = res.Defaults
}
for topic, val := range settingshash {
data := map[string]any{}
err := json.Unmarshal(val, &data)
if err != nil {
return fmt.Errorf("failed to unmarshall setting for topic %s: %s", topic, err)
}
paths := getJsonPath(data, topic)
slog.Debug("settings", topic, paths)
for setting, value := range paths {
entries = append(entries, []string{setting, fmt.Sprintf("%v", value)})
}
}
table.entries = entries
table.Sort()
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}
func ClusterSettingsSet(conf *cfg.Config, args cli.Args) error {
put := conf.Clusters[conf.CurrentCluster].ES.Cluster.PutSettings()
for _, arg := range args.Slice() {
setting, value := splitArg(arg)
switch {
case conf.Transient:
message, err := json.Marshal(value)
if err != nil {
return fmt.Errorf("failed to marshall transient value <%v> to valid JSON: %s", value, err)
}
put.AddTransient(setting, message)
case conf.Persistent:
fallthrough
default:
message, err := json.Marshal(value)
if err != nil {
return fmt.Errorf("failed to marshall persistent value <%v> to valid JSON: %s", value, err)
}
put.AddPersistent(setting, message)
}
}
_, err := put.
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to set settings: %s", err)
}
return nil
}
func ClusterSettingsSetSingle(conf *cfg.Config, setting, value string) error {
message, err := json.Marshal(value)
if err != nil {
return fmt.Errorf("failed to marshall persistent value <%v> to valid JSON: %s", value, err)
}
_, err = conf.Clusters[conf.CurrentCluster].ES.Cluster.PutSettings().
AddPersistent(setting, message).
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to set %s: %s", setting, err)
}
return nil
}

385
pkg/es/cluster_util.go Normal file
View File

@@ -0,0 +1,385 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"errors"
"fmt"
"regexp"
"strings"
"sync"
"codeberg.org/scip/esctl/pkg/cfg"
"github.com/elastic/go-elasticsearch/v9"
"github.com/elastic/go-elasticsearch/v9/typedapi/cluster/health"
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
)
const (
DefaultExclude = `(part|monitoring|.internal|metrics-endpoint)`
)
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").
Do(context.Background())
if err != nil {
fmt.Printf("failed to get ccr stats from %s: %s", leader, err)
return false
}
if len(stats.AutoFollowStats.AutoFollowedClusters) == 0 {
// is not following anyone
return true
}
if stats.AutoFollowStats.AutoFollowedClusters[0].ClusterName != "" {
fmt.Println("leader/follower attribution is invalid, reverse cluster attribution and retry")
return false
}
return true
}
// checks if both clusters are green
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").
Do(context.Background())
if err != nil {
fmt.Printf("failed to get health from %s: %s", cluster, err)
return false
}
status[cluster] = st
}
table := NewTable(3, 5)
table.Addheaders("setting", "leader:"+leader, "follower:"+follower)
table.entries = [][]string{
{"Cluster Name",
Colorize(status[leader].Status.Name, status[leader].ClusterName),
Colorize(status[follower].Status.Name, status[follower].ClusterName),
},
{"Active Shards",
fmt.Sprintf("%d", status[leader].ActiveShards),
fmt.Sprintf("%d", status[follower].ActiveShards),
},
{"Active Primary Shards",
fmt.Sprintf("%d", status[leader].ActivePrimaryShards),
fmt.Sprintf("%d", status[follower].ActivePrimaryShards),
},
{"Indicies",
fmt.Sprintf("%d", len(status[leader].Indices)),
fmt.Sprintf("%d", len(status[follower].Indices)),
},
{"Nodes",
fmt.Sprintf("%d", status[leader].NumberOfNodes),
fmt.Sprintf("%d", status[follower].NumberOfNodes),
},
}
if err := table.PrintMarkdown(); err != nil {
fmt.Println(err)
return false
}
if status[leader].Status.Name == "green" && status[follower].Status.Name == "green" {
return true
}
return false
}
// finds indices on both clusters which have ilm errors
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").
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, err)
return false
}
failed[cluster] = map[string]string{}
for name, ilmstate := range ilm.Indices {
count := ilmstate.(*types.LifecycleExplainManaged).FailedStepRetryCount
if count != nil && *count > 0 {
failed[cluster][name] = fmt.Sprintf("%d", *count)
}
}
}
if len(failed[leader]) == 0 && len(failed[follower]) == 0 {
return true
}
for idx, cluster := range []string{leader, follower} {
which := "leader"
if idx > 0 {
which = "follower"
}
if len(failed[cluster]) > 0 {
idx := 0
table := NewTable(2, len(failed[cluster]))
table.Addheaders("ilm errors on "+which, "errors")
for name, count := range failed[cluster] {
table.entries[idx] = []string{name, count}
idx++
}
table.Sort()
if err := table.PrintMarkdown(); err != nil {
fmt.Println(err)
return false
}
}
}
return false
}
// find unsynchronized indicies only present on leader
func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool {
exclude := regexp.MustCompile(DefaultExclude)
if conf.Exclude != "" {
exclude = regexp.MustCompile(conf.Exclude)
}
indexOnlyOnLeader := map[string]*types.IndicesRecord{}
for name, index := range indices[leader] {
if exclude.MatchString(name) {
continue
}
_, 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").
Do(context.Background())
if err != nil {
continue // ignore it then
}
_, defined := res[name]
if !defined {
// json response map didn't contain the index
continue
}
isWritable := false
for _, alias := range res[name].Aliases {
if *alias.IsWriteIndex {
isWritable = true
break
}
}
if isWritable {
// ignore index if associated alias index is writing
continue
}
indexOnlyOnLeader[name] = index
}
}
if len(indexOnlyOnLeader) == 0 {
return true
}
idx := 0
table := NewTable(3, len(indexOnlyOnLeader))
table.Addheaders("index only on leader", "size", "docscount")
for name, index := range indexOnlyOnLeader {
name := Colorize("red", name)
table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount}
idx++
}
table.Sort()
if err := table.PrintMarkdown(); err != nil {
fmt.Println(err)
return false
}
return false
}
// find indices only present on follower
func findOrphanedIndices(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool {
orphaned := map[string]*types.IndicesRecord{}
for name, index := range indices[follower] {
_, 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").
Do(context.Background())
if err != nil {
continue // ignore it then
}
_, defined := res[name]
if !defined {
// json response map didn't contain the index
continue
}
orphaned[name] = index
}
}
if len(orphaned) == 0 {
return true
}
idx := 0
table := NewTable(3, len(orphaned))
table.Addheaders("orphaned index on follower", "size", "docscount")
for name, index := range orphaned {
name := Colorize("red", name)
table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount}
idx++
}
table.Sort()
if err := table.PrintMarkdown(); err != nil {
fmt.Println(err)
return false
}
return false
}
// find red indices on follower
func findFailedFollowerIndices(conf *cfg.Config, indices ClusterIndices, follower string) bool {
red := map[string]*types.IndicesRecord{}
for name, index := range indices[follower] {
if *index.Health == "red" {
red[name] = index
}
}
if len(red) == 0 {
return true
}
idx := 0
table := NewTable(3, len(red))
table.Addheaders("red index on follower", "size", "docscount")
for name, index := range red {
name := Colorize("red", name)
table.entries[idx] = []string{name, *index.DatasetSize, *index.DocsCount}
idx++
}
table.Sort()
if err := table.PrintMarkdown(); err != nil {
fmt.Println(err)
return false
}
return false
}
func splitArg(arg string) (string, string) {
parts := strings.Split(arg, ":")
switch len(parts) {
case 0:
fallthrough
case 1:
return arg, ""
default:
return parts[0], parts[1]
}
}
func getClusterData(es *elasticsearch.TypedClient, wg *sync.WaitGroup, reschan chan apiResponse, which string) {
defer wg.Done()
ar := apiResponse{}
arerr := errors.New("")
switch which {
case "health":
res, err := es.Cluster.Health().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
ar.health = res
ar.which = ResponseHealth
arerr = err
case "info":
res, err := es.Info().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
ar.info = res
ar.which = ResponseInfo
arerr = err
case "ccrstats":
res, err := es.Ccr.Stats().
Header("content-type", "application/json").
Header("accept", "application/json").
Do(context.Background())
ar.ccr = res
ar.which = ResponseCcr
arerr = err
}
if arerr != nil {
ar.error = fmt.Errorf("failed to get cluster health: %s", arerr)
}
reschan <- ar
}

55
pkg/es/doc.go Normal file
View File

@@ -0,0 +1,55 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"encoding/json"
"fmt"
"time"
"codeberg.org/scip/esctl/pkg/cfg"
)
// create index with:
//
// esctl index create foo [id:keyword user:text age:integer]
//
// then add a doc:
//
// esctl doc add foo2 '{"id":"d8d8d","user":"scip"}'
func DocAdd(conf *cfg.Config, index, jsondoc string) error {
data := map[string]any{}
err := json.Unmarshal([]byte(jsondoc), &data)
if err != nil {
return fmt.Errorf("supplied document was not valid JSON: %s", err)
}
now := fmt.Sprintf("%d", time.Now().Unix())
_, err = conf.DefaultCluster.ES.Create(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", index, err)
}
return nil
}

View File

@@ -21,6 +21,7 @@ import (
"fmt"
"log/slog"
"strconv"
"strings"
"time"
"codeberg.org/scip/esctl/pkg/cfg"
@@ -28,12 +29,6 @@ import (
"github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/healthstatus"
)
type Settings struct {
wait_for_active_shards string
number_of_shards int
number_of_replicas int
}
// used for completion
func IndexNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES.Cat.Indices().
@@ -42,7 +37,7 @@ func IndexNames(conf *cfg.Config) ([]string, error) {
Do(context.Background())
if err != nil {
return nil, fmt.Errorf("Error getting indicies: %s", err)
return nil, fmt.Errorf("failed to get indicies: %s", err)
}
indices := make([]string, len(res))
@@ -65,7 +60,7 @@ func IndexList(conf *cfg.Config) error {
res, err := cat.Do(context.Background())
if err != nil {
return fmt.Errorf("Error getting indicies: %s", err)
return fmt.Errorf("failed to get indicies: %s", err)
}
slog.Debug("ES result", "indicies", res)
@@ -92,7 +87,9 @@ func IndexList(conf *cfg.Config) error {
}
table.Sort()
table.PrintMarkdown()
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}
@@ -131,16 +128,14 @@ func IndexShow(conf *cfg.Config, index string) error {
{"uuid", *res[index].Settings.Index.Uuid},
}
table.PrintMarkdown()
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}
func IndexCreate(conf *cfg.Config, index string) error {
if index == "" {
return fmt.Errorf("no index specified")
}
func IndexCreate(conf *cfg.Config, index string, mappings []string) error {
settings := esdsl.NewIndexSettings()
create := conf.DefaultCluster.ES.Indices.Create(index).
@@ -158,6 +153,30 @@ func IndexCreate(conf *cfg.Config, index string) error {
settings = settings.NumberOfReplicas(strconv.Itoa(conf.Replicas))
}
if len(mappings) > 0 {
maps := esdsl.NewTypeMapping()
for _, mapping := range mappings {
parts := strings.Split(mapping, ":")
if len(parts) != 2 {
return fmt.Errorf("invalid mapping %s, expect <name:type> (type: integer, text, date, keyword)", mapping)
}
switch parts[1] {
case "text":
maps.AddProperty(parts[0], esdsl.NewTextProperty())
case "integer":
maps.AddProperty(parts[0], esdsl.NewIntegerNumberProperty())
case "date":
maps.AddProperty(parts[0], esdsl.NewDateProperty())
case "keyword":
maps.AddProperty(parts[0], esdsl.NewKeywordProperty())
}
}
create.Mappings(maps)
}
_, err := create.Settings(settings).
Do(context.Background())

57
pkg/es/node.go Normal file
View File

@@ -0,0 +1,57 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"fmt"
"log/slog"
"codeberg.org/scip/esctl/pkg/cfg"
)
func NodeList(conf *cfg.Config) error {
// get nodes
nodes, err := conf.DefaultCluster.ES.Cat.Nodes().Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get nodes: %s", err)
}
slog.Debug("ES result", "nodes", nodes)
table := NewTable(7, len(nodes))
table.Addheaders("name", "ip", "load1m", "load5m", "load15m", "ram %", "heap %")
for idx, node := range nodes {
table.entries[idx] = []string{
*node.Name,
*node.Ip,
*node.Load1M,
*node.Load5M,
*node.Load15M,
node.RamPercent.(string),
node.HeapPercent.(string),
}
}
table.Sort()
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}

View File

@@ -34,9 +34,17 @@ Execute an ES search.
q is the actual search query given as arg to the 'search' cmd
additional filters can be given as -F key=value
*/
func Search(conf *cfg.Config, q string) error {
query := esdsl.NewBoolQuery().
Must(esdsl.NewMatchQuery("message", q))
func Search(conf *cfg.Config, queries []string) error {
query := esdsl.NewBoolQuery()
for _, q := range queries {
parts := strings.Split(q, "=")
if len(parts) != 2 {
return fmt.Errorf("search queries must be in the form field=pattern")
}
query.Must(esdsl.NewMatchQuery(parts[0], parts[1]))
}
filters := make([]types.QueryVariant, len(conf.Filter))
@@ -62,7 +70,7 @@ func Search(conf *cfg.Config, q string) error {
}).
Do(context.Background())
if err != nil {
return fmt.Errorf("Error running search (esdsl): %s", err)
return fmt.Errorf("failed to run search (esdsl): %s", err)
}
slog.Debug("ES result", "search", res)

View File

@@ -44,7 +44,7 @@ func SnapshotList(conf *cfg.Config) error {
// get partial indicies
ires, err := conf.DefaultCluster.ES.Cat.Indices().Do(context.Background())
if err != nil {
return fmt.Errorf("Error getting indicies: %s", err)
return fmt.Errorf("failed to get indicies: %s", err)
}
indicies := map[string]int{}
@@ -57,7 +57,7 @@ func SnapshotList(conf *cfg.Config) error {
// get snapshots
sres, err := conf.DefaultCluster.ES.Cat.Snapshots().Do(context.Background())
if err != nil {
return fmt.Errorf("Error getting snapshots: %s", err)
return fmt.Errorf("failed to get snapshots: %s", err)
}
slog.Debug("ES result", "indicies", sres)
@@ -97,7 +97,9 @@ func SnapshotList(conf *cfg.Config) error {
}
table.Sort()
table.PrintMarkdown()
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}
@@ -135,7 +137,9 @@ func SnapshotShow(conf *cfg.Config, snapshot string) error {
{"shards-successful", fmt.Sprintf("%d", snap.Shards.Successful)},
}
table.PrintMarkdown()
if err := table.PrintMarkdown(); err != nil {
return err
}
return nil
}

View File

@@ -84,11 +84,11 @@ func (data *Table) PrintMarkdown() error {
table.Header(data.headers)
if err := table.Bulk(data.entries); err != nil {
return fmt.Errorf("Failed to add data to table renderer: %s", err)
return fmt.Errorf("failed to add data to table renderer: %s", err)
}
if err := table.Render(); err != nil {
fmt.Errorf("Failed to render table: %s", err)
return fmt.Errorf("failed to render table: %s", err)
}
fmt.Println(tableString.String())