feat: add monitoring retention controls

This commit is contained in:
sagitchu
2026-04-28 11:41:07 +08:00
parent e8d5687419
commit 023be27287
15 changed files with 430 additions and 13 deletions
+17 -3
View File
@@ -6,6 +6,7 @@ import (
"sync"
"time"
"go-backend/internal/monitoring"
"go-backend/internal/store/model"
"go-backend/internal/store/repo"
)
@@ -31,7 +32,6 @@ type IngestionService struct {
nodeBuffer []*model.NodeMetric
nodeBufferMu sync.Mutex
flushInterval time.Duration
retentionDays int
}
func NewIngestionService(repo *repo.Repository) *IngestionService {
@@ -39,7 +39,6 @@ func NewIngestionService(repo *repo.Repository) *IngestionService {
repo: repo,
nodeBuffer: make([]*model.NodeMetric, 0, 500),
flushInterval: 30 * time.Second,
retentionDays: 7,
}
}
@@ -111,7 +110,22 @@ func (s *IngestionService) flushNodeMetrics() {
}
func (s *IngestionService) pruneMetrics() {
cutoff := time.Now().Add(-time.Duration(s.retentionDays) * 24 * time.Hour).UnixMilli()
s.pruneMetricsAt(time.Now())
}
func (s *IngestionService) retentionDaysFromConfig() int {
if s == nil || s.repo == nil {
return monitoring.DefaultMonitorRetentionDays
}
cfg, err := s.repo.GetConfigsByNames([]string{monitoring.ConfigMonitorRetentionDays})
if err != nil {
return monitoring.DefaultMonitorRetentionDays
}
return monitoring.MonitoringRetentionDaysFromConfigMap(cfg)
}
func (s *IngestionService) pruneMetricsAt(now time.Time) {
cutoff := now.Add(-time.Duration(s.retentionDaysFromConfig()) * 24 * time.Hour).UnixMilli()
if s.repo == nil {
return
}
+34 -1
View File
@@ -5,6 +5,7 @@ import (
"testing"
"time"
"go-backend/internal/store/model"
"go-backend/internal/store/repo"
)
@@ -215,7 +216,6 @@ func TestPruneMetrics(t *testing.T) {
defer r.Close()
svc := NewIngestionService(r)
svc.retentionDays = 1
info := SystemInfo{CPUUsage: 50.0, MemoryUsage: 60.0, DiskUsage: 30.0}
@@ -233,6 +233,39 @@ func TestPruneMetrics(t *testing.T) {
}
}
func TestPruneMetricsUsesConfiguredRetentionDays(t *testing.T) {
r, err := repo.Open(":memory:")
if err != nil {
t.Fatalf("open repo: %v", err)
}
defer r.Close()
now := time.Now().UnixMilli()
if err := r.UpsertConfig("monitor_retention_days", "2", now); err != nil {
t.Fatalf("upsert retention config: %v", err)
}
oldMetric := &model.NodeMetric{NodeID: 1, Timestamp: now - int64(3*24*time.Hour/time.Millisecond), CPUUsage: 10}
newMetric := &model.NodeMetric{NodeID: 1, Timestamp: now - int64(1*24*time.Hour/time.Millisecond), CPUUsage: 20}
if err := r.InsertNodeMetric(oldMetric); err != nil {
t.Fatalf("insert old metric: %v", err)
}
if err := r.InsertNodeMetric(newMetric); err != nil {
t.Fatalf("insert new metric: %v", err)
}
svc := NewIngestionService(r)
svc.pruneMetricsAt(time.UnixMilli(now))
metrics, err := r.GetNodeMetrics(1, now-int64(4*24*time.Hour/time.Millisecond), now+1000)
if err != nil {
t.Fatalf("get node metrics: %v", err)
}
if len(metrics) != 1 || metrics[0].CPUUsage != 20 {
t.Fatalf("expected only newer metric to remain, got %#v", metrics)
}
}
func TestMultipleNodes(t *testing.T) {
r, err := repo.Open(":memory:")
if err != nil {