/* 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" "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/types" "github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/sortorder" ) 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 { 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 conf.Tail { case true: return searchTail(conf, searchEs) default: if conf.To > MAXPAGE { return searchPit(conf, req) } else { return searchOnce(conf, searchEs) } } } 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) for _, hit := range res.Hits.Hits { printer.PrintDoc(conf, hit) } 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 } for _, hit := range res.Hits.Hits { printer.PrintDoc(conf, hit) } 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) docs[*hit.Id_] = 1 } time.Sleep(100 * time.Millisecond) } }