diff --git a/kubecost/config.go b/kubecost/config.go index 49535f0..c51c3e7 100644 --- a/kubecost/config.go +++ b/kubecost/config.go @@ -71,6 +71,9 @@ type EmitterConfig struct { EmitKubeModelMinuteResolution bool HeartbeatExportEnabled bool DiagnosticsExportEnabled bool + HeartbeatStorageRetention time.Duration + DiagnosticsStorageRetention time.Duration + StorageCleanupInterval time.Duration EmitLegacyDateModels bool EmitKubeModel bool KubernetesResourcesRequired []string @@ -95,6 +98,9 @@ func NewEmitterConfigFromEnv(clusterUID string) *EmitterConfig { EmitKubeModelMinuteResolution: kcenv.IsMinuteMetricsEnabled(), HeartbeatExportEnabled: kcenv.IsHeartbeatExportEnabled(), DiagnosticsExportEnabled: kcenv.IsDiagnosticsExportEnabled(), + HeartbeatStorageRetention: kcenv.GetHeartbeatStorageRetention(), + DiagnosticsStorageRetention: kcenv.GetDiagnosticsStorageRetention(), + StorageCleanupInterval: kcenv.GetStorageCleanupInterval(), EmitLegacyDateModels: coreenv.IsLegacyDataModelExported(), EmitKubeModel: kcenv.IsFinOpsAgentKubeModelExported(), // Kubecost emitter requires all kubernetes resources to be enabled diff --git a/kubecost/emitter.go b/kubecost/emitter.go index 9a56af5..588ff8e 100644 --- a/kubecost/emitter.go +++ b/kubecost/emitter.go @@ -27,6 +27,7 @@ type KubecostEmitter struct { pipelineControllers *exporter.PipelineExportControllers heartbeatController ocexporter.ExportController diagController ocexporter.ExportController + storageCleaner *EventStorageCleaner diag diagnostics.DiagnosticService config *EmitterConfig @@ -123,12 +124,32 @@ func (ke *KubecostEmitter) Init(snapshot *emitter.ClusterSnapshot) error { diagnosticsExporter.Start(ke.config.ExportIntervals.DiagnosticsInterval) } + // Clean up expired heartbeat/diagnostics objects written by this agent. + storageCleaner := NewEventStorageCleaner( + bucketStore, + ke.config.AppName, + ke.config.ClusterName, + ke.config.HeartbeatStorageRetention, + ke.config.DiagnosticsStorageRetention, + ) + if storageCleaner.Enabled() { + if storageCleaner.Start(ke.config.StorageCleanupInterval) { + log.Infof( + "Started federated storage cleanup for heartbeat retention=%s diagnostics retention=%s interval=%s", + ke.config.HeartbeatStorageRetention, + ke.config.DiagnosticsStorageRetention, + ke.config.StorageCleanupInterval, + ) + } + } + // initialize emitter's internal state ke.dataSource = dataSource ke.costModel = costModel ke.pipelineControllers = pipelineControllers ke.heartbeatController = agentHeartbeat ke.diagController = diagnosticsExporter + ke.storageCleaner = storageCleaner return nil } diff --git a/kubecost/env/emitterenv.go b/kubecost/env/emitterenv.go index 48eb808..c0103d2 100644 --- a/kubecost/env/emitterenv.go +++ b/kubecost/env/emitterenv.go @@ -21,8 +21,15 @@ const ( KubeModelExportIntervalEnvVar = "KUBEMODEL_EXPORT_INTERVAL" HeartbeatExportIntervalEnvVar = "HEARTBEAT_EXPORT_INTERVAL" DiagnosticsExportIntervalEnvVar = "DIAGNOSTICS_EXPORT_INTERVAL" + HeartbeatStorageRetentionEnvVar = "HEARTBEAT_STORAGE_RETENTION" + DiagnosticsStorageRetentionEnvVar = "DIAGNOSTICS_STORAGE_RETENTION" + StorageCleanupIntervalEnvVar = "STORAGE_CLEANUP_INTERVAL" StreamingExportEnabledEnvVar = "STREAMING_EXPORT_ENABLED" StreamingExportCompressionLevelEnvVar = "STREAMING_EXPORT_COMPRESSION_LEVEL" + + DefaultHeartbeatStorageRetention = 7 * 24 * time.Hour + DefaultDiagnosticsStorageRetention = 7 * 24 * time.Hour + DefaultStorageCleanupInterval = 1 * time.Hour ) // IsMinuteMetricsEnabled returns true if the 10m resolution emitter for kubecost @@ -71,6 +78,36 @@ func IsDiagnosticsExportEnabled() bool { return coreenv.GetBool(DiagnosticsExportEnabledEnvVar, true) } +// GetHeartbeatStorageRetention returns how long heartbeat objects are retained +// in federated storage before the agent deletes them. A value of 0 disables cleanup. +func GetHeartbeatStorageRetention() time.Duration { + return getStorageRetention(HeartbeatStorageRetentionEnvVar, DefaultHeartbeatStorageRetention) +} + +// GetDiagnosticsStorageRetention returns how long diagnostics objects are retained +// in federated storage before the agent deletes them. A value of 0 disables cleanup. +func GetDiagnosticsStorageRetention() time.Duration { + return getStorageRetention(DiagnosticsStorageRetentionEnvVar, DefaultDiagnosticsStorageRetention) +} + +// GetStorageCleanupInterval returns how often the agent scans federated storage +// for expired heartbeat and diagnostics objects. +func GetStorageCleanupInterval() time.Duration { + return coreenv.GetDuration(StorageCleanupIntervalEnvVar, DefaultStorageCleanupInterval) +} + +func getStorageRetention(envVar string, defaultValue time.Duration) time.Duration { + raw := coreenv.Get(envVar, "") + if raw == "" { + return defaultValue + } + // Allow a bare "0" to disable cleanup; duration parsers require a unit. + if raw == "0" { + return 0 + } + return coreenv.GetDuration(envVar, defaultValue) +} + // IsStreamingExportEnabled returns true if the bingen pipeline exporters should use a streaming io.Writer // when exporting data, as opposed to encoding a []byte, then uploading. func IsStreamingExportEnabled() bool { diff --git a/kubecost/storagecleanup.go b/kubecost/storagecleanup.go new file mode 100644 index 0000000..c53efb4 --- /dev/null +++ b/kubecost/storagecleanup.go @@ -0,0 +1,186 @@ +package kubecost + +import ( + "fmt" + "path" + "strings" + "time" + + "github.com/opencost/opencost/core/pkg/diagnostics" + "github.com/opencost/opencost/core/pkg/exporter/pathing" + "github.com/opencost/opencost/core/pkg/heartbeat" + "github.com/opencost/opencost/core/pkg/log" + "github.com/opencost/opencost/core/pkg/storage" + "github.com/opencost/opencost/core/pkg/util/atomic" +) + +// EventStorageCleaner periodically removes expired heartbeat and diagnostics +// objects from federated storage for a single cluster prefix. +type EventStorageCleaner struct { + store storage.Storage + appName string + clusterName string + heartbeatRetention time.Duration + diagnosticsRetention time.Duration + runState atomic.AtomicRunState +} + +// NewEventStorageCleaner creates a cleaner scoped to the agent's own +// app/cluster heartbeat and diagnostics prefixes. +func NewEventStorageCleaner( + store storage.Storage, + appName string, + clusterName string, + heartbeatRetention time.Duration, + diagnosticsRetention time.Duration, +) *EventStorageCleaner { + return &EventStorageCleaner{ + store: store, + appName: appName, + clusterName: clusterName, + heartbeatRetention: heartbeatRetention, + diagnosticsRetention: diagnosticsRetention, + } +} + +// Enabled reports whether either retention window is configured. +func (c *EventStorageCleaner) Enabled() bool { + return c.heartbeatRetention > 0 || c.diagnosticsRetention > 0 +} + +// Start begins periodic cleanup on the provided interval. Returns false if the +// cleaner is already running or no retention windows are configured. +func (c *EventStorageCleaner) Start(interval time.Duration) bool { + if !c.Enabled() { + return false + } + if interval <= 0 { + log.Warnf("EventStorageCleaner: invalid cleanup interval %s; cleanup will not start", interval) + return false + } + + c.runState.WaitForReset() + if !c.runState.Start() { + return false + } + + go func() { + // Run once immediately so upgrades begin reclaiming space without waiting + // for the first interval to elapse. + c.Cleanup() + + for { + select { + case <-c.runState.OnStop(): + c.runState.Reset() + return + case <-time.After(interval): + c.Cleanup() + } + } + }() + + return true +} + +// Stop halts the cleanup loop. +func (c *EventStorageCleaner) Stop() { + c.runState.Stop() +} + +// Cleanup deletes expired heartbeat and diagnostics objects for this cluster. +// Individual delete failures are logged and do not stop processing. +func (c *EventStorageCleaner) Cleanup() { + now := time.Now().UTC() + + if c.heartbeatRetention > 0 { + dir := path.Join(c.appName, c.clusterName, heartbeat.HeartbeatEventName) + deleted, skipped, errCount := cleanupExpiredObjects(c.store, dir, now.Add(-c.heartbeatRetention)) + logCleanupSummary(heartbeat.HeartbeatEventName, dir, deleted, skipped, errCount) + } + + if c.diagnosticsRetention > 0 { + dir := path.Join(c.appName, c.clusterName, diagnostics.DiagnosticsEventName) + deleted, skipped, errCount := cleanupExpiredObjects(c.store, dir, now.Add(-c.diagnosticsRetention)) + logCleanupSummary(diagnostics.DiagnosticsEventName, dir, deleted, skipped, errCount) + } +} + +func logCleanupSummary(valueType, dir string, deleted, skipped, errCount int) { + if deleted == 0 && errCount == 0 { + log.Debugf("EventStorageCleaner: %s cleanup complete for %s (deleted=0, skipped=%d)", valueType, dir, skipped) + return + } + log.Infof("EventStorageCleaner: %s cleanup complete for %s (deleted=%d, skipped=%d, errors=%d)", valueType, dir, deleted, skipped, errCount) +} + +// cleanupExpiredObjects lists objects under dir and removes those whose age is +// older than cutoff. Age prefers the event timestamp encoded in the filename +// (YYYYMMDDHHmmss.json); ModTime is used when the filename cannot be parsed. +func cleanupExpiredObjects(store storage.Storage, dir string, cutoff time.Time) (deleted, skipped, errCount int) { + files, err := store.List(dir) + if err != nil { + log.Errorf("EventStorageCleaner: failed to list %s: %v", dir, err) + return 0, 0, 1 + } + + for _, file := range files { + if file == nil || file.Name == "" || strings.HasSuffix(file.Name, "/") { + skipped++ + continue + } + + objectAge, ok := objectAge(file) + if !ok { + log.Debugf("EventStorageCleaner: skipping unparseable object %s/%s", dir, file.Name) + skipped++ + continue + } + + if !objectAge.Before(cutoff) { + skipped++ + continue + } + + objectPath := path.Join(dir, file.Name) + if err := store.Remove(objectPath); err != nil { + log.Errorf("EventStorageCleaner: failed to delete %s: %v", objectPath, err) + errCount++ + continue + } + deleted++ + } + + return deleted, skipped, errCount +} + +func objectAge(file *storage.StorageInfo) (time.Time, bool) { + if ts, err := parseEventFilenameTimestamp(file.Name); err == nil { + return ts, true + } + if !file.ModTime.IsZero() { + return file.ModTime.UTC(), true + } + return time.Time{}, false +} + +func parseEventFilenameTimestamp(name string) (time.Time, error) { + base := path.Base(name) + // Expected: YYYYMMDDHHmmss.json or optionally prefix.YYYYMMDDHHmmss.json + parts := strings.Split(base, ".") + if len(parts) < 2 { + return time.Time{}, fmt.Errorf("unexpected filename format: %s", name) + } + + // Prefer the segment immediately before the extension when present. + timestampPart := parts[0] + if len(parts) >= 2 { + timestampPart = parts[len(parts)-2] + } + + ts, err := time.Parse(pathing.EventStorageTimeFormat, timestampPart) + if err != nil { + return time.Time{}, err + } + return ts.UTC(), nil +} diff --git a/kubecost/storagecleanup_test.go b/kubecost/storagecleanup_test.go new file mode 100644 index 0000000..ecc9bcf --- /dev/null +++ b/kubecost/storagecleanup_test.go @@ -0,0 +1,137 @@ +package kubecost + +import ( + "os" + "path" + "path/filepath" + "testing" + "time" + + "github.com/opencost/opencost/core/pkg/diagnostics" + "github.com/opencost/opencost/core/pkg/exporter/pathing" + "github.com/opencost/opencost/core/pkg/heartbeat" + "github.com/opencost/opencost/core/pkg/storage" +) + +var testFileContent = []byte(`{"test":true}`) + +func writeEventFile(t *testing.T, store storage.Storage, baseDir, dir, fileName string, modTime time.Time) string { + t.Helper() + + objectPath := path.Join(dir, fileName) + if err := store.Write(objectPath, testFileContent); err != nil { + t.Fatalf("failed to write %s: %v", objectPath, err) + } + + fullPath := filepath.Join(baseDir, objectPath) + if err := os.Chtimes(fullPath, modTime, modTime); err != nil { + t.Fatalf("failed to set modtime for %s: %v", fullPath, err) + } + + return objectPath +} + +func TestEventStorageCleaner_Cleanup_DeletesExpiredObjects(t *testing.T) { + baseDir := t.TempDir() + store := storage.NewFileStorage(baseDir) + + appName := "finops-agent" + clusterName := "cluster-1" + retention := 7 * 24 * time.Hour + now := time.Now().UTC() + + heartbeatDir := path.Join(appName, clusterName, heartbeat.HeartbeatEventName) + diagnosticsDir := path.Join(appName, clusterName, diagnostics.DiagnosticsEventName) + otherClusterDir := path.Join(appName, "other-cluster", heartbeat.HeartbeatEventName) + + oldHeartbeat := writeEventFile(t, store, baseDir, heartbeatDir, now.Add(-8*24*time.Hour).Format(pathing.EventStorageTimeFormat)+".json", now.Add(-8*24*time.Hour)) + recentHeartbeat := writeEventFile(t, store, baseDir, heartbeatDir, now.Add(-1*time.Hour).Format(pathing.EventStorageTimeFormat)+".json", now.Add(-1*time.Hour)) + oldDiagnostics := writeEventFile(t, store, baseDir, diagnosticsDir, now.Add(-10*24*time.Hour).Format(pathing.EventStorageTimeFormat)+".json", now.Add(-10*24*time.Hour)) + recentDiagnostics := writeEventFile(t, store, baseDir, diagnosticsDir, now.Add(-2*time.Hour).Format(pathing.EventStorageTimeFormat)+".json", now.Add(-2*time.Hour)) + otherClusterOld := writeEventFile(t, store, baseDir, otherClusterDir, now.Add(-30*24*time.Hour).Format(pathing.EventStorageTimeFormat)+".json", now.Add(-30*24*time.Hour)) + + cleaner := NewEventStorageCleaner(store, appName, clusterName, retention, retention) + cleaner.Cleanup() + + assertExists(t, store, oldHeartbeat, false) + assertExists(t, store, recentHeartbeat, true) + assertExists(t, store, oldDiagnostics, false) + assertExists(t, store, recentDiagnostics, true) + assertExists(t, store, otherClusterOld, true) +} + +func TestEventStorageCleaner_Cleanup_RetentionZeroDisablesPrefix(t *testing.T) { + baseDir := t.TempDir() + store := storage.NewFileStorage(baseDir) + + appName := "finops-agent" + clusterName := "cluster-1" + now := time.Now().UTC() + + heartbeatDir := path.Join(appName, clusterName, heartbeat.HeartbeatEventName) + diagnosticsDir := path.Join(appName, clusterName, diagnostics.DiagnosticsEventName) + + oldHeartbeat := writeEventFile(t, store, baseDir, heartbeatDir, now.Add(-30*24*time.Hour).Format(pathing.EventStorageTimeFormat)+".json", now.Add(-30*24*time.Hour)) + oldDiagnostics := writeEventFile(t, store, baseDir, diagnosticsDir, now.Add(-30*24*time.Hour).Format(pathing.EventStorageTimeFormat)+".json", now.Add(-30*24*time.Hour)) + + cleaner := NewEventStorageCleaner(store, appName, clusterName, 0, 7*24*time.Hour) + cleaner.Cleanup() + + assertExists(t, store, oldHeartbeat, true) + assertExists(t, store, oldDiagnostics, false) +} + +func TestEventStorageCleaner_Cleanup_FallsBackToModTime(t *testing.T) { + baseDir := t.TempDir() + store := storage.NewFileStorage(baseDir) + + appName := "finops-agent" + clusterName := "cluster-1" + retention := 7 * 24 * time.Hour + now := time.Now().UTC() + + heartbeatDir := path.Join(appName, clusterName, heartbeat.HeartbeatEventName) + oldUnparseable := writeEventFile(t, store, baseDir, heartbeatDir, "not-a-timestamp.json", now.Add(-10*24*time.Hour)) + recentUnparseable := writeEventFile(t, store, baseDir, heartbeatDir, "also-not-a-timestamp.json", now.Add(-1*time.Hour)) + + cleaner := NewEventStorageCleaner(store, appName, clusterName, retention, 0) + cleaner.Cleanup() + + assertExists(t, store, oldUnparseable, false) + assertExists(t, store, recentUnparseable, true) +} + +func TestParseEventFilenameTimestamp(t *testing.T) { + ts, err := parseEventFilenameTimestamp("20251022153000.json") + if err != nil { + t.Fatalf("expected parse success, got %v", err) + } + expected := time.Date(2025, 10, 22, 15, 30, 0, 0, time.UTC) + if !ts.Equal(expected) { + t.Fatalf("expected %s, got %s", expected, ts) + } + + ts, err = parseEventFilenameTimestamp("prefix.20251022153000.json") + if err != nil { + t.Fatalf("expected prefixed parse success, got %v", err) + } + if !ts.Equal(expected) { + t.Fatalf("expected %s, got %s", expected, ts) + } + + if _, err := parseEventFilenameTimestamp("badname"); err == nil { + t.Fatal("expected parse failure for badname") + } +} + +func assertExists(t *testing.T, store storage.Storage, objectPath string, want bool) { + t.Helper() + + exists, err := store.Exists(objectPath) + if err != nil { + t.Fatalf("Exists(%s) failed: %v", objectPath, err) + } + if exists != want { + t.Fatalf("Exists(%s)=%v, want %v", objectPath, exists, want) + } +}