Files
esctl/pkg/es/datastream.go

248 lines
5.8 KiB
Go
Raw Normal View History

/*
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"
"regexp"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer"
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
)
func DatastreamNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES().Indices.GetDataStream().
Do(context.Background())
if err != nil {
return nil, fmt.Errorf("failed to get data streams: %w", esErrorString(err))
}
dss := make([]string, len(res.DataStreams))
for idx, ds := range res.DataStreams {
dss[idx] = ds.Name
}
return dss, nil
}
func DatastreamList(conf *cfg.Config) error {
res, err := conf.DefaultCluster.ES().Indices.GetDataStream().
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get data streams: %w", esErrorString(err))
}
slog.Debug("ES result", "data streams", res)
list := filterDatastreams(conf, res.DataStreams)
size := len(list)
if conf.MaxItems > 0 {
if size > conf.MaxItems {
size = conf.MaxItems
}
}
table := printer.NewTable(conf, 7, size)
table.Addheaders("name", "ilm policy", "hidden", "system", "replicated", "generation", "timestamp field")
2026-07-07 23:46:43 +02:00
for idx, datastream := range list {
name := printer.Colorize(conf, datastream.Status.String(), datastream.Name)
policy := ""
2026-07-07 23:46:43 +02:00
if datastream.IlmPolicy != nil {
policy = *datastream.IlmPolicy
}
2026-07-07 07:29:03 +02:00
table.Entries[idx] = []any{
name,
policy,
2026-07-07 23:46:43 +02:00
datastream.Hidden,
*datastream.System,
*datastream.Replicated,
datastream.Generation,
datastream.TimestampField.Name,
}
if idx == size-1 {
break
}
}
table.Sort()
2026-07-07 23:46:43 +02:00
if err := table.Print(); err != nil {
return err
}
return nil
}
func filterDatastreams(conf *cfg.Config, list []types.DataStream) []types.DataStream {
selected := []types.DataStream{}
for _, ds := range list {
if !conf.Hidden && ds.Hidden {
continue
}
selected = append(selected, ds)
}
if len(conf.Filter) == 0 {
return selected
}
// we support just one filter here, for now
filter := *regexp.MustCompile(conf.Filter[0])
newlist := []types.DataStream{}
for _, ds := range selected {
if filter.MatchString(ds.Name) {
newlist = append(newlist, ds)
}
}
return newlist
}
func DatastreamShow(conf *cfg.Config, dsname string) error {
res, err := conf.DefaultCluster.ES().Indices.GetDataStream().
Name(dsname).
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get data stream: %w", esErrorString(err))
}
slog.Debug("ES result", "data stream", res)
if len(res.DataStreams) == 0 {
return errors.New("no data stream found with that name")
}
stats, err := conf.DefaultCluster.ES().Indices.DataStreamsStats().
Name(dsname).
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get data stream stats: %w", esErrorString(err))
}
table := printer.NewTable(conf, 11, 2)
table.Addheaders("data stream property", "value")
2026-07-07 23:46:43 +02:00
datastream := res.DataStreams[0]
name := printer.Colorize(conf, datastream.Status.String(), datastream.Name)
policy := ""
2026-07-07 23:46:43 +02:00
if datastream.IlmPolicy != nil {
policy = *datastream.IlmPolicy
}
2026-07-07 07:29:03 +02:00
table.Entries = [][]any{
{"name", name},
{"ilm policy", policy},
2026-07-07 23:46:43 +02:00
{"hidden", datastream.Hidden},
{"system", *datastream.System},
{"replicated", *datastream.Replicated},
{"generation", datastream.Generation},
{"timestamp field", datastream.TimestampField.Name},
2026-07-07 07:29:03 +02:00
{"backing indices", stats.BackingIndices},
{"size", printer.Bytes(stats.TotalStoreSizeBytes)},
{"shards-failed", stats.Shards_.Failed},
{"shards-successful", stats.Shards_.Successful},
}
if err := table.Print(); err != nil {
return err
}
2026-07-07 23:46:43 +02:00
table = printer.NewTable(conf, 5, len(datastream.Indices))
table.Addheaders("backing index name", "uuid", "prefer ilm", "ilm policy", "managed by")
2026-07-07 23:46:43 +02:00
for idx, index := range datastream.Indices {
policy := ""
2026-07-07 23:46:43 +02:00
if datastream.IlmPolicy != nil {
policy = *datastream.IlmPolicy
}
2026-07-07 07:29:03 +02:00
table.Entries[idx] = []any{
index.IndexName,
index.IndexUuid,
2026-07-07 07:29:03 +02:00
*index.PreferIlm,
policy,
index.ManagedBy.Name,
}
2026-07-07 23:46:43 +02:00
if idx < len(datastream.Indices) {
fmt.Println()
}
}
table.Sort()
2026-07-07 23:46:43 +02:00
if err := table.Print(); err != nil {
return err
}
return nil
}
func DatastreamCreate(conf *cfg.Config, dsname string) error {
_, err := conf.DefaultCluster.ES().Indices.CreateDataStream(dsname).
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to create datastream: %w", esErrorString(err))
}
return nil
}
func DatastreamDelete(conf *cfg.Config, dsname string) error {
_, err := conf.DefaultCluster.ES().Indices.DeleteDataStream(dsname).
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to delete datastream: %w", esErrorString(err))
}
return nil
}
func DatastreamRollover(conf *cfg.Config, ds string) error {
res, err := RolloverAlias(conf, ds)
if err != nil {
return fmt.Errorf("failed to rollover data stream: %w", esErrorString(err))
}
table := printer.NewTable(conf, 2, 5)
table.Addheaders("rollover response", "value")
2026-07-07 07:29:03 +02:00
table.Entries = [][]any{
{"acknowledged", res.Acknowledged},
{"rolled over", res.RolledOver},
{"shards acknowledged", res.ShardsAcknowledged},
{"old index", res.OldIndex},
{"new index", res.NewIndex},
}
return table.Print()
}