This commit is contained in:
@@ -0,0 +1,152 @@
|
||||
package rules
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/actions"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/alerts"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/analysis"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/loki"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/store"
|
||||
)
|
||||
|
||||
type Store interface {
|
||||
Load() (store.Data, error)
|
||||
UpdateRule(id string, fn func(*store.Rule)) error
|
||||
AddIncident(incident store.Incident) error
|
||||
}
|
||||
|
||||
type Engine struct {
|
||||
store Store
|
||||
loki *loki.Client
|
||||
analyzer analysis.Analyzer
|
||||
actions *actions.Runner
|
||||
dispatcher *alerts.Dispatcher
|
||||
}
|
||||
|
||||
func NewEngine(s Store, lokiClient *loki.Client, analyzer analysis.Analyzer, actionRunner *actions.Runner, dispatcher *alerts.Dispatcher) *Engine {
|
||||
return &Engine{store: s, loki: lokiClient, analyzer: analyzer, actions: actionRunner, dispatcher: dispatcher}
|
||||
}
|
||||
|
||||
func (e *Engine) Start(ctx context.Context, interval time.Duration) {
|
||||
if interval <= 0 {
|
||||
interval = time.Minute
|
||||
}
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
e.CheckAll(ctx)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (e *Engine) CheckAll(ctx context.Context) {
|
||||
data, err := e.store.Load()
|
||||
if err != nil {
|
||||
log.Printf("load rules failed: %v", err)
|
||||
return
|
||||
}
|
||||
channels := map[string]store.AlertChannel{}
|
||||
for _, channel := range data.AlertChannels {
|
||||
channels[channel.ID] = channel
|
||||
}
|
||||
for _, rule := range data.Rules {
|
||||
if !rule.Enabled {
|
||||
continue
|
||||
}
|
||||
if err := e.checkRule(ctx, rule, channels); err != nil {
|
||||
log.Printf("rule %q check failed: %v", rule.Name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (e *Engine) checkRule(ctx context.Context, rule store.Rule, channels map[string]store.AlertChannel) error {
|
||||
window, err := time.ParseDuration(rule.Window)
|
||||
if err != nil {
|
||||
window = 5 * time.Minute
|
||||
}
|
||||
result, err := e.loki.QueryRange(ctx, rule.LogQL, window, 200)
|
||||
if err != nil {
|
||||
e.record(rule.ID, 0, err)
|
||||
return err
|
||||
}
|
||||
if result.Count < rule.Threshold {
|
||||
return e.record(rule.ID, result.Count, nil)
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
cooldown, err := time.ParseDuration(rule.Cooldown)
|
||||
if err != nil || cooldown <= 0 {
|
||||
cooldown = time.Hour
|
||||
}
|
||||
if !rule.LastAlertedAt.IsZero() && now.Sub(rule.LastAlertedAt) < cooldown {
|
||||
return e.store.UpdateRule(rule.ID, func(r *store.Rule) {
|
||||
r.LastCheckedAt = now
|
||||
r.LastMatchedAt = now
|
||||
r.LastMatchCount = result.Count
|
||||
r.SuppressedCount++
|
||||
r.LastError = ""
|
||||
})
|
||||
}
|
||||
finding, err := e.analyzer.Analyze(ctx, rule, result.Matches)
|
||||
if err != nil {
|
||||
e.record(rule.ID, result.Count, err)
|
||||
return err
|
||||
}
|
||||
selected := make([]store.AlertChannel, 0, len(rule.AlertChannels))
|
||||
for _, id := range rule.AlertChannels {
|
||||
if channel, ok := channels[id]; ok {
|
||||
selected = append(selected, channel)
|
||||
}
|
||||
}
|
||||
alertResults, alertErr := e.dispatcher.Send(ctx, rule, finding, selected)
|
||||
var alertEvidence []string
|
||||
for _, result := range alertResults {
|
||||
if result.ChannelID == "" {
|
||||
alertEvidence = append(alertEvidence, result.Detail)
|
||||
} else {
|
||||
alertEvidence = append(alertEvidence, result.ChannelID+": "+result.Detail)
|
||||
}
|
||||
}
|
||||
var actionEvidence []string
|
||||
for _, action := range rule.Actions {
|
||||
result, err := e.actions.Run(ctx, rule, action)
|
||||
if err != nil {
|
||||
log.Printf("action %q for rule %q failed: %v", action.Type, rule.Name, err)
|
||||
actionEvidence = append(actionEvidence, action.Type+": "+err.Error())
|
||||
continue
|
||||
}
|
||||
actionEvidence = append(actionEvidence, result.Action+": "+result.Detail)
|
||||
}
|
||||
incident := store.Incident{RuleID: rule.ID, RuleName: rule.Name, Severity: rule.Severity, Count: result.Count, Summary: finding.Summary, AlertResults: alertEvidence, RemediationEvidence: actionEvidence, CreatedAt: now}
|
||||
if err := e.store.AddIncident(incident); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := e.store.UpdateRule(rule.ID, func(r *store.Rule) {
|
||||
r.LastCheckedAt = now
|
||||
r.LastMatchedAt = now
|
||||
r.LastAlertedAt = now
|
||||
r.LastMatchCount = result.Count
|
||||
r.LastError = ""
|
||||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
return alertErr
|
||||
}
|
||||
|
||||
func (e *Engine) record(id string, count int, err error) error {
|
||||
return e.store.UpdateRule(id, func(r *store.Rule) {
|
||||
r.LastCheckedAt = time.Now().UTC()
|
||||
r.LastMatchCount = count
|
||||
if err != nil {
|
||||
r.LastError = err.Error()
|
||||
} else {
|
||||
r.LastError = ""
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
package rules
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/actions"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/alerts"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/analysis"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/loki"
|
||||
"gitea.wayfinderak.com/wayfinderak/log-guardian/internal/store"
|
||||
)
|
||||
|
||||
func TestCooldownSuppressesDuplicateIncident(t *testing.T) {
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
_, _ = w.Write([]byte(`{"status":"success","data":{"result":[{"stream":{},"values":[["1700000000000000000","error"]]}]}}`))
|
||||
}))
|
||||
defer ts.Close()
|
||||
|
||||
s := store.NewFileStore(t.TempDir() + "/config.json")
|
||||
if err := s.UpsertRule(store.Rule{Name: "Errors", Enabled: true, LogQL: `{service="api"}`, Threshold: 1, Window: "5m", Cooldown: "1h", Actions: []store.Action{{Type: "record_recommendation", Enabled: true}}}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
engine := NewEngine(s, loki.New(ts.URL, "", "", ""), analysis.NoopAnalyzer{}, actions.NewRunner(true), alerts.NewDispatcher())
|
||||
engine.CheckAll(t.Context())
|
||||
engine.CheckAll(t.Context())
|
||||
data, err := s.Load()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(data.Incidents) != 1 {
|
||||
t.Fatalf("expected one incident, got %d", len(data.Incidents))
|
||||
}
|
||||
if data.Rules[0].SuppressedCount != 1 {
|
||||
t.Fatalf("expected suppressed count 1, got %d", data.Rules[0].SuppressedCount)
|
||||
}
|
||||
if time.Since(data.Rules[0].LastAlertedAt) > time.Minute {
|
||||
t.Fatalf("last alerted not set: %s", data.Rules[0].LastAlertedAt)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user