Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 122 additions & 0 deletions app/vtselect/common/extra_filters.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
package common

import (
"fmt"
"regexp"
"strings"

"github.com/VictoriaMetrics/VictoriaLogs/lib/logstorage"
"github.com/valyala/fastjson"
)

// ParseExtraFilters parses extra_filters from either LogsQL or JSON format.
func ParseExtraFilters(s string) (*logstorage.Filter, error) {
if s == "" {
return nil, nil
}
if !strings.HasPrefix(s, `{"`) {
return logstorage.ParseFilter(s)
}

// Extra filters in the form {"field":"value",...}.
filters, err := parseExtraFiltersJSON(s)
if err != nil {
return nil, err
}

result := make([]string, len(filters))
for i, f := range filters {
if len(f.values) == 1 {
result[i] = fmt.Sprintf("%q:=%q", f.key, f.values[0])
} else {
orValues := make([]string, len(f.values))
for j, v := range f.values {
orValues[j] = fmt.Sprintf("%q", v)
}
result[i] = fmt.Sprintf("%q:in(%s)", f.key, strings.Join(orValues, ","))
}
}
return logstorage.ParseFilter(strings.Join(result, " "))
}

// ParseExtraStreamFilters parses extra_stream_filters from either LogsQL or JSON format.
func ParseExtraStreamFilters(s string) (*logstorage.Filter, error) {
if s == "" {
return nil, nil
}
if !strings.HasPrefix(s, `{"`) {
return logstorage.ParseFilter(s)
}

// Extra stream filters in the form {"field":"value",...}.
filters, err := parseExtraFiltersJSON(s)
if err != nil {
return nil, err
}

result := make([]string, len(filters))
for i, f := range filters {
if len(f.values) == 1 {
result[i] = fmt.Sprintf("%q=%q", f.key, f.values[0])
} else {
orValues := make([]string, len(f.values))
for j, v := range f.values {
orValues[j] = regexp.QuoteMeta(v)
}
result[i] = fmt.Sprintf("%q=~%q", f.key, strings.Join(orValues, "|"))
}
}
return logstorage.ParseFilter("{" + strings.Join(result, ",") + "}")
}

type extraFilter struct {
key string
values []string
}

func parseExtraFiltersJSON(s string) ([]extraFilter, error) {
v, err := fastjson.Parse(s)
if err != nil {
return nil, err
}
o := v.GetObject()

var errOuter error
var filters []extraFilter
o.Visit(func(k []byte, v *fastjson.Value) {
if errOuter != nil {
return
}
switch v.Type() {
case fastjson.TypeString:
filters = append(filters, extraFilter{
key: string(k),
values: []string{string(v.GetStringBytes())},
})
case fastjson.TypeArray:
a := v.GetArray()
if len(a) == 0 {
return
}
orValues := make([]string, len(a))
for i, av := range a {
ov, err := av.StringBytes()
if err != nil {
errOuter = fmt.Errorf("cannot obtain string item at the array for key %q; item: %s", k, av)
return
}
orValues[i] = string(ov)
}
filters = append(filters, extraFilter{
key: string(k),
values: orValues,
})
default:
errOuter = fmt.Errorf("unexpected type of value for key %q: %s; value: %s", k, v.Type(), v)
}
})
if errOuter != nil {
return nil, errOuter
}
return filters, nil
}
62 changes: 62 additions & 0 deletions app/vtselect/common/extra_filters_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
package common

import "testing"

func TestParseExtraFilters(t *testing.T) {
tests := []struct {
name string
input string
want string
}{
{"empty", "", ""},
{"json string", `{"foo":"bar"}`, `foo:=bar`},
{"json array", `{"foo":["bar","baz"]}`, `foo:in(bar,baz)`},
{"json mixed", `{"z":"=b ","c":["d","e,"],"a":[],"_msg":"x"}`, `z:="=b " c:in(d,"e,") =x`},
{"logsql", `foo:(bar or baz) error _time:5m {"foo"=bar,baz="z"}`, `{foo="bar",baz="z"} (foo:bar or foo:baz) error _time:5m`},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
f, err := ParseExtraFilters(tt.input)
if err != nil {
t.Fatal(err)
}
if got := f.String(); got != tt.want {
t.Fatalf("got %q; want %q", got, tt.want)
}
})
}

for _, input := range []string{`{"foo"}`, `[1,2]`, `{"foo":[1]}`, `foo:(bar`, `foo | count()`} {
if _, err := ParseExtraFilters(input); err == nil {
t.Fatalf("expected error for %q", input)
}
}
}

func TestParseExtraStreamFilters(t *testing.T) {
tests := []struct {
input string
want string
}{
{"", ""},
{`{"foo":"bar"}`, `{foo="bar"}`},
{`{"foo":["bar","baz"]}`, `{foo=~"bar|baz"}`},
{`{"z":"b","c":["d","e|\""],"a":[],"_msg":"x"}`, `{z="b",c=~"d|e\\|\"",_msg="x"}`},
{`foo:(bar or baz) error _time:5m {"foo"=bar,baz="z"}`, `{foo="bar",baz="z"} (foo:bar or foo:baz) error _time:5m`},
}
for _, tt := range tests {
f, err := ParseExtraStreamFilters(tt.input)
if err != nil {
t.Fatal(err)
}
if got := f.String(); got != tt.want {
t.Fatalf("got %q; want %q", got, tt.want)
}
}

for _, input := range []string{`{"foo"}`, `[1,2]`, `{"foo":[1]}`, `foo:(bar`, `foo | count()`} {
if _, err := ParseExtraStreamFilters(input); err == nil {
t.Fatalf("expected error for %q", input)
}
}
}
119 changes: 3 additions & 116 deletions app/vtselect/logsql/logsql.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ import (
"io"
"math"
"net/http"
"regexp"
"slices"
"sort"
"strconv"
Expand All @@ -26,9 +25,9 @@ import (
"github.com/VictoriaMetrics/VictoriaMetrics/lib/logger"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/timeutil"
"github.com/VictoriaMetrics/metrics"
"github.com/valyala/fastjson"
"github.com/valyala/quicktemplate"

"github.com/VictoriaMetrics/VictoriaTraces/app/vtselect/common"
"github.com/VictoriaMetrics/VictoriaTraces/app/vtstorage"
)

Expand Down Expand Up @@ -1553,7 +1552,7 @@ func parseCommonArgsWithConfig(r *http.Request, skipMaxRangeCheck bool) (*common

// Parse optional extra_filters
for _, extraFiltersStr := range r.Form["extra_filters"] {
extraFilters, err := parseExtraFilters(extraFiltersStr)
extraFilters, err := common.ParseExtraFilters(extraFiltersStr)
if err != nil {
return nil, err
}
Expand All @@ -1562,7 +1561,7 @@ func parseCommonArgsWithConfig(r *http.Request, skipMaxRangeCheck bool) (*common

// Parse optional extra_stream_filters
for _, extraStreamFiltersStr := range r.Form["extra_stream_filters"] {
extraStreamFilters, err := parseExtraStreamFilters(extraStreamFiltersStr)
extraStreamFilters, err := common.ParseExtraStreamFilters(extraStreamFiltersStr)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -1652,118 +1651,6 @@ func getTimeNsec(r *http.Request, argName string) (int64, bool, error) {
return nsecs, true, nil
}

func parseExtraFilters(s string) (*logstorage.Filter, error) {
if s == "" {
return nil, nil
}
if !strings.HasPrefix(s, `{"`) {
return logstorage.ParseFilter(s)
}

// Extra filters in the form {"field":"value",...}.
kvs, err := parseExtraFiltersJSON(s)
if err != nil {
return nil, err
}

filters := make([]string, len(kvs))
for i, kv := range kvs {
if len(kv.values) == 1 {
filters[i] = fmt.Sprintf("%q:=%q", kv.key, kv.values[0])
} else {
orValues := make([]string, len(kv.values))
for j, v := range kv.values {
orValues[j] = fmt.Sprintf("%q", v)
}
filters[i] = fmt.Sprintf("%q:in(%s)", kv.key, strings.Join(orValues, ","))
}
}
s = strings.Join(filters, " ")
return logstorage.ParseFilter(s)
}

func parseExtraStreamFilters(s string) (*logstorage.Filter, error) {
if s == "" {
return nil, nil
}
if !strings.HasPrefix(s, `{"`) {
return logstorage.ParseFilter(s)
}

// Extra stream filters in the form {"field":"value",...}.
kvs, err := parseExtraFiltersJSON(s)
if err != nil {
return nil, err
}

filters := make([]string, len(kvs))
for i, kv := range kvs {
if len(kv.values) == 1 {
filters[i] = fmt.Sprintf("%q=%q", kv.key, kv.values[0])
} else {
orValues := make([]string, len(kv.values))
for j, v := range kv.values {
orValues[j] = regexp.QuoteMeta(v)
}
filters[i] = fmt.Sprintf("%q=~%q", kv.key, strings.Join(orValues, "|"))
}
}
s = "{" + strings.Join(filters, ",") + "}"
return logstorage.ParseFilter(s)
}

type extraFilter struct {
key string
values []string
}

func parseExtraFiltersJSON(s string) ([]extraFilter, error) {
v, err := fastjson.Parse(s)
if err != nil {
return nil, err
}
o := v.GetObject()

var errOuter error
var filters []extraFilter
o.Visit(func(k []byte, v *fastjson.Value) {
if errOuter != nil {
return
}
switch v.Type() {
case fastjson.TypeString:
filters = append(filters, extraFilter{
key: string(k),
values: []string{string(v.GetStringBytes())},
})
case fastjson.TypeArray:
a := v.GetArray()
if len(a) == 0 {
return
}
orValues := make([]string, len(a))
for i, av := range a {
ov, err := av.StringBytes()
if err != nil {
errOuter = fmt.Errorf("cannot obtain string item at the array for key %q; item: %s", k, av)
return
}
orValues[i] = string(ov)
}
filters = append(filters, extraFilter{
key: string(k),
values: orValues,
})
default:
errOuter = fmt.Errorf("unexpected type of value for key %q: %s; value: %s", k, v.Type(), v)
}
})
if errOuter != nil {
return nil, errOuter
}
return filters, nil
}

func getPositiveInt(r *http.Request, argName string) (int, error) {
n, err := httputil.GetInt(r, argName)
if err != nil {
Expand Down
Loading