From 2017c95a7a4c8d1893ab237752a0c7eecf1245a0 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Thu, 17 Sep 2026 08:59:24 +0000 Subject: [PATCH] feat(monitors): sweep metric monitors per server with incidents per breach --- server/cmd/main.go | 4 + server/internal/metricsched/scheduler.go | 1 + server/internal/services/metricsweep.go | 144 +++++++++++++++++++ server/internal/services/mfa_scoped_test.go | 10 ++ server/internal/services/migrate_instance.go | 1 + 5 files changed, 160 insertions(+) create mode 100644 server/internal/services/metricsweep.go diff --git a/server/cmd/main.go b/server/cmd/main.go index 4f86806..e6da41b 100644 --- a/server/cmd/main.go +++ b/server/cmd/main.go @@ -175,6 +175,10 @@ func runSchemaSetup() { log.Printf("warning: failed to ensure monitor sample indexes: %v", err) } + if err := services.EnsureMonitorServerStateIndexes(); err != nil { + log.Printf("warning: %v", err) + } + if err := services.EnsureVulnIndexes(); err != nil { log.Printf("warning: failed to ensure vuln indexes: %v", err) } diff --git a/server/internal/metricsched/scheduler.go b/server/internal/metricsched/scheduler.go index 24da575..94c9279 100644 --- a/server/internal/metricsched/scheduler.go +++ b/server/internal/metricsched/scheduler.go @@ -37,4 +37,5 @@ func sweep(now time.Time) { } }() services.SweepHeartbeats(now) + services.SweepMetricMonitors(now) } diff --git a/server/internal/services/metricsweep.go b/server/internal/services/metricsweep.go new file mode 100644 index 0000000..5e44be9 --- /dev/null +++ b/server/internal/services/metricsweep.go @@ -0,0 +1,144 @@ +package services + +import ( + "context" + "log" + "time" + + "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/db" + "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/models" + "go.mongodb.org/mongo-driver/v2/bson" + "go.mongodb.org/mongo-driver/v2/mongo" + "go.mongodb.org/mongo-driver/v2/mongo/options" +) + +// EnsureMonitorServerStateIndexes declares the one key every read and upsert +// uses. Unique, so a double sweep across a leader handover cannot fork a +// server's state into two documents. +func EnsureMonitorServerStateIndexes() error { + _, err := db.Col("monitor_server_states").Indexes().CreateOne(context.Background(), mongo.IndexModel{ + Keys: bson.D{{Key: "monitor_id", Value: 1}, {Key: "server_id", Value: 1}}, + Options: options.Index().SetUnique(true), + }) + if err != nil { + log.Printf("warning: monitor_server_states indexes: %v", err) + } + return nil +} + +func ListMonitorServerStates(instanceID, monitorID string) ([]models.MonitorServerState, error) { + ctx, cancel := monCtx() + defer cancel() + cur, err := db.Col("monitor_server_states").Find(ctx, bson.M{"instance_id": instanceID, "monitor_id": monitorID}) + if err != nil { + return nil, err + } + out := []models.MonitorServerState{} + if err := cur.All(ctx, &out); err != nil { + return nil, err + } + return out, nil +} + +// ponytail: every metric monitor re-lists its servers each sweep; cache the +// fleet per instance per sweep if large fleets show up in the sweep log. +func SweepMetricMonitors(now time.Time) { + ctx, cancel := context.WithTimeout(context.Background(), 25*time.Second) + defer cancel() + cur, err := db.Col("monitors").Find(ctx, bson.M{"type": models.MonitorMetric, "enabled": true}) + if err != nil { + log.Printf("metrics: list: %v", err) + return + } + var monitors []models.Monitor + if err := cur.All(ctx, &monitors); err != nil { + log.Printf("metrics: decode: %v", err) + return + } + for i := range monitors { + sweepOneMetric(ctx, &monitors[i], now) + } +} + +func sweepOneMetric(ctx context.Context, m *models.Monitor, now time.Time) { + defer func() { + if r := recover(); r != nil { + log.Printf("metrics: monitor %s panic: %v", m.MonitorID, r) + } + }() + + servers, err := ListServersFiltered(m.InstanceID, m.Target.Selector) + if err != nil { + log.Printf("metrics: servers for %s: %v", m.MonitorID, err) + return + } + prevStates, err := ListMonitorServerStates(m.InstanceID, m.MonitorID) + if err != nil { + log.Printf("metrics: states for %s: %v", m.MonitorID, err) + return + } + prevByServer := map[string]*models.MonitorServerState{} + for i := range prevStates { + prevByServer[prevStates[i].ServerID] = &prevStates[i] + } + + col := db.Col("monitor_server_states") + current := make([]models.MonitorServerState, 0, len(servers)) + matched := map[string]bool{} + for _, srv := range servers { + matched[srv.ServerID] = true + prev := prevByServer[srv.ServerID] + + var wls []models.Workload + if metricNeedsWorkloads(m.Target.Metric) { + if sw, err := GetWorkloads(m.InstanceID, srv.ServerID); err == nil && sw != nil { + wls = sw.Workloads + } + } + breach, value, msg, ok := EvaluateMetric(m.Target, srv, wls, now) + if !ok { + if prev != nil { + current = append(current, *prev) + } + continue + } + + status, since := nextServerState(prev, breach, m.ForSec, now) + st := models.MonitorServerState{ + InstanceID: m.InstanceID, MonitorID: m.MonitorID, ServerID: srv.ServerID, + Status: status, BreachSince: since, Value: value, Message: msg, UpdatedAt: now, + } + if _, err := col.ReplaceOne(ctx, bson.M{"monitor_id": m.MonitorID, "server_id": srv.ServerID}, st, + options.Replace().SetUpsert(true)); err != nil { + log.Printf("metrics: save state %s/%s: %v", m.MonitorID, srv.ServerID, err) + continue + } + current = append(current, st) + + // A server seen for the first time has no previous status to leave, so + // it can open an incident but never announce a recovery. + prevStatus := models.StatusPending + if prev != nil { + prevStatus = prev.Status + } + applyTransition(ctx, m, srv.ServerID, srv.Hostname, prevStatus, status, msg, now) + } + + // Servers that left the selector, or the fleet, did not recover: their + // incidents close quietly and their state goes. + for _, prev := range prevStates { + if matched[prev.ServerID] { + continue + } + if prev.Status == models.StatusDown { + resolveIncident(ctx, m, prev.ServerID, now) + } + col.DeleteOne(ctx, bson.M{"monitor_id": m.MonitorID, "server_id": prev.ServerID}) + } + + status, msg := rollupParent(current) + db.Col("monitors").UpdateOne(ctx, bson.M{"monitor_id": m.MonitorID}, bson.M{"$set": bson.M{ + "state.status": status, "state.message": msg, "state.last_check_at": now, + }}) + recordSample(ctx, m, status != models.StatusDown, 0, now) +} diff --git a/server/internal/services/mfa_scoped_test.go b/server/internal/services/mfa_scoped_test.go index a50ec60..c670d46 100644 --- a/server/internal/services/mfa_scoped_test.go +++ b/server/internal/services/mfa_scoped_test.go @@ -19,3 +19,13 @@ func TestMFACollectionsAreScoped(t *testing.T) { } } } + +// Per-server metric state is tenant data and must be purged with its instance. +func TestMonitorServerStatesAreScoped(t *testing.T) { + for _, got := range ScopedCollections { + if got == "monitor_server_states" { + return + } + } + t.Error("monitor_server_states is not in ScopedCollections") +} diff --git a/server/internal/services/migrate_instance.go b/server/internal/services/migrate_instance.go index aa1825e..25a4b46 100644 --- a/server/internal/services/migrate_instance.go +++ b/server/internal/services/migrate_instance.go @@ -37,6 +37,7 @@ var ScopedCollections = []string{ "incidents", "monitor_rollups", "monitor_samples", + "monitor_server_states", "notification_channels", "console_sessions", "audit_logs",