mirror of
https://codeberg.org/scip/esctl.git
synced 2026-08-24 14:14:23 +02:00
more search enhancements (#22)
This commit is contained in:
160
pkg/es/search.go
160
pkg/es/search.go
@@ -19,10 +19,22 @@ package es
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"log/slog"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"codeberg.org/scip/esctl/pkg/cfg"
|
||||
"github.com/tidwall/gjson"
|
||||
"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
|
||||
)
|
||||
|
||||
/*
|
||||
@@ -35,36 +47,146 @@ func Search(conf *cfg.Config, queries []string) error {
|
||||
search := conf.DefaultCluster.ES.Search().
|
||||
Index(conf.Index)
|
||||
|
||||
if len(queries) > 0 {
|
||||
req, err := prepareQuery(conf, queries)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
search.Request(req)
|
||||
req, err := prepareQuery(conf, queries)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
res, err := search.Do(context.Background())
|
||||
search.Request(req)
|
||||
|
||||
switch conf.Tail {
|
||||
case true:
|
||||
return searchTail(conf, search)
|
||||
case false:
|
||||
if conf.To > MAXPAGE {
|
||||
return searchPit(conf, req)
|
||||
} else {
|
||||
return searchOnce(conf, search)
|
||||
}
|
||||
}
|
||||
|
||||
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 {
|
||||
if strings.Contains(err.Error(), "reason: all shards failed") {
|
||||
return nil
|
||||
}
|
||||
|
||||
return fmt.Errorf("failed to run search (esdsl): %s", err)
|
||||
}
|
||||
|
||||
slog.Debug("ES result", "search", res)
|
||||
|
||||
for _, hit := range res.Hits.Hits {
|
||||
docjson := fmt.Sprintf(`{"id":%s, "score":%0.4f, "index":"%s", "source":%s}`,
|
||||
*hit.Id_,
|
||||
*hit.Score_,
|
||||
hit.Index_,
|
||||
hit.Source_)
|
||||
printer.PrintDoc(conf, hit)
|
||||
}
|
||||
|
||||
if conf.Path != "" {
|
||||
value := gjson.Get(docjson, conf.Path)
|
||||
fmt.Println(value.String())
|
||||
} else {
|
||||
fmt.Println(docjson)
|
||||
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)
|
||||
|
||||
for {
|
||||
res, err := search.Do(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to run search (esdsl pit): %s", 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 {
|
||||
if strings.Contains(err.Error(), "reason: all shards failed") {
|
||||
return nil
|
||||
}
|
||||
|
||||
return fmt.Errorf("failed to run search (esdsl): %s", 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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,13 +17,16 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
package es
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strings"
|
||||
|
||||
"codeberg.org/scip/esctl/pkg/cfg"
|
||||
"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/operator"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -40,7 +43,22 @@ type filter struct {
|
||||
multi bool // message,title=foo => filter:foo, multi: []string{"message","title"}
|
||||
}
|
||||
|
||||
// build a new filter object
|
||||
// Build a new filter object. We use this to build our elastic query
|
||||
// out of it. We support differnt types of queries:
|
||||
//
|
||||
// - nop filter: no query at all, just return the first N documents.
|
||||
//
|
||||
// - simple ones like: "authenticated", which match across all fields
|
||||
//
|
||||
// - simple ones combined like "authenticated monitoring", if -O was set,
|
||||
// apply an OR logic
|
||||
//
|
||||
// - specific field queries: "message=authenticated" (can be combined too like above)
|
||||
//
|
||||
// - and specific field queries with negation: "message!=authenticated"
|
||||
//
|
||||
// All queries support additional filters using -F field=value, which
|
||||
// must match literally, and range filters using -r "@timestamp:2026-05-28T10:00:00 to now"
|
||||
func NewFilter(query string) (*filter, error) {
|
||||
var separator string
|
||||
var criteria int // we use the constants on top for this
|
||||
@@ -49,9 +67,6 @@ func NewFilter(query string) (*filter, error) {
|
||||
case strings.Contains(query, "!="):
|
||||
criteria = Fmustnot
|
||||
separator = "!="
|
||||
case strings.Contains(query, "?"):
|
||||
criteria = Fshould
|
||||
separator = "?"
|
||||
default:
|
||||
criteria = Fmust
|
||||
separator = "="
|
||||
@@ -59,48 +74,30 @@ func NewFilter(query string) (*filter, error) {
|
||||
|
||||
part := strings.Split(query, separator)
|
||||
if len(part) != 2 {
|
||||
return nil, fmt.Errorf("search queries must be in the form field<sep>pattern where <sep> must be one of: =, !=, ?")
|
||||
return nil, fmt.Errorf("search queries must be in the form field<sep>pattern where <sep> must be one of: = or !=")
|
||||
}
|
||||
|
||||
f := &filter{term: part[0], filter: part[1], criteria: criteria}
|
||||
flt := &filter{term: part[0], filter: part[1], criteria: criteria}
|
||||
|
||||
if strings.Contains(part[0], ",") {
|
||||
// a MultiMatchQuery, match across multiple fields at once
|
||||
multi := strings.Split(part[0], ",")
|
||||
f.multi = true
|
||||
f.mterm = multi
|
||||
flt.multi = true
|
||||
flt.mterm = multi
|
||||
}
|
||||
|
||||
return f, nil
|
||||
return flt, nil
|
||||
}
|
||||
|
||||
// prepare q user search query and turn it into a proper search request
|
||||
func prepareQuery(conf *cfg.Config, queries []string) (*search.Request, error) {
|
||||
if len(queries) == 0 {
|
||||
// nothing given, just return all docs, if any
|
||||
return &search.Request{
|
||||
Query: esdsl.NewMatchAllQuery().QueryCaster(),
|
||||
From: &conf.From,
|
||||
Size: &conf.To,
|
||||
}, nil
|
||||
}
|
||||
|
||||
if len(queries) == 1 && !strings.ContainsAny(queries[0], "!=?") {
|
||||
// a general query w/o fields, search across all fields
|
||||
return &search.Request{
|
||||
Query: esdsl.NewSimpleQueryStringQuery(queries[0]).QueryCaster(),
|
||||
From: &conf.From,
|
||||
Size: &conf.To,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// complex query, form a proper query struct
|
||||
query := esdsl.NewBoolQuery()
|
||||
// Build a complex query set for queries containing field[!]=pattern
|
||||
func mkMatchQueries(queries []string, op operator.Operator) ([]types.QueryVariant, []types.QueryVariant, error) {
|
||||
matchqueries := []types.QueryVariant{}
|
||||
matchNotqueries := []types.QueryVariant{}
|
||||
|
||||
for _, q := range queries {
|
||||
filter, err := NewFilter(q)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// by default we match on a single field
|
||||
@@ -108,23 +105,92 @@ func prepareQuery(conf *cfg.Config, queries []string) (*search.Request, error) {
|
||||
|
||||
if filter.multi {
|
||||
// ok, match across multiple given fields
|
||||
match = esdsl.NewMultiMatchQuery(filter.filter).Fields(filter.mterm...)
|
||||
match = esdsl.NewMultiMatchQuery(filter.filter).Fields(filter.mterm...).Operator(op)
|
||||
}
|
||||
|
||||
// apply logic
|
||||
switch filter.criteria {
|
||||
case Fmustnot:
|
||||
query.MustNot(match)
|
||||
case Fmust:
|
||||
query.Must(match)
|
||||
case Fshould:
|
||||
query.Should(match)
|
||||
if filter.criteria == Fmustnot {
|
||||
matchNotqueries = append(matchNotqueries, match)
|
||||
} else {
|
||||
matchqueries = append(matchqueries, match)
|
||||
}
|
||||
}
|
||||
|
||||
// there might be boolean filters as well
|
||||
filters := make([]types.QueryVariant, len(conf.Filter))
|
||||
return matchqueries, matchNotqueries, nil
|
||||
}
|
||||
|
||||
// Prepare user search query and turn it into a proper search request
|
||||
func prepareQuery(conf *cfg.Config, queries []string) (*search.Request, error) {
|
||||
// logical operator for simple and multimatch queries
|
||||
op := operator.Operator{Name: "AND"}
|
||||
if conf.Or {
|
||||
op = operator.Operator{Name: "OR"}
|
||||
}
|
||||
|
||||
query := esdsl.NewBoolQuery()
|
||||
wholeQuery := strings.Join(queries, " ")
|
||||
|
||||
switch {
|
||||
case len(queries) == 0:
|
||||
// nothing provided via ARGs, so match any docs
|
||||
query.Must(esdsl.NewMatchAllQuery())
|
||||
|
||||
case !strings.ContainsAny(wholeQuery, "!="):
|
||||
// simple query w/o any field[!]=pattern style
|
||||
query.Must(esdsl.NewSimpleQueryStringQuery(wholeQuery).DefaultOperator(op))
|
||||
|
||||
default:
|
||||
// a complex query
|
||||
matchqueries, matchNotqueries, err := mkMatchQueries(queries, op)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// apply boolean logic
|
||||
if conf.Or {
|
||||
query.Should(matchqueries...)
|
||||
} else {
|
||||
// by default we use AND
|
||||
query.Must(matchqueries...)
|
||||
}
|
||||
|
||||
if matchNotqueries != nil {
|
||||
query.MustNot(matchNotqueries...)
|
||||
}
|
||||
}
|
||||
|
||||
slog.Debug("complex query", "query", query)
|
||||
|
||||
// there might be boolean or range filters like -Ffield=value as well
|
||||
filters, err := addFilters(conf)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if filters != nil {
|
||||
query.Filter(filters...)
|
||||
}
|
||||
|
||||
// this being fed into ES.Serach().Req()
|
||||
return &search.Request{
|
||||
Query: query.QueryCaster(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Build a slice of filters, to be fed into query.Filter() including static ones
|
||||
// like -Ffield=value and ranges like -r "@timestamp:2026-05-28T10:00:00 to now"
|
||||
func addFilters(conf *cfg.Config) ([]types.QueryVariant, error) {
|
||||
count := len(conf.Filter)
|
||||
if conf.Range != "" {
|
||||
count++
|
||||
}
|
||||
|
||||
if count == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
filters := make([]types.QueryVariant, count)
|
||||
|
||||
// static filters
|
||||
for idx, filter := range conf.Filter {
|
||||
parts := strings.Split(filter, "=")
|
||||
if len(parts) != 2 {
|
||||
@@ -134,13 +200,43 @@ func prepareQuery(conf *cfg.Config, queries []string) (*search.Request, error) {
|
||||
filters[idx] = esdsl.NewTermQuery(parts[0], esdsl.NewFieldValue().String(parts[1]))
|
||||
}
|
||||
|
||||
if len(filters) > 0 {
|
||||
query.Filter(filters...)
|
||||
// range filters. format: @timestamp:2026-05-05 to 2026-05-15
|
||||
// see: https://www.elastic.co/docs/reference/elasticsearch/rest-apis/common-options#date-math
|
||||
if conf.Range != "" {
|
||||
parts := strings.SplitN(conf.Range, ":", 2)
|
||||
if len(parts) != 2 {
|
||||
return nil, errors.New("invalid range format, expected field:range")
|
||||
}
|
||||
|
||||
field := parts[0]
|
||||
|
||||
parts = strings.Split(parts[1], " to ")
|
||||
if len(parts) != 2 {
|
||||
return nil, errors.New("invalid date range format, expected '<start> to <end>'")
|
||||
}
|
||||
|
||||
from := parts[0]
|
||||
to := parts[1]
|
||||
|
||||
slog.Debug("time range filter",
|
||||
"from", from,
|
||||
"to", to,
|
||||
"format", conf.TimestampFormat,
|
||||
"count", count,
|
||||
)
|
||||
|
||||
rng := esdsl.NewDateRangeQuery(field).
|
||||
Gte(from).
|
||||
Lte(to)
|
||||
|
||||
if strings.Contains(parts[1], ":") {
|
||||
// specific format not needed as long as there are no times specified
|
||||
rng.Format(conf.TimestampFormat)
|
||||
}
|
||||
|
||||
filters[count-1] = rng
|
||||
|
||||
}
|
||||
|
||||
return &search.Request{
|
||||
Query: query.QueryCaster(),
|
||||
From: &conf.From,
|
||||
Size: &conf.To,
|
||||
}, nil
|
||||
return filters, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user