Skip to content

Commit dbcdcbd

Browse files
committed
fix(alerts): resolve lastEvent grouping and add scoped filter contracts
1 parent 6c3af7e commit dbcdcbd

8 files changed

Lines changed: 636 additions & 19 deletions

File tree

Lines changed: 205 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,205 @@
1+
package main
2+
3+
import (
4+
"encoding/json"
5+
"fmt"
6+
"os"
7+
"path/filepath"
8+
"regexp"
9+
"strings"
10+
"testing"
11+
"time"
12+
13+
"github.com/threatwinds/go-sdk/plugins"
14+
"github.com/threatwinds/go-sdk/utils"
15+
"google.golang.org/protobuf/encoding/protojson"
16+
"google.golang.org/protobuf/reflect/protoreflect"
17+
)
18+
19+
func contractPaths(d protoreflect.MessageDescriptor, prefix string, out map[string]bool) {
20+
for i := 0; i < d.Fields().Len(); i++ {
21+
f := d.Fields().Get(i)
22+
p := prefix + f.JSONName()
23+
out[p] = true
24+
if f.Message() != nil && !f.IsMap() && !strings.HasPrefix(string(f.Message().FullName()), "google.protobuf.") {
25+
contractPaths(f.Message(), p+".", out)
26+
}
27+
}
28+
}
29+
30+
// Manifests let independent technology fixes validate their own producers and
31+
// consumers without requiring unrelated, not-yet-merged normalization fixes.
32+
func TestFilterAndRuleContracts(t *testing.T) {
33+
selected := map[string]bool{}
34+
for _, manifest := range loadFilterContracts(t) {
35+
for _, p := range append(manifest.Filters, manifest.Rules...) {
36+
if _, err := os.Stat(filepath.Join("../..", p)); err != nil {
37+
t.Fatal(err)
38+
}
39+
selected[filepath.Clean(filepath.Join("../..", p))] = true
40+
}
41+
}
42+
all := os.Getenv("UTMSTACK_CONTRACT_ALL") == "1"
43+
44+
eventPaths := map[string]bool{}
45+
alertPaths := map[string]bool{}
46+
contractPaths(new(plugins.Event).ProtoReflect().Descriptor(), "", eventPaths)
47+
contractPaths(new(plugins.Alert).ProtoReflect().Descriptor(), "", alertPaths)
48+
arrayIndex := regexp.MustCompile(`\.[0-9]+(\.|$)`)
49+
eventPath := func(p string) bool {
50+
p = strings.TrimSuffix(p, ".keyword")
51+
p = arrayIndex.ReplaceAllString(p, "$1")
52+
return eventPaths[p] || strings.HasPrefix(p, "log.") || strings.HasPrefix(p, "compliance.")
53+
}
54+
alertPath := func(p string) bool {
55+
p = strings.TrimSuffix(p, ".keyword")
56+
if strings.HasPrefix(p, "lastEvent.") {
57+
return eventPath(strings.TrimPrefix(p, "lastEvent."))
58+
}
59+
p = arrayIndex.ReplaceAllString(p, "$1")
60+
return alertPaths[p]
61+
}
62+
cache := plugins.NewCELCache("filter-rule-contract-test")
63+
sample := `{"log":{"messageId":0,"severity":0},"origin":{},"target":{},"action":"","actionResult":"","protocol":"","severity":"","connectionStatus":"","raw":"","dataType":"","dataSource":"","deviceTime":"","tenantId":"","tenantName":"","statusCode":0}`
64+
var expressions func(*testing.T, any)
65+
expressions = func(t *testing.T, v any) {
66+
switch n := v.(type) {
67+
case map[string]any:
68+
for k, x := range n {
69+
if k == "where" {
70+
if w, ok := x.(string); ok && w != "" {
71+
_, err := cache.Eval(w, sample)
72+
// Direct selectors can fail on this empty sample after a successful compile.
73+
// Their presence is not a syntax error or evidence that parsed logs fail.
74+
if err != nil && !strings.Contains(err.Error(), "failed to evaluate program") {
75+
t.Errorf("CEL compilation: %v", err)
76+
}
77+
}
78+
} else {
79+
expressions(t, x)
80+
}
81+
}
82+
case []any:
83+
for _, x := range n {
84+
expressions(t, x)
85+
}
86+
}
87+
}
88+
for _, dir := range []string{"filters", "rules"} {
89+
err := filepath.WalkDir(filepath.Join("../..", dir), func(path string, d os.DirEntry, err error) error {
90+
if err != nil {
91+
return err
92+
}
93+
if d.IsDir() || (filepath.Ext(path) != ".yml" && filepath.Ext(path) != ".yaml") {
94+
return nil
95+
}
96+
if !all && !selected[filepath.Clean(path)] {
97+
return nil
98+
}
99+
t.Run(path, func(t *testing.T) {
100+
b, err := utils.ReadPbYaml(path)
101+
if err != nil {
102+
t.Fatal(err)
103+
}
104+
var doc any
105+
if err = json.Unmarshal(b, &doc); err != nil {
106+
t.Fatal(err)
107+
}
108+
expressions(t, doc)
109+
if dir == "filters" {
110+
cfg := new(plugins.Config)
111+
if err = protojson.Unmarshal(b, cfg); err != nil {
112+
t.Fatal(err)
113+
}
114+
for _, stage := range cfg.Pipeline {
115+
for _, step := range stage.Steps {
116+
fields := []string{}
117+
if s := step.Rename; s != nil {
118+
fields = append(fields, s.To)
119+
}
120+
if s := step.Grok; s != nil {
121+
for _, p := range s.Patterns {
122+
if p.FieldName != "" { // Empty grok names are non-capturing separators.
123+
fields = append(fields, p.FieldName)
124+
}
125+
}
126+
}
127+
if s := step.Csv; s != nil {
128+
fields = append(fields, s.Headers...)
129+
}
130+
if s := step.Add; s != nil {
131+
fields = append(fields, s.Params["key"].GetStringValue())
132+
}
133+
if s := step.Cast; s != nil {
134+
fields = append(fields, s.Fields...)
135+
}
136+
for _, p := range fields {
137+
if !eventPath(p) {
138+
t.Errorf("unknown event write/cast field %q", p)
139+
}
140+
}
141+
}
142+
}
143+
} else {
144+
rule := new(plugins.Rule)
145+
if err = protojson.Unmarshal(b, rule); err != nil {
146+
t.Fatal(err)
147+
}
148+
rule.Normalize()
149+
if len(rule.GroupBy) > 0 && len(rule.DeduplicateBy) > 0 {
150+
t.Error("groupBy and deduplicateBy are mutually exclusive")
151+
}
152+
if rule.Adversary != "" && rule.Adversary != "origin" && rule.Adversary != "target" {
153+
t.Errorf("unknown adversary %s", rule.Adversary)
154+
}
155+
for _, p := range append(rule.GroupBy, rule.DeduplicateBy...) {
156+
if !alertPath(p) {
157+
t.Errorf("unknown alert grouping field %q", p)
158+
}
159+
}
160+
var checkSearch func([]*plugins.SearchRequest)
161+
checkSearch = func(searches []*plugins.SearchRequest) {
162+
for _, s := range searches {
163+
if s.Within != "" {
164+
if _, err := time.ParseDuration(s.Within); err != nil {
165+
t.Error(err)
166+
}
167+
}
168+
for _, x := range s.With {
169+
valid := eventPath(x.Field)
170+
if strings.Contains(s.IndexPattern, "-alert-") {
171+
valid = alertPath(x.Field)
172+
}
173+
if !valid {
174+
t.Errorf("unknown search field %q", x.Field)
175+
}
176+
if x.Value == nil {
177+
t.Errorf("missing search value for %s", x.Field)
178+
continue
179+
}
180+
v := x.Value.GetStringValue()
181+
if strings.HasPrefix(v, "{{.") && strings.HasSuffix(v, "}}") {
182+
p := strings.TrimSuffix(strings.TrimPrefix(v, "{{."), "}}")
183+
if !eventPath(p) {
184+
t.Errorf("unknown event placeholder %q", p)
185+
}
186+
}
187+
switch x.Operator {
188+
case "filter_term", "filter_match", "must_not_term", "must_not_match":
189+
default:
190+
t.Error(fmt.Sprintf("unknown search operator %s", x.Operator))
191+
}
192+
}
193+
checkSearch(s.Or)
194+
}
195+
}
196+
checkSearch(rule.Correlation)
197+
}
198+
})
199+
return nil
200+
})
201+
if err != nil {
202+
t.Fatal(err)
203+
}
204+
}
205+
}
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
package main
2+
3+
import (
4+
"encoding/json"
5+
"os"
6+
"path/filepath"
7+
"testing"
8+
)
9+
10+
type filterContract struct {
11+
Technology string `json:"technology"`
12+
Filters []string `json:"filters"`
13+
Rules []string `json:"rules"`
14+
Fixtures []Fixture `json:"fixtures"`
15+
}
16+
17+
func loadFilterContracts(t *testing.T) []filterContract {
18+
t.Helper()
19+
paths, err := filepath.Glob("testdata/filter-contracts/*.json")
20+
if err != nil {
21+
t.Fatal(err)
22+
}
23+
if len(paths) == 0 {
24+
t.Fatal("no filter contract manifests")
25+
}
26+
var out []filterContract
27+
for _, p := range paths {
28+
data, err := os.ReadFile(p)
29+
if err != nil {
30+
t.Fatal(err)
31+
}
32+
var manifest filterContract
33+
if err := json.Unmarshal(data, &manifest); err != nil {
34+
t.Fatalf("%s: %v", p, err)
35+
}
36+
if manifest.Technology == "" || len(manifest.Filters) == 0 || len(manifest.Fixtures) == 0 {
37+
t.Fatalf("%s: technology, filters and fixtures are required", p)
38+
}
39+
out = append(out, manifest)
40+
}
41+
return out
42+
}

0 commit comments

Comments
 (0)