Files
esctl/pkg/es/search.go

188 lines
4.4 KiB
Go
Raw Permalink Normal View History

2026-04-21 10:50:09 +02:00
/*
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"
"fmt"
2026-05-29 10:19:18 +02:00
"log"
"log/slog"
2026-05-29 10:19:18 +02:00
"time"
2026-04-21 10:50:09 +02:00
"codeberg.org/scip/esctl/pkg/cfg"
2026-05-29 10:19:18 +02:00
"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
2026-04-21 10:50:09 +02:00
)
/*
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 {
2026-05-29 11:05:57 +02:00
searchEs := conf.DefaultCluster.ES.Search().
Index(conf.Index)
2026-05-29 11:05:57 +02:00
queryCaster, err := prepareQuery(conf, queries)
2026-05-29 10:19:18 +02:00
if err != nil {
return err
}
2026-05-29 11:05:57 +02:00
req := &search.Request{Query: queryCaster}
searchEs.Request(req)
2026-05-29 10:19:18 +02:00
searchEs = addSort(conf, searchEs)
2026-05-29 10:19:18 +02:00
switch conf.Tail {
case true:
2026-05-29 11:05:57 +02:00
return searchTail(conf, searchEs)
default:
2026-05-29 10:19:18 +02:00
if conf.To > MAXPAGE {
return searchPit(conf, req)
} else {
2026-05-29 11:05:57 +02:00
return searchOnce(conf, searchEs)
}
2026-05-29 10:19:18 +02:00
}
}
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
}
2026-04-21 10:50:09 +02:00
2026-05-29 10:19:18 +02:00
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())
2026-04-21 10:50:09 +02:00
if err != nil {
return fmt.Errorf("failed to run search (esdsl): %s", esErrorString(err))
2026-04-21 10:50:09 +02:00
}
slog.Debug("ES result", "search", res)
2026-04-21 10:50:09 +02:00
for _, hit := range res.Hits.Hits {
2026-05-29 10:19:18 +02:00
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)
2026-05-29 10:19:18 +02:00
for {
res, err := search.Do(ctx)
if err != nil {
return fmt.Errorf("failed to run search (esdsl pit): %s", esErrorString(err))
2026-05-29 10:19:18 +02:00
}
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")))
}
2026-04-21 10:50:09 +02:00
}
return nil
}
2026-05-29 10:19:18 +02:00
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))
2026-05-29 10:19:18 +02:00
}
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)
}
}