/* 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" "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") for idx, ds := range list { name := printer.Colorize(conf, ds.Status.String(), ds.Name) policy := "" if ds.IlmPolicy != nil { policy = *ds.IlmPolicy } table.Entries[idx] = []any{ name, policy, ds.Hidden, *ds.System, *ds.Replicated, ds.Generation, ds.TimestampField.Name, } if idx == size-1 { break } } table.Sort() 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") ds := res.DataStreams[0] name := printer.Colorize(conf, ds.Status.String(), ds.Name) policy := "" if ds.IlmPolicy != nil { policy = *ds.IlmPolicy } table.Entries = [][]any{ {"name", name}, {"ilm policy", policy}, {"hidden", ds.Hidden}, {"system", *ds.System}, {"replicated", *ds.Replicated}, {"generation", ds.Generation}, {"timestamp field", ds.TimestampField.Name}, {"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 } table = printer.NewTable(conf, 5, len(ds.Indices)) table.Addheaders("backing index name", "uuid", "prefer ilm", "ilm policy", "managed by") for idx, index := range ds.Indices { policy := "" if ds.IlmPolicy != nil { policy = *ds.IlmPolicy } table.Entries[idx] = []any{ index.IndexName, index.IndexUuid, *index.PreferIlm, policy, index.ManagedBy.Name, } if idx < len(ds.Indices) { fmt.Println() } } table.Sort() 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") 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() }