diff --git a/Makefile b/Makefile index 746a571..69b273d 100644 --- a/Makefile +++ b/Makefile @@ -110,3 +110,12 @@ profile: buildlocal ./esctl api ls --profile-file cpu.profile go tool pprof -text esctl cpu.profile go tool pprof --http localhost:8888 ./esctl cpu.profile + +docker-up: + make -C t up + +docker-waitup: + make -C t up + +docker-down: + make -C t down diff --git a/README.md b/README.md index 9a823b6..ac71275 100644 --- a/README.md +++ b/README.md @@ -216,7 +216,7 @@ Pending Tasks 0 Nodes 3 Red Indices 0 Long Running Tasks 2 -Indicies 221 +indices 221 Docs 9854777 Total Size 3.4 GB Total Queries 4210411 @@ -534,7 +534,7 @@ cluster - manage cluster[s] allocate-empty-primary - allocate an empty primary shard to a node allocate-stale-primary - allocate a stale primary shard to a node datastream - manage data streams - list - list indicies + list - list indices show - show details about an data stream create - create a new data stream delete - delete a data stream @@ -554,8 +554,8 @@ ilm - manage index lifecycle list - list index rollover config show - show rollover forecast over all indices explain - explain ilm condition of an index -index - manage indicies - list - list indicies +index - manage indices + list - list indices show - show details about an index create - create a new index update - update an index @@ -564,6 +564,7 @@ index - manage indicies fields - show info about field capabilities ilm - show ilm status du - show index disk usage + copy - copy (reindex) documents from one index to another alias - manage index aliases create - create an index alias list - list index aliases @@ -587,6 +588,7 @@ role - manage roles show - show details about a role diff - show differences between roles and CSV baseline search - search within an index +searchql - search using ES/QL language shard - manage shards list - list shards show - show details about a shard diff --git a/cmd/ccr.go b/cmd/ccr.go index 8a28c87..16a19ee 100644 --- a/cmd/ccr.go +++ b/cmd/ccr.go @@ -52,7 +52,7 @@ func CcrStatus(conf *cfg.Config) *cli.Command { Flags: []cli.Flag{ &cli.StringFlag{ Name: "exclude", - Usage: "regexp of indicies to exclude", + Usage: "regexp of indices to exclude", Destination: &conf.Exclude, Aliases: []string{"e"}, }, diff --git a/cmd/datastream.go b/cmd/datastream.go index 55be913..0b9bb94 100644 --- a/cmd/datastream.go +++ b/cmd/datastream.go @@ -46,7 +46,7 @@ func DatastreamList(conf *cfg.Config) *cli.Command { return &cli.Command{ Name: "list", Aliases: []string{"ls"}, - Usage: "list indicies", + Usage: "list indices", Flags: []cli.Flag{ &cli.IntFlag{ @@ -169,7 +169,7 @@ func DatastreamRollover(conf *cfg.Config) *cli.Command { Usage: "roll over after max age (eg: 7d, 2m, 8h)", Destination: &conf.MaxAge, }, - &cli.IntFlag{ + &cli.Int64Flag{ Name: "max-docs", Usage: "roll over after max docs", Destination: &conf.MaxDocs, diff --git a/cmd/ilm_forecast.go b/cmd/ilm_forecast.go index 1408545..6cfb481 100644 --- a/cmd/ilm_forecast.go +++ b/cmd/ilm_forecast.go @@ -89,7 +89,7 @@ func IlmForecastList(conf *cfg.Config) *cli.Command { }, &cli.BoolFlag{ Name: "hidden", - Usage: "include hidden indicies", + Usage: "include hidden indices", Destination: &conf.Hidden, Aliases: []string{"H"}, }, diff --git a/cmd/index.go b/cmd/index.go index 09fc656..a7d0320 100644 --- a/cmd/index.go +++ b/cmd/index.go @@ -30,7 +30,7 @@ func Index(conf *cfg.Config) *cli.Command { return &cli.Command{ Name: "index", Aliases: []string{"i"}, - Usage: "manage indicies", + Usage: "manage indices", Commands: []*cli.Command{ IndexList(conf), @@ -42,6 +42,7 @@ func Index(conf *cfg.Config) *cli.Command { IndexFields(conf), IndexIlm(conf), IndexDu(conf), + IndexCopy(conf), // sub commands IndexAlias(conf), @@ -54,7 +55,7 @@ func IndexList(conf *cfg.Config) *cli.Command { return &cli.Command{ Name: "list", Aliases: []string{"ls"}, - Usage: "list indicies", + Usage: "list indices", Flags: []cli.Flag{ &cli.IntFlag{ @@ -65,25 +66,25 @@ func IndexList(conf *cfg.Config) *cli.Command { }, &cli.BoolFlag{ Name: "partials", - Usage: "include partial indicies", + Usage: "include partial indices", Destination: &conf.Partials, Aliases: []string{"p"}, }, &cli.BoolFlag{ Name: "hidden", - Usage: "include hidden indicies", + Usage: "include hidden indices", Destination: &conf.Hidden, Aliases: []string{"H"}, }, &cli.BoolFlag{ Name: "failed", - Usage: "include only red failed indicies", + Usage: "include only red failed indices", Destination: &conf.Failed, Aliases: []string{"f"}, }, &cli.StringSliceFlag{ Name: "filter", - Usage: "show only indicies matching the filter", + Usage: "show only indices matching the filter", Destination: &conf.Filter, Aliases: []string{"F"}, }, @@ -299,3 +300,82 @@ func IndexDu(conf *cfg.Config) *cli.Command { }, } } + +func IndexCopy(conf *cfg.Config) *cli.Command { + return &cli.Command{ + Name: "copy", + Aliases: []string{"alias", "cp"}, + Usage: "copy (reindex) documents from one index to another", + UsageText: "index copy [options] -s -t ", + + Flags: []cli.Flag{ + // FIXME: not implemented by typed API + // &cli.BoolFlag{ + // Name: "missing", // op_type + // Usage: "copy only missing docs", + // Destination: &conf.Missing, + // Aliases: []string{"m"}, + // }, + // &cli.BoolFlag{ + // Name: "sync", // version_type to external + // Usage: "create missing docs and update outdated docs", + // Destination: &conf.Sync, + // Aliases: []string{"S"}, + // }, + &cli.BoolFlag{ + Name: "force", // conflicts to proceed + Usage: "continue reindexing even when conflicts happen", + Destination: &conf.Force, + Aliases: []string{"f"}, + }, + &cli.Float64Flag{ + Name: "requests-per-second", + Usage: "the maximum number of documents to index per second (-1 turns off throttling)", + Destination: &conf.RequestsPerSecond, + Aliases: []string{"R"}, + }, + &cli.Int64Flag{ + Name: "max-docs", + Usage: "the maximum number of documents to reindex", + Destination: &conf.MaxDocs, + Aliases: []string{"m"}, + }, + &cli.DurationFlag{ + Name: "timeout", + Usage: "timeout for write operations (e.g. 300m or 120s)", + Destination: &conf.Timeout, + Aliases: []string{"T"}, + }, + &cli.BoolFlag{ + Name: "refresh", + Usage: "refresh affected shards to make this operation visible to search", + Destination: &conf.Refresh, + Aliases: []string{"r"}, + }, + &cli.BoolFlag{ + Name: "wait", + Usage: "wait for active shards", + Destination: &conf.Wait, + Aliases: []string{"w"}, + }, + &cli.StringSliceFlag{ + Name: "source", + Usage: "source index (multiple supported)", + Destination: &conf.SourceIndices, + Aliases: []string{"s"}, + Required: true, + }, + &cli.StringFlag{ + Name: "target", + Usage: "target index", + Destination: &conf.Index, + Aliases: []string{"t"}, + Required: true, + }, + }, + + Action: func(ctx context.Context, cmd *cli.Command) error { + return es.IndexCopy(conf) + }, + } +} diff --git a/cmd/index_alias.go b/cmd/index_alias.go index 7d9ab02..0e6a6e7 100644 --- a/cmd/index_alias.go +++ b/cmd/index_alias.go @@ -121,7 +121,7 @@ func IndexAliasRollover(conf *cfg.Config) *cli.Command { Usage: "roll over after max age (eg: 7d, 2m, 8h)", Destination: &conf.MaxAge, }, - &cli.IntFlag{ + &cli.Int64Flag{ Name: "max-docs", Usage: "roll over after max docs", Destination: &conf.MaxDocs, @@ -159,7 +159,7 @@ func IndexAliasList(conf *cfg.Config) *cli.Command { Flags: []cli.Flag{ &cli.StringSliceFlag{ Name: "filter", - Usage: "show only aliases for indicies matching the filter", + Usage: "show only aliases for indices matching the filter", Destination: &conf.Filter, Aliases: []string{"F"}, }, diff --git a/pkg/cfg/config.go b/pkg/cfg/config.go index 47626d3..0d9c344 100644 --- a/pkg/cfg/config.go +++ b/pkg/cfg/config.go @@ -22,6 +22,7 @@ import ( "os" "path/filepath" "reflect" + "time" "github.com/alecthomas/repr" "gopkg.in/yaml.v3" @@ -68,6 +69,13 @@ type Config struct { Retention string // index template create: -r Rollover bool // index template create: -R + Missing bool // index copy: -m + Sync bool // index copy: -S + RequestsPerSecond float64 // index copy: -R + Timeout time.Duration // index copy: -T + Refresh bool // index copy: -r + SourceIndices []string // index copy: -s + From, To, MaxItems int // search: flags Filter []string // search: -F Path string // search+doc sh: -p @@ -98,9 +106,10 @@ type Config struct { Hidden bool // ds ls: -H // rollover - MaxAge string - MaxDocs, MaxShardSize, MaxShardDocs int // roll over - DryRun bool // rollover: -n + MaxAge string + MaxDocs int64 // roll over, plus others + MaxShardSize, MaxShardDocs int // roll over + DryRun bool // rollover: -n Tag string // api ls: -t HumanCat bool // api repl: -H diff --git a/pkg/es/ccr.go b/pkg/es/ccr.go index 26aa3b2..671f724 100644 --- a/pkg/es/ccr.go +++ b/pkg/es/ccr.go @@ -57,7 +57,7 @@ func CcrStatus(conf *cfg.Config, leader, follower string) error { res, err := conf.Clusters[alias].ES().Cat.Indices(). Do(context.Background()) if err != nil { - return fmt.Errorf("failed to get indicies on %s: %w", alias, esErrorString(err)) + return fmt.Errorf("failed to get indices on %s: %w", alias, esErrorString(err)) } indices[alias] = make(map[string]*types.IndicesRecord, len(res)) diff --git a/pkg/es/cluster.go b/pkg/es/cluster.go index e13bec8..a7af1d6 100644 --- a/pkg/es/cluster.go +++ b/pkg/es/cluster.go @@ -288,7 +288,7 @@ func gatherClusterStats(clusterstats *clusterstats.Response, table *printer.Tabl } table.Entries = append(table.Entries, [][]any{ - {"Indicies", clusterstats.Indices.Count}, + {"indices", clusterstats.Indices.Count}, {"Docs", clusterstats.Indices.Docs.Count}, {"Total Size", printer.Bytes(clusterstats.Indices.Docs.TotalSizeInBytes)}, {"Total Queries", "%d", querycount}, diff --git a/pkg/es/cluster_util.go b/pkg/es/cluster_util.go index 8f7b9ae..67fabcb 100644 --- a/pkg/es/cluster_util.go +++ b/pkg/es/cluster_util.go @@ -89,7 +89,7 @@ func checkClusterStatus(conf *cfg.Config, leader, follower string) bool { status[leader].ActivePrimaryShards, status[follower].ActivePrimaryShards, }, - {"Indicies", + {"indices", len(status[leader].Indices), len(status[follower].Indices), }, @@ -172,7 +172,7 @@ func findIlmErrors(conf *cfg.Config, leader, follower string) bool { return false } -// find unsynchronized indicies only present on leader +// find unsynchronized indices only present on leader func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool { exclude := regexp.MustCompile(DefaultExclude) if conf.Exclude != "" { diff --git a/pkg/es/index.go b/pkg/es/index.go index 1c23061..6aca80c 100644 --- a/pkg/es/index.go +++ b/pkg/es/index.go @@ -43,7 +43,7 @@ func IndexNames(conf *cfg.Config) ([]string, error) { res, err := conf.DefaultCluster.ES().Cat.Indices(). Do(context.Background()) if err != nil { - return nil, fmt.Errorf("failed to get indicies: %w", esErrorString(err)) + return nil, fmt.Errorf("failed to get indices: %w", esErrorString(err)) } indices := make([]string, len(res)) @@ -89,10 +89,10 @@ func IndexList(conf *cfg.Config) error { res, err := cat.Do(context.Background()) if err != nil { - return fmt.Errorf("failed to get indicies: %w", esErrorString(err)) + return fmt.Errorf("failed to get indices: %w", esErrorString(err)) } - slog.Debug("ES result", "indicies", res) + slog.Debug("ES result", "indices", res) list := filterIndices(conf, res) diff --git a/pkg/es/index_copy.go b/pkg/es/index_copy.go new file mode 100644 index 0000000..63c5db6 --- /dev/null +++ b/pkg/es/index_copy.go @@ -0,0 +1,83 @@ +/* +Copyright © 2026 Thomas von Dein + +This program is free software: you can redistribute it and/or modify +it under the terms of the GNU General Public License as published by +the Free Software Foundation, either version 3 of the License, or +(at your option) any later version. + +This program is distributed in the hope that it will be useful, +but WITHOUT ANY WARRANTY; without even the implied warranty of +MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +GNU General Public License for more details. + +You should have received a copy of the GNU General Public License +along with this program. If not, see . +*/ +package es + +import ( + "context" + "fmt" + "time" + + "codeberg.org/scip/esctl/pkg/cfg" + "codeberg.org/scip/esctl/pkg/printer" + "github.com/elastic/go-elasticsearch/v9/typedapi/esdsl" + "github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/conflicts" +) + +func IndexCopy(conf *cfg.Config) error { + copy := conf.DefaultCluster.ES().Reindex() + + if conf.Force { + copy.Conflicts(conflicts.Proceed) + } + + if conf.RequestsPerSecond != 0 { + copy.RequestsPerSecond(fmt.Sprintf("%.2f", conf.RequestsPerSecond)) + } + + if conf.MaxDocs > 0 { + copy.MaxDocs(conf.MaxDocs) + } + + if conf.Timeout > 0 { + copy.Timeout(formatDuration(conf.Timeout)) + } + + if conf.Wait { + copy.WaitForActiveShards("all") + } + + if conf.Refresh { + copy.Refresh(true) + } + + copy.Source(esdsl.NewReindexSource().Index(conf.SourceIndices...)) + copy.Dest(esdsl.NewReindexDestination().Index(conf.Index)) + + res, err := copy.Do(context.Background()) + if err != nil { + return fmt.Errorf("failed to copy indices: %w", esErrorString(err)) + } + + table := printer.NewTable(conf, 2, 0). + WithHeaders("Reindex metrtic", "value") + + table.Entries = [][]any{ + {"Source indices", conf.SourceIndices}, + {"Target index", conf.Index}, + {"Batches", *res.Batches}, + {"Documents total", *res.Total}, + {"Documents created", *res.Created}, + {"Documents deleted", *res.Deleted}, + {"Documents updated", *res.Updated}, + {"Requests/s", *res.RequestsPerSecond}, + {"Timed out", *res.TimedOut}, + {"Time elapsed", time.Duration(*res.Took) * time.Millisecond}, + {"Version conflicts", *res.VersionConflicts}, + } + + return table.Print() +} diff --git a/pkg/es/rollover.go b/pkg/es/rollover.go index 5ee79e9..cc4643d 100644 --- a/pkg/es/rollover.go +++ b/pkg/es/rollover.go @@ -33,7 +33,7 @@ func RolloverConditions(conf *cfg.Config) types.RolloverConditionsVariant { } if conf.MaxDocs > 0 { - cond.MaxDocs(int64(conf.MaxDocs)) + cond.MaxDocs(conf.MaxDocs) } if conf.MaxShardSize > 0 { diff --git a/pkg/es/snapshot.go b/pkg/es/snapshot.go index aa2fb6d..4f13e49 100644 --- a/pkg/es/snapshot.go +++ b/pkg/es/snapshot.go @@ -42,17 +42,17 @@ type Snapshot struct { } func SnapshotList(conf *cfg.Config) error { - // get partial indicies + // get partial indices ires, err := conf.DefaultCluster.ES().Cat.Indices().Do(context.Background()) if err != nil { - return fmt.Errorf("failed to get indicies: %w", esErrorString(err)) + return fmt.Errorf("failed to get indices: %w", esErrorString(err)) } - indicies := map[string]int{} + indices := map[string]int{} for _, index := range ires { name := strings.ReplaceAll(*index.Index, "partial-", "") - indicies[name] = 1 + indices[name] = 1 } // get snapshots @@ -61,7 +61,7 @@ func SnapshotList(conf *cfg.Config) error { return fmt.Errorf("failed to get snapshots: %w", esErrorString(err)) } - slog.Debug("ES result", "indicies", sres) + slog.Debug("ES result", "indices", sres) snapshots := []*Snapshot{} // original snapshot names @@ -74,7 +74,7 @@ func SnapshotList(conf *cfg.Config) error { Orphaned: "no", }) - _, exists := indicies[snap.Forindex] + _, exists := indices[snap.Forindex] if !exists { snap.Orphaned = "orphaned" }