/* 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" "encoding/json" "fmt" "log" "log/slog" "time" "codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/printer" "github.com/alecthomas/repr" "github.com/elastic/go-elasticsearch/v9/typedapi/core/search" "github.com/elastic/go-elasticsearch/v9/typedapi/esdsl" "github.com/elastic/go-elasticsearch/v9/typedapi/indices/validatequery" "github.com/elastic/go-elasticsearch/v9/typedapi/types" "github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/sortorder" "github.com/tidwall/gjson" ) const ( MAXPAGE = 5000 ) /* 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, queries []string) error { if conf.Validate { return validateSearch(conf, queries) } searchEs := conf.DefaultCluster.ES().Search().Index(conf.Index) queryCaster, err := prepareQuery(conf, queries) if err != nil { return err } req := &search.Request{Query: queryCaster} searchEs.Request(req) searchEs = addSort(conf, searchEs) switch { case conf.Tail: return searchTail(conf, searchEs) case conf.Explain: return explainSearch(conf, searchEs) default: if conf.To > MAXPAGE { return searchPit(conf, req) } else { return searchOnce(conf, searchEs) } } } func explainSearch(conf *cfg.Config, search *search.Search) error { res, err := search. Explain(true). Size(1). // one's enough for explain Do(context.Background()) if err != nil { return fmt.Errorf("failed to call explain search (esdsl): %s", esErrorString(err)) } if conf.Debug { raw, err := json.Marshal(res) if err != nil { return fmt.Errorf("failed to marshal explain result: %s", err) } value := gjson.Get(string(raw), "hits.hits.0._explanation") fmt.Println(value.String()) } if len(res.Hits.Hits) > 0 { ex := res.Hits.Hits[0].Explanation_ fmt.Println(ex.Description) fmt.Println(ex.Value) // recurse into explanation details (it's a tree) for _, ex := range ex.Details { explain(&ex, " ") } } return nil } func explain(res *types.ExplanationDetail, indent string) { fmt.Println(indent + "- " + res.Description) for _, ex := range res.Details { fmt.Println(indent + " - " + ex.Description) fmt.Println(indent + fmt.Sprintf(" score: %f", res.Value)) explain(&ex, indent+" ") } } func validateSearch(conf *cfg.Config, queries []string) error { validate := conf.DefaultCluster.ES().Indices.ValidateQuery() queryCaster, err := prepareQuery(conf, queries) if err != nil { return err } req := &validatequery.Request{Query: queryCaster} validate.Request(req) res, err := validate. Do(context.Background()) if err != nil { return fmt.Errorf("failed to validate search (esdsl): %s", esErrorString(err)) } slog.Debug("ES result", "search", res) if res.Valid { fmt.Println(printer.Colorize(conf, "green", "valid")) } else { fmt.Println(printer.Colorize(conf, "red", "invalid")) } return nil } func Debug(conf *cfg.Config) error { res, err := conf.DefaultCluster.ES().Search(). Index(conf.Index). Size(0). Aggregations(map[string]types.Aggregations{ "min_ts": *esdsl.NewMinAggregation().Field("@timestamp").AggregationsCaster(), "max_ts": *esdsl.NewMaxAggregation().Field("@timestamp").AggregationsCaster(), }). Do(context.Background()) if err != nil { return err } repr.Println(res) return nil } func searchOnce(conf *cfg.Config, search *search.Search) error { res, err := search. From(conf.From). Size(conf.To). Do(context.Background()) if err != nil { return fmt.Errorf("failed to run search (esdsl): %s", esErrorString(err)) } slog.Debug("ES result", "search", res) printer.PrintDocs(conf, res.Hits.Hits) return nil } // 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) 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) if err != nil { log.Fatalf("failed to close PIT: %s", err) } }() search := conf.DefaultCluster.ES().Search(). Request(req). Pit(esdsl.NewPointInTimeReference(). Id(pit.Id). KeepAlive(esdsl.NewDuration().String("1m"))). Sort(esdsl.NewSortOptions(). AddSortOption("_shard_doc", esdsl.NewFieldSort(sortorder.Asc))). Size(conf.To) search = addSort(conf, search) for { res, err := search.Do(ctx) if err != nil { return fmt.Errorf("failed to run search (esdsl pit): %s", esErrorString(err)) } if len(res.Hits.Hits) == 0 { break } printer.PrintDocs(conf, res.Hits.Hits) last := res.Hits.Hits[len(res.Hits.Hits)-1] search = search.SearchAfterValues(last.Sort) if res.PitId != nil { search = search.Pit(esdsl.NewPointInTimeReference(). Id(*res.PitId). KeepAlive(esdsl.NewDuration().String("1m"))) } } return nil } func searchTail(conf *cfg.Config, search *search.Search) error { docs := map[string]int{} fmt.Println("enter ctrl-c to abort...") for { res, err := search.Do(context.Background()) if err != nil { return fmt.Errorf("failed to run search (esdsl): %s", esErrorString(err)) } slog.Debug("ES result", "search", res) for _, hit := range res.Hits.Hits { _, exists := docs[*hit.Id_] if exists { continue } printer.PrintDoc(conf, hit) fmt.Println() docs[*hit.Id_] = 1 } time.Sleep(100 * time.Millisecond) } }