mirror of
https://codeberg.org/scip/esctl.git
synced 2026-08-25 13:14:18 +02:00
Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f1877c6903 | |||
|
|
2b2094b155 | ||
| 1c72b5ec9a | |||
| aeded16b89 | |||
| 57bbe680a7 | |||
|
|
130d0b879d | ||
|
|
f7302ea506 | ||
| b2da6f0f29 | |||
| c92f10fb5d | |||
| 501deec538 |
65
cmd/doc.go
65
cmd/doc.go
@@ -34,7 +34,7 @@ func Doc(conf *cfg.Config) *cli.Command {
|
|||||||
Commands: []*cli.Command{
|
Commands: []*cli.Command{
|
||||||
DocAdd(conf),
|
DocAdd(conf),
|
||||||
DocShow(conf),
|
DocShow(conf),
|
||||||
//Delete(conf),
|
DocDelete(conf),
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -112,3 +112,66 @@ func DocShow(conf *cfg.Config) *cli.Command {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func DocDelete(conf *cfg.Config) *cli.Command {
|
||||||
|
return &cli.Command{
|
||||||
|
Name: "delete",
|
||||||
|
Aliases: []string{"rm"},
|
||||||
|
Usage: "delete JSON document[s] from index[es]",
|
||||||
|
UsageText: "delete [options] -i <index> [<[field<sep>]pattern> ...]\n" + SearchUsage,
|
||||||
|
|
||||||
|
Flags: []cli.Flag{
|
||||||
|
&cli.BoolFlag{
|
||||||
|
Name: "all",
|
||||||
|
Usage: "delete all docs (dangerous!)",
|
||||||
|
Destination: &conf.All,
|
||||||
|
Aliases: []string{"a"},
|
||||||
|
},
|
||||||
|
&cli.StringFlag{
|
||||||
|
Name: "index",
|
||||||
|
Usage: "index to work on",
|
||||||
|
Sources: cli.EnvVars("ES_INDEX"),
|
||||||
|
Destination: &conf.Index,
|
||||||
|
Aliases: []string{"i"},
|
||||||
|
},
|
||||||
|
&cli.StringSliceFlag{
|
||||||
|
Name: "filter",
|
||||||
|
Usage: "additional boolean filters. format: key=value",
|
||||||
|
Destination: &conf.Filter,
|
||||||
|
Aliases: []string{"F"},
|
||||||
|
},
|
||||||
|
&cli.StringFlag{
|
||||||
|
Name: "timerange",
|
||||||
|
Usage: "field:<date> to <date> (e.g. @timestamp:2026-05-05 to 2026-05-15)",
|
||||||
|
Destination: &conf.Range,
|
||||||
|
Aliases: []string{"r"},
|
||||||
|
},
|
||||||
|
&cli.StringFlag{
|
||||||
|
Name: "timestamp-format",
|
||||||
|
Usage: "a valid ES builtin timestamp or custom format",
|
||||||
|
Destination: &conf.TimestampFormat,
|
||||||
|
Value: "strict_date_hour_minute",
|
||||||
|
},
|
||||||
|
&cli.BoolFlag{
|
||||||
|
Name: "or",
|
||||||
|
Usage: "logical operator (default: and)",
|
||||||
|
Destination: &conf.Or,
|
||||||
|
Aliases: []string{"O"},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
|
||||||
|
Action: func(ctx context.Context, cmd *cli.Command) error {
|
||||||
|
if conf.Subhelp {
|
||||||
|
return showJsonPathHelp()
|
||||||
|
}
|
||||||
|
|
||||||
|
args := cmd.Args()
|
||||||
|
|
||||||
|
if args.Len() == 0 && !conf.All {
|
||||||
|
return errors.New("missing arguments: <query>")
|
||||||
|
}
|
||||||
|
|
||||||
|
return es.DocDelete(conf, cmd.Args().Slice())
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
29
cmd/root.go
29
cmd/root.go
@@ -22,6 +22,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"codeberg.org/scip/esctl/pkg/cfg"
|
||||||
|
"codeberg.org/scip/esctl/pkg/es"
|
||||||
"codeberg.org/scip/esctl/pkg/log"
|
"codeberg.org/scip/esctl/pkg/log"
|
||||||
|
|
||||||
"github.com/urfave/cli/v3"
|
"github.com/urfave/cli/v3"
|
||||||
@@ -54,6 +55,12 @@ func Main() int {
|
|||||||
Sources: cli.EnvVars("ES_DEBUG"),
|
Sources: cli.EnvVars("ES_DEBUG"),
|
||||||
Destination: &conf.Debug,
|
Destination: &conf.Debug,
|
||||||
},
|
},
|
||||||
|
&cli.BoolFlag{
|
||||||
|
Name: "debug-http",
|
||||||
|
Value: false,
|
||||||
|
Usage: "enable HTTP debugging",
|
||||||
|
Destination: &conf.DebugHTTP,
|
||||||
|
},
|
||||||
&cli.StringFlag{
|
&cli.StringFlag{
|
||||||
Name: "config",
|
Name: "config",
|
||||||
Aliases: []string{"c"},
|
Aliases: []string{"c"},
|
||||||
@@ -89,6 +96,7 @@ func Main() int {
|
|||||||
Doc(conf),
|
Doc(conf),
|
||||||
Repl(conf),
|
Repl(conf),
|
||||||
Version(conf),
|
Version(conf),
|
||||||
|
Debug(conf),
|
||||||
},
|
},
|
||||||
|
|
||||||
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {
|
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {
|
||||||
@@ -122,3 +130,24 @@ func Version(conf *cfg.Config) *cli.Command {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func Debug(conf *cfg.Config) *cli.Command {
|
||||||
|
return &cli.Command{
|
||||||
|
Name: "debug",
|
||||||
|
Usage: "developer only",
|
||||||
|
|
||||||
|
Flags: []cli.Flag{
|
||||||
|
&cli.StringFlag{
|
||||||
|
Name: "index",
|
||||||
|
Usage: "index to search within",
|
||||||
|
Sources: cli.EnvVars("ES_INDEX"),
|
||||||
|
Destination: &conf.Index,
|
||||||
|
Aliases: []string{"i"},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
|
||||||
|
Action: func(ctx context.Context, cmd *cli.Command) error {
|
||||||
|
return es.Debug(conf)
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -26,22 +26,33 @@ import (
|
|||||||
"github.com/urfave/cli/v3"
|
"github.com/urfave/cli/v3"
|
||||||
)
|
)
|
||||||
|
|
||||||
func Search(conf *cfg.Config) *cli.Command {
|
const SearchUsage = `<sep> might be one of:
|
||||||
return &cli.Command{
|
|
||||||
Name: "search",
|
|
||||||
Aliases: []string{"/"},
|
|
||||||
Usage: "search within an index",
|
|
||||||
UsageText: `search [options] [<[field<sep>]pattern> ...]
|
|
||||||
|
|
||||||
<sep> might be one of:
|
|
||||||
=: Must match
|
=: Must match
|
||||||
!=: Must not match
|
!=: Must not match
|
||||||
?: Should match
|
|
||||||
|
|
||||||
You can omit a field spec and thereby search across all fields.
|
You can omit a field spec and thereby search across all fields.
|
||||||
|
|
||||||
|
By default all queries contribute to matches (logical AND), use
|
||||||
|
-O to apply a logical OR operator.
|
||||||
|
|
||||||
You can also search multiple fields by separating them with comma, eg:
|
You can also search multiple fields by separating them with comma, eg:
|
||||||
user,group=root`,
|
user,group=root
|
||||||
|
|
||||||
|
Use filters to further restrict results, they must match literally.
|
||||||
|
|
||||||
|
For datetime range format refer to:
|
||||||
|
https://www.elastic.co/docs/reference/elasticsearch/rest-apis/common-options#date-math
|
||||||
|
|
||||||
|
For timestamp formats refer to:
|
||||||
|
https://www.elastic.co/docs/reference/elasticsearch/mapping-reference/mapping-date-format
|
||||||
|
`
|
||||||
|
|
||||||
|
func Search(conf *cfg.Config) *cli.Command {
|
||||||
|
return &cli.Command{
|
||||||
|
Name: "search",
|
||||||
|
Aliases: []string{"/"},
|
||||||
|
Usage: "search within an index",
|
||||||
|
UsageText: "search [options] [<[field<sep>]pattern> ...]\n" + SearchUsage,
|
||||||
|
|
||||||
Flags: []cli.Flag{
|
Flags: []cli.Flag{
|
||||||
&cli.StringFlag{
|
&cli.StringFlag{
|
||||||
@@ -53,17 +64,17 @@ user,group=root`,
|
|||||||
},
|
},
|
||||||
&cli.IntFlag{
|
&cli.IntFlag{
|
||||||
Name: "from",
|
Name: "from",
|
||||||
Usage: "show results FROM (default 0)",
|
Usage: "show results starting at <from>",
|
||||||
Destination: &conf.From,
|
Destination: &conf.From,
|
||||||
Value: 0,
|
Value: 0,
|
||||||
Aliases: []string{"f"},
|
Aliases: []string{"f"},
|
||||||
},
|
},
|
||||||
&cli.IntFlag{
|
&cli.IntFlag{
|
||||||
Name: "to",
|
Name: "len",
|
||||||
Usage: "show results to (default 10)",
|
Usage: "number of results to show (-1: all[max:10k], caution: might be slow)",
|
||||||
Destination: &conf.To,
|
Destination: &conf.To,
|
||||||
Value: 10,
|
Value: 20,
|
||||||
Aliases: []string{"t"},
|
Aliases: []string{"l"},
|
||||||
},
|
},
|
||||||
&cli.StringSliceFlag{
|
&cli.StringSliceFlag{
|
||||||
Name: "filter",
|
Name: "filter",
|
||||||
@@ -77,12 +88,36 @@ user,group=root`,
|
|||||||
Destination: &conf.Path,
|
Destination: &conf.Path,
|
||||||
Aliases: []string{"p"},
|
Aliases: []string{"p"},
|
||||||
},
|
},
|
||||||
|
&cli.StringFlag{
|
||||||
|
Name: "timerange",
|
||||||
|
Usage: "field:<date> to <date> (e.g. @timestamp:2026-05-05 to 2026-05-15)",
|
||||||
|
Destination: &conf.Range,
|
||||||
|
Aliases: []string{"r"},
|
||||||
|
},
|
||||||
|
&cli.StringFlag{
|
||||||
|
Name: "timestamp-format",
|
||||||
|
Usage: "a valid ES builtin timestamp or custom format",
|
||||||
|
Destination: &conf.TimestampFormat,
|
||||||
|
Value: "strict_date_hour_minute",
|
||||||
|
},
|
||||||
&cli.BoolFlag{
|
&cli.BoolFlag{
|
||||||
Name: "help-jsonpath",
|
Name: "help-jsonpath",
|
||||||
Usage: "show jsonPath help",
|
Usage: "show jsonPath help",
|
||||||
Destination: &conf.Subhelp,
|
Destination: &conf.Subhelp,
|
||||||
Aliases: []string{"H"},
|
Aliases: []string{"H"},
|
||||||
},
|
},
|
||||||
|
&cli.BoolFlag{
|
||||||
|
Name: "tail",
|
||||||
|
Usage: "follow search live, like tail -f",
|
||||||
|
Destination: &conf.Tail,
|
||||||
|
Aliases: []string{"T"},
|
||||||
|
},
|
||||||
|
&cli.BoolFlag{
|
||||||
|
Name: "or",
|
||||||
|
Usage: "logical operator (default: and)",
|
||||||
|
Destination: &conf.Or,
|
||||||
|
Aliases: []string{"O"},
|
||||||
|
},
|
||||||
},
|
},
|
||||||
|
|
||||||
Action: func(ctx context.Context, cmd *cli.Command) error {
|
Action: func(ctx context.Context, cmd *cli.Command) error {
|
||||||
@@ -92,6 +127,10 @@ user,group=root`,
|
|||||||
|
|
||||||
args := cmd.Args()
|
args := cmd.Args()
|
||||||
|
|
||||||
|
if conf.To == -1 {
|
||||||
|
conf.To = 10000
|
||||||
|
}
|
||||||
|
|
||||||
return es.Search(conf, args.Slice())
|
return es.Search(conf, args.Slice())
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|||||||
2
go.mod
2
go.mod
@@ -26,7 +26,7 @@ require (
|
|||||||
github.com/mattn/go-isatty v0.0.22
|
github.com/mattn/go-isatty v0.0.22
|
||||||
github.com/olekukonko/tablewriter v1.1.4
|
github.com/olekukonko/tablewriter v1.1.4
|
||||||
github.com/tlinden/yadu v0.1.3
|
github.com/tlinden/yadu v0.1.3
|
||||||
github.com/urfave/cli/v3 v3.9.0
|
github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de
|
||||||
gopkg.in/yaml.v3 v3.0.1
|
gopkg.in/yaml.v3 v3.0.1
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
2
go.sum
2
go.sum
@@ -70,6 +70,8 @@ github.com/urfave/cli/v3 v3.8.0 h1:XqKPrm0q4P0q5JpoclYoCAv0/MIvH/jZ2umzuf8pNTI=
|
|||||||
github.com/urfave/cli/v3 v3.8.0/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso=
|
github.com/urfave/cli/v3 v3.8.0/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso=
|
||||||
github.com/urfave/cli/v3 v3.9.0 h1:AV9lIiPv3ukYnxunaCUsHnEozptYmDN2F0+yWqLMn/c=
|
github.com/urfave/cli/v3 v3.9.0 h1:AV9lIiPv3ukYnxunaCUsHnEozptYmDN2F0+yWqLMn/c=
|
||||||
github.com/urfave/cli/v3 v3.9.0/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso=
|
github.com/urfave/cli/v3 v3.9.0/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso=
|
||||||
|
github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de h1:ESKPiS7inVoBnv4FmgGNZdjWBI/wmvaragoyD3D9nM4=
|
||||||
|
github.com/urfave/cli/v3 v3.9.1-0.20260524212652-be8b79d0c8de/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso=
|
||||||
go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA=
|
go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA=
|
||||||
go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A=
|
go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A=
|
||||||
go.opentelemetry.io/otel v1.35.0 h1:xKWKPxrxB6OtMCbmMY021CqC45J+3Onta9MqjhnusiQ=
|
go.opentelemetry.io/otel v1.35.0 h1:xKWKPxrxB6OtMCbmMY021CqC45J+3Onta9MqjhnusiQ=
|
||||||
|
|||||||
@@ -17,10 +17,13 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
|||||||
package cfg
|
package cfg
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"crypto/tls"
|
"crypto/tls"
|
||||||
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
|
|
||||||
@@ -31,7 +34,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
Version string = `v0.0.13`
|
Version string = `v0.0.14`
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -60,11 +63,16 @@ type Config struct {
|
|||||||
Filter []string // search: -F
|
Filter []string // search: -F
|
||||||
Path string // search+doc sh: -p
|
Path string // search+doc sh: -p
|
||||||
Subhelp bool // search+doc sh: -H
|
Subhelp bool // search+doc sh: -H
|
||||||
|
Tail bool // search: -f [tail]
|
||||||
|
Or bool // search: -O
|
||||||
|
Range string // search: -r
|
||||||
|
TimestampFormat string // search: --timestamp-format
|
||||||
Exclude string // cluster compare: -e (regexp)
|
Exclude string // cluster compare: -e (regexp)
|
||||||
All, Verbose bool // cluster status: -a -v
|
All, Verbose bool // cluster status: -a -v
|
||||||
Persistent, Transient, Default bool // -p -t -D cluster settings set
|
Persistent, Transient, Default bool // -p -t -D cluster settings set
|
||||||
Force bool // ccr follower renew: -f
|
Force bool // ccr follower renew: -f
|
||||||
HaveJQ bool // determined at runtime by ourselfes
|
HaveJQ bool // determined at runtime by ourselfes
|
||||||
|
DebugHTTP bool // root: --debug-http
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewConfig() *Config {
|
func NewConfig() *Config {
|
||||||
@@ -200,18 +208,26 @@ func (conf *Config) LoadConfig() error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (conf *Config) getTransport() elastictransport.Option {
|
||||||
|
transport := &http.Transport{
|
||||||
|
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
||||||
|
}
|
||||||
|
|
||||||
|
if conf.DebugHTTP {
|
||||||
|
return elastictransport.WithTransport(
|
||||||
|
&DebugTransport{Transport: transport},
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
return elastictransport.WithTransport(transport)
|
||||||
|
}
|
||||||
|
|
||||||
func (conf *Config) SetupES() error {
|
func (conf *Config) SetupES() error {
|
||||||
for _, cluster := range conf.Clusters {
|
for _, cluster := range conf.Clusters {
|
||||||
es, err := elasticsearch.NewTyped(
|
es, err := elasticsearch.NewTyped(
|
||||||
elasticsearch.WithAddresses(cluster.Uri),
|
elasticsearch.WithAddresses(cluster.Uri),
|
||||||
elasticsearch.WithBasicAuth(cluster.User, cluster.Pass),
|
elasticsearch.WithBasicAuth(cluster.User, cluster.Pass),
|
||||||
elasticsearch.WithTransportOptions(
|
elasticsearch.WithTransportOptions(conf.getTransport()),
|
||||||
elastictransport.WithTransport(
|
|
||||||
&http.Transport{
|
|
||||||
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
|
||||||
},
|
|
||||||
),
|
|
||||||
),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -224,6 +240,39 @@ func (conf *Config) SetupES() error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// used to print uri, path and body of a request made by the go-client
|
||||||
|
type DebugTransport struct {
|
||||||
|
Transport http.RoundTripper
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t *DebugTransport) RoundTrip(req *http.Request) (*http.Response, error) {
|
||||||
|
content := ""
|
||||||
|
contentline := ""
|
||||||
|
|
||||||
|
if req.ContentLength > 0 {
|
||||||
|
buf := new(bytes.Buffer)
|
||||||
|
body, _ := req.GetBody()
|
||||||
|
|
||||||
|
_, err := buf.ReadFrom(body)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
var pretty bytes.Buffer
|
||||||
|
err = json.Indent(&pretty, buf.Bytes(), "", "\t")
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("json parse error: %s", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
content = pretty.String()
|
||||||
|
contentline = buf.String()
|
||||||
|
}
|
||||||
|
|
||||||
|
slog.Info("req", "host", req.URL.Host, "uri", req.URL.Path, "body", content, "bodyline", contentline)
|
||||||
|
|
||||||
|
return t.Transport.RoundTrip(req)
|
||||||
|
}
|
||||||
|
|
||||||
func fileExists(filename string) bool {
|
func fileExists(filename string) bool {
|
||||||
info, err := os.Stat(filename)
|
info, err := os.Stat(filename)
|
||||||
|
|
||||||
|
|||||||
@@ -163,11 +163,12 @@ func ClusterStatus(conf *cfg.Config) error {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
table := printer.NewTable(conf, 2, 5)
|
table := printer.NewTable(conf, 2, 7)
|
||||||
table.Addheaders(cluster, "status")
|
table.Addheaders(cluster, "status")
|
||||||
|
|
||||||
table.Entries = [][]string{
|
table.Entries = [][]string{
|
||||||
{"Cluster Name", printer.Colorize(conf, clusterhealth.Status.Name, clusterhealth.ClusterName)},
|
{"Cluster Name", clusterhealth.ClusterName},
|
||||||
|
{"ES Status", printer.Colorize(conf, clusterhealth.Status.Name, clusterhealth.Status.Name)},
|
||||||
{"ES Version", info.Version.Int},
|
{"ES Version", info.Version.Int},
|
||||||
{"Active Shards", fmt.Sprintf("%d", clusterhealth.ActiveShards)},
|
{"Active Shards", fmt.Sprintf("%d", clusterhealth.ActiveShards)},
|
||||||
{"Active Primary Shards", fmt.Sprintf("%d", clusterhealth.ActivePrimaryShards)},
|
{"Active Primary Shards", fmt.Sprintf("%d", clusterhealth.ActivePrimaryShards)},
|
||||||
|
|||||||
@@ -25,7 +25,7 @@ import (
|
|||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"codeberg.org/scip/esctl/pkg/cfg"
|
||||||
"codeberg.org/scip/esctl/pkg/printer"
|
"codeberg.org/scip/esctl/pkg/printer"
|
||||||
"github.com/elastic/go-elasticsearch/v9"
|
"github.com/elastic/go-elasticsearch/v9"
|
||||||
"github.com/elastic/go-elasticsearch/v9/typedapi/cluster/health"
|
"github.com/elastic/go-elasticsearch/v9/typedapi/cluster/health"
|
||||||
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
|
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
|
||||||
@@ -75,14 +75,15 @@ func checkClusterStatus(conf *cfg.Config, leader, follower string) bool {
|
|||||||
status[cluster] = st
|
status[cluster] = st
|
||||||
}
|
}
|
||||||
|
|
||||||
table := printer.NewTable(conf, 3, 5)
|
table := printer.NewTable(conf, 3, 6)
|
||||||
|
|
||||||
table.Addheaders("setting", "leader:"+leader, "follower:"+follower)
|
table.Addheaders("setting", "leader:"+leader, "follower:"+follower)
|
||||||
|
|
||||||
table.Entries = [][]string{
|
table.Entries = [][]string{
|
||||||
{"Cluster Name",
|
{"Cluster Name", status[leader].ClusterName, status[follower].ClusterName},
|
||||||
printer.Colorize(conf, status[leader].Status.Name, status[leader].ClusterName),
|
{"Cluster Status",
|
||||||
printer.Colorize(conf, status[follower].Status.Name, status[follower].ClusterName),
|
printer.Colorize(conf, status[leader].Status.Name, status[leader].Status.Name),
|
||||||
|
printer.Colorize(conf, status[follower].Status.Name, status[follower].Status.Name),
|
||||||
},
|
},
|
||||||
{"Active Shards",
|
{"Active Shards",
|
||||||
fmt.Sprintf("%d", status[leader].ActiveShards),
|
fmt.Sprintf("%d", status[leader].ActiveShards),
|
||||||
@@ -107,6 +108,8 @@ func checkClusterStatus(conf *cfg.Config, leader, follower string) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fmt.Println()
|
||||||
|
|
||||||
if status[leader].Status.Name == "green" && status[follower].Status.Name == "green" {
|
if status[leader].Status.Name == "green" && status[follower].Status.Name == "green" {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,9 +21,12 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"time"
|
"math/rand/v2"
|
||||||
|
"strings"
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"codeberg.org/scip/esctl/pkg/cfg"
|
||||||
|
"github.com/elastic/go-elasticsearch/v9/typedapi/core/deletebyquery"
|
||||||
|
"github.com/elastic/go-elasticsearch/v9/typedapi/esdsl"
|
||||||
"github.com/tidwall/gjson"
|
"github.com/tidwall/gjson"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -42,7 +45,7 @@ func DocAdd(conf *cfg.Config, jsondoc string) error {
|
|||||||
return fmt.Errorf("supplied document was not valid JSON: %s", err)
|
return fmt.Errorf("supplied document was not valid JSON: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
now := fmt.Sprintf("%d", time.Now().Unix())
|
now := fmt.Sprintf("%d", rand.Int64())
|
||||||
|
|
||||||
res, err := conf.DefaultCluster.ES.Create(conf.Index, now).
|
res, err := conf.DefaultCluster.ES.Create(conf.Index, now).
|
||||||
Document(data).
|
Document(data).
|
||||||
@@ -87,3 +90,43 @@ func DocShow(conf *cfg.Config, id string) error {
|
|||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func DocDelete(conf *cfg.Config, queries []string) error {
|
||||||
|
if len(queries) == 1 && strings.Contains(queries[0], "id=") {
|
||||||
|
id, _ := strings.CutPrefix(queries[0], "id=")
|
||||||
|
|
||||||
|
_, err := conf.DefaultCluster.ES.Delete(conf.Index, id).
|
||||||
|
Header("content-type", "application/json").
|
||||||
|
Header("accept", "application/json").
|
||||||
|
Do(context.Background())
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("failed to delete doc in index %s: %s", conf.Index, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
req := &deletebyquery.Request{}
|
||||||
|
|
||||||
|
if len(queries) == 0 && conf.All {
|
||||||
|
req.Query = esdsl.NewMatchAllQuery().QueryCaster()
|
||||||
|
} else {
|
||||||
|
queryCaster, err := prepareQuery(conf, queries)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
req.Query = queryCaster
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err := conf.DefaultCluster.ES.DeleteByQuery(conf.Index).
|
||||||
|
Request(req).
|
||||||
|
Header("content-type", "application/json").
|
||||||
|
Header("accept", "application/json").
|
||||||
|
Do(context.Background())
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("failed to delete docs in index %s: %s", conf.Index, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|||||||
182
pkg/es/repl.go
182
pkg/es/repl.go
@@ -17,6 +17,7 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
|||||||
package es
|
package es
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bufio"
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"crypto/tls"
|
"crypto/tls"
|
||||||
@@ -34,19 +35,98 @@ import (
|
|||||||
"github.com/chzyer/readline"
|
"github.com/chzyer/readline"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const intro = `Input format: verb path [data]"
|
||||||
|
|
||||||
|
Example:
|
||||||
|
|
||||||
|
post /yourindex/_ccr/pause_follow
|
||||||
|
put /yourindex/_settings {"number_of_replicas": 1}
|
||||||
|
|
||||||
|
You can also put multiline JSON after the path like:
|
||||||
|
|
||||||
|
put /yourindex/_settings
|
||||||
|
{
|
||||||
|
"number_of_replicas": 1
|
||||||
|
}
|
||||||
|
|
||||||
|
If you do NOT supply a JSON in the first line, you need to hit ENTER
|
||||||
|
twice to complete.`
|
||||||
|
|
||||||
|
func Repl(conf *cfg.Config) error {
|
||||||
|
verbs := []string{"post", "get", "put", "delete"}
|
||||||
|
|
||||||
|
fmt.Println(intro)
|
||||||
|
fmt.Println()
|
||||||
|
|
||||||
|
reader, err := readline.NewEx(&readline.Config{
|
||||||
|
Prompt: "> ",
|
||||||
|
HistoryFile: os.Getenv("HOME") + "/.config/esctl/history",
|
||||||
|
HistoryLimit: 500,
|
||||||
|
InterruptPrompt: "^C",
|
||||||
|
EOFPrompt: "exit",
|
||||||
|
HistorySearchFold: true,
|
||||||
|
})
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("failed to initialize readline lib: %s", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
for {
|
||||||
|
text, err := reader.Readline()
|
||||||
|
if err != nil {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
|
||||||
|
text = strings.TrimSpace(text)
|
||||||
|
|
||||||
|
if text == "" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
parts := strings.SplitN(strings.TrimSpace(text), " ", 3)
|
||||||
|
if len(parts) < 2 {
|
||||||
|
fmt.Println("error: you need to input a verb, uri [and post data]")
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
if !slices.Contains(verbs, strings.ToLower(parts[0])) {
|
||||||
|
fmt.Println("error: verb must be one of " + strings.Join(verbs, ","))
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
if !strings.HasPrefix(parts[1], "/") {
|
||||||
|
parts[1] = "/" + parts[1]
|
||||||
|
}
|
||||||
|
|
||||||
|
json := ""
|
||||||
|
if len(parts) == 3 {
|
||||||
|
// put /uri {json}
|
||||||
|
json = parts[2]
|
||||||
|
}
|
||||||
|
|
||||||
|
// put /uri/<Ret> [json]
|
||||||
|
data, err := readJSON(json)
|
||||||
|
if err != nil {
|
||||||
|
fmt.Println(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = CallAPI(conf, parts[0], parts[1], data)
|
||||||
|
if err != nil {
|
||||||
|
fmt.Printf("failed to call API: %s\n", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
reader.SetPrompt("> ")
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func encodeAuth(username, password string) string {
|
func encodeAuth(username, password string) string {
|
||||||
return base64.StdEncoding.EncodeToString([]byte(username + ":" + password))
|
return base64.StdEncoding.EncodeToString([]byte(username + ":" + password))
|
||||||
}
|
}
|
||||||
|
|
||||||
func CallAPI(conf *cfg.Config, input []string) error {
|
func CallAPI(conf *cfg.Config, verb, path, data string) error {
|
||||||
var data string
|
verb = strings.ToUpper(verb)
|
||||||
|
|
||||||
verb := strings.ToUpper(input[0])
|
|
||||||
path := input[1]
|
|
||||||
|
|
||||||
if len(input) == 3 {
|
|
||||||
data = input[2]
|
|
||||||
}
|
|
||||||
|
|
||||||
// we're using port-forwards anyway
|
// we're using port-forwards anyway
|
||||||
tr := &http.Transport{
|
tr := &http.Transport{
|
||||||
@@ -108,69 +188,37 @@ func prettyfiJson(conf *cfg.Config, raw []byte) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func Repl(conf *cfg.Config) error {
|
// interactively read arbitrary JSON data from STDIN, which is
|
||||||
verbs := []string{"post", "get", "put", "delete"}
|
// virtually a repl inside the primary repl
|
||||||
|
func readJSON(input string) (string, error) {
|
||||||
|
data := ""
|
||||||
|
|
||||||
fmt.Println("Input format: verb path [data]")
|
if input != "" {
|
||||||
fmt.Println("example: post /yourindex/_ccr/pause_follow")
|
data = input
|
||||||
|
} else {
|
||||||
|
scanner := bufio.NewScanner(os.Stdin)
|
||||||
|
for scanner.Scan() {
|
||||||
|
line := strings.TrimSpace(scanner.Text())
|
||||||
|
|
||||||
reader, err := readline.NewEx(&readline.Config{
|
if line == "" {
|
||||||
Prompt: "> ",
|
break
|
||||||
HistoryFile: os.Getenv("HOME") + "/.config/esctl/history",
|
|
||||||
HistoryLimit: 500,
|
|
||||||
InterruptPrompt: "^C",
|
|
||||||
EOFPrompt: "exit",
|
|
||||||
HistorySearchFold: true,
|
|
||||||
})
|
|
||||||
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("failed to initialize readline lib: %s", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
for {
|
|
||||||
text, err := reader.Readline()
|
|
||||||
if err != nil {
|
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
text = strings.TrimSpace(text)
|
|
||||||
|
|
||||||
if text == "" {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
parts := strings.SplitN(strings.TrimSpace(text), " ", 3)
|
|
||||||
if len(parts) < 2 {
|
|
||||||
fmt.Println("error: you need to input a verb, uri [and post data]")
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if !slices.Contains(verbs, strings.ToLower(parts[0])) {
|
|
||||||
fmt.Println("error: verb must be one of " + strings.Join(verbs, ","))
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if !strings.HasPrefix(parts[1], "/") {
|
|
||||||
fmt.Println("error: url path must start with /")
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(parts) == 3 {
|
|
||||||
data := map[string]any{}
|
|
||||||
err := json.Unmarshal([]byte(parts[2]), &data)
|
|
||||||
if err != nil {
|
|
||||||
fmt.Printf("error: input data is not proper JSON: %s", err)
|
|
||||||
continue
|
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
err = CallAPI(conf, parts)
|
data += line
|
||||||
if err != nil {
|
|
||||||
fmt.Printf("failed to call API: %s\n", err)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
reader.SetPrompt("> ")
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
if data == "" {
|
||||||
|
return data, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// validate
|
||||||
|
check := map[string]any{}
|
||||||
|
err := json.Unmarshal([]byte(data), &check)
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("error: input data is not proper JSON: %s", err)
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
return data, nil
|
||||||
}
|
}
|
||||||
|
|||||||
168
pkg/es/search.go
168
pkg/es/search.go
@@ -19,10 +19,22 @@ package es
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"log"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"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
|
||||||
)
|
)
|
||||||
|
|
||||||
/*
|
/*
|
||||||
@@ -32,39 +44,155 @@ Execute an ES search.
|
|||||||
additional filters can be given as -F key=value
|
additional filters can be given as -F key=value
|
||||||
*/
|
*/
|
||||||
func Search(conf *cfg.Config, queries []string) error {
|
func Search(conf *cfg.Config, queries []string) error {
|
||||||
search := conf.DefaultCluster.ES.Search().
|
searchEs := conf.DefaultCluster.ES.Search().
|
||||||
Index(conf.Index)
|
Index(conf.Index)
|
||||||
|
|
||||||
if len(queries) > 0 {
|
queryCaster, err := prepareQuery(conf, queries)
|
||||||
req, err := prepareQuery(conf, queries)
|
if err != nil {
|
||||||
if err != nil {
|
return err
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
search.Request(req)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
res, err := search.Do(context.Background())
|
req := &search.Request{Query: queryCaster}
|
||||||
|
|
||||||
|
searchEs.Request(req)
|
||||||
|
|
||||||
|
switch conf.Tail {
|
||||||
|
case true:
|
||||||
|
return searchTail(conf, searchEs)
|
||||||
|
case false:
|
||||||
|
if conf.To > MAXPAGE {
|
||||||
|
return searchPit(conf, req)
|
||||||
|
} else {
|
||||||
|
return searchOnce(conf, searchEs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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 {
|
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)
|
return fmt.Errorf("failed to run search (esdsl): %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
slog.Debug("ES result", "search", res)
|
slog.Debug("ES result", "search", res)
|
||||||
|
|
||||||
for _, hit := range res.Hits.Hits {
|
for _, hit := range res.Hits.Hits {
|
||||||
docjson := fmt.Sprintf(`{"id":%s, "score":%0.4f, "index":"%s", "source":%s}`,
|
printer.PrintDoc(conf, hit)
|
||||||
*hit.Id_,
|
}
|
||||||
*hit.Score_,
|
|
||||||
hit.Index_,
|
|
||||||
hit.Source_)
|
|
||||||
|
|
||||||
if conf.Path != "" {
|
return nil
|
||||||
value := gjson.Get(docjson, conf.Path)
|
}
|
||||||
fmt.Println(value.String())
|
|
||||||
} else {
|
// https://www.elastic.co/docs/reference/elasticsearch/clients/go/using-the-api/searching#_pit_search_after
|
||||||
fmt.Println(docjson)
|
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 {
|
||||||
|
if strings.Contains(err.Error(), "reason: all shards failed") {
|
||||||
|
return 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
|
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,15 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
|||||||
package es
|
package es
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"codeberg.org/scip/esctl/pkg/cfg"
|
"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/esdsl"
|
||||||
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
|
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
|
||||||
|
"github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/operator"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -40,7 +42,22 @@ type filter struct {
|
|||||||
multi bool // message,title=foo => filter:foo, multi: []string{"message","title"}
|
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) {
|
func NewFilter(query string) (*filter, error) {
|
||||||
var separator string
|
var separator string
|
||||||
var criteria int // we use the constants on top for this
|
var criteria int // we use the constants on top for this
|
||||||
@@ -49,9 +66,6 @@ func NewFilter(query string) (*filter, error) {
|
|||||||
case strings.Contains(query, "!="):
|
case strings.Contains(query, "!="):
|
||||||
criteria = Fmustnot
|
criteria = Fmustnot
|
||||||
separator = "!="
|
separator = "!="
|
||||||
case strings.Contains(query, "?"):
|
|
||||||
criteria = Fshould
|
|
||||||
separator = "?"
|
|
||||||
default:
|
default:
|
||||||
criteria = Fmust
|
criteria = Fmust
|
||||||
separator = "="
|
separator = "="
|
||||||
@@ -59,48 +73,30 @@ func NewFilter(query string) (*filter, error) {
|
|||||||
|
|
||||||
part := strings.Split(query, separator)
|
part := strings.Split(query, separator)
|
||||||
if len(part) != 2 {
|
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], ",") {
|
if strings.Contains(part[0], ",") {
|
||||||
// a MultiMatchQuery, match across multiple fields at once
|
// a MultiMatchQuery, match across multiple fields at once
|
||||||
multi := strings.Split(part[0], ",")
|
multi := strings.Split(part[0], ",")
|
||||||
f.multi = true
|
flt.multi = true
|
||||||
f.mterm = multi
|
flt.mterm = multi
|
||||||
}
|
}
|
||||||
|
|
||||||
return f, nil
|
return flt, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// prepare q user search query and turn it into a proper search request
|
// Build a complex query set for queries containing field[!]=pattern
|
||||||
func prepareQuery(conf *cfg.Config, queries []string) (*search.Request, error) {
|
func mkMatchQueries(queries []string, op operator.Operator) ([]types.QueryVariant, []types.QueryVariant, error) {
|
||||||
if len(queries) == 0 {
|
matchqueries := []types.QueryVariant{}
|
||||||
// nothing given, just return all docs, if any
|
matchNotqueries := []types.QueryVariant{}
|
||||||
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()
|
|
||||||
|
|
||||||
for _, q := range queries {
|
for _, q := range queries {
|
||||||
filter, err := NewFilter(q)
|
filter, err := NewFilter(q)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// by default we match on a single field
|
// by default we match on a single field
|
||||||
@@ -108,23 +104,89 @@ func prepareQuery(conf *cfg.Config, queries []string) (*search.Request, error) {
|
|||||||
|
|
||||||
if filter.multi {
|
if filter.multi {
|
||||||
// ok, match across multiple given fields
|
// 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
|
if filter.criteria == Fmustnot {
|
||||||
switch filter.criteria {
|
matchNotqueries = append(matchNotqueries, match)
|
||||||
case Fmustnot:
|
} else {
|
||||||
query.MustNot(match)
|
matchqueries = append(matchqueries, match)
|
||||||
case Fmust:
|
|
||||||
query.Must(match)
|
|
||||||
case Fshould:
|
|
||||||
query.Should(match)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// there might be boolean filters as well
|
return matchqueries, matchNotqueries, nil
|
||||||
filters := make([]types.QueryVariant, len(conf.Filter))
|
}
|
||||||
|
|
||||||
|
// Prepare user search query and turn it into a proper search request
|
||||||
|
func prepareQuery(conf *cfg.Config, queries []string) (*types.Query, 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...)
|
||||||
|
}
|
||||||
|
|
||||||
|
return 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 {
|
for idx, filter := range conf.Filter {
|
||||||
parts := strings.Split(filter, "=")
|
parts := strings.Split(filter, "=")
|
||||||
if len(parts) != 2 {
|
if len(parts) != 2 {
|
||||||
@@ -134,13 +196,43 @@ func prepareQuery(conf *cfg.Config, queries []string) (*search.Request, error) {
|
|||||||
filters[idx] = esdsl.NewTermQuery(parts[0], esdsl.NewFieldValue().String(parts[1]))
|
filters[idx] = esdsl.NewTermQuery(parts[0], esdsl.NewFieldValue().String(parts[1]))
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(filters) > 0 {
|
// range filters. format: @timestamp:2026-05-05 to 2026-05-15
|
||||||
query.Filter(filters...)
|
// 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{
|
return filters, nil
|
||||||
Query: query.QueryCaster(),
|
|
||||||
From: &conf.From,
|
|
||||||
Size: &conf.To,
|
|
||||||
}, nil
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ import (
|
|||||||
|
|
||||||
var (
|
var (
|
||||||
green = color.New(color.BgGreen, color.FgHiWhite).SprintFunc()
|
green = color.New(color.BgGreen, color.FgHiWhite).SprintFunc()
|
||||||
|
yellow = color.New(color.BgYellow, color.FgHiWhite).SprintFunc()
|
||||||
orange = color.BgRGB(255, 128, 0).AddRGB(0, 0, 0).SprintFunc()
|
orange = color.BgRGB(255, 128, 0).AddRGB(0, 0, 0).SprintFunc()
|
||||||
red = color.New(color.BgRed, color.FgHiWhite).SprintFunc()
|
red = color.New(color.BgRed, color.FgHiWhite).SprintFunc()
|
||||||
bold = color.New(color.Bold).SprintFunc()
|
bold = color.New(color.Bold).SprintFunc()
|
||||||
@@ -42,6 +43,8 @@ func Colorize(conf *cfg.Config, col, what string) string {
|
|||||||
what = orange(what)
|
what = orange(what)
|
||||||
case "red":
|
case "red":
|
||||||
what = red(what)
|
what = red(what)
|
||||||
|
case "yellow":
|
||||||
|
what = yellow(what)
|
||||||
}
|
}
|
||||||
|
|
||||||
return what
|
return what
|
||||||
|
|||||||
45
pkg/printer/doc.go
Normal file
45
pkg/printer/doc.go
Normal file
@@ -0,0 +1,45 @@
|
|||||||
|
/*
|
||||||
|
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 printer
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
|
||||||
|
"codeberg.org/scip/esctl/pkg/cfg"
|
||||||
|
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
|
||||||
|
"github.com/tidwall/gjson"
|
||||||
|
)
|
||||||
|
|
||||||
|
func PrintDoc(conf *cfg.Config, hit types.Hit) {
|
||||||
|
var score types.Float64
|
||||||
|
if hit.Score_ != nil {
|
||||||
|
score = *hit.Score_
|
||||||
|
}
|
||||||
|
|
||||||
|
docjson := fmt.Sprintf(`{"id":"%s", "score":%0.4f, "index":"%s", "source":%s}`,
|
||||||
|
*hit.Id_,
|
||||||
|
score,
|
||||||
|
hit.Index_,
|
||||||
|
hit.Source_)
|
||||||
|
|
||||||
|
if conf.Path != "" {
|
||||||
|
value := gjson.Get(docjson, conf.Path)
|
||||||
|
fmt.Println(value.String())
|
||||||
|
} else {
|
||||||
|
fmt.Println(docjson)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user