diff --git a/pkg/metricstore/archive.go b/pkg/metricstore/archive.go index 3938ac82..fcf3124a 100644 --- a/pkg/metricstore/archive.go +++ b/pkg/metricstore/archive.go @@ -94,16 +94,32 @@ var ErrNoNewArchiveData error = errors.New("all data already archived") // CleanupCheckpoints deletes or archives all checkpoint files older than `from`. // When archiving, consolidates all hosts per cluster into a single Parquet file. +// +// Hosts currently in use by running jobs (per the MemoryStore's NodeProvider) +// are skipped entirely; their files are handled by a later cycle once the job +// has ended. If the provider query fails, the whole cycle is aborted: deleting +// without knowing which nodes are in use risks data loss. Without a provider, +// everything older than `from` is cleaned (old behavior). func CleanupCheckpoints(checkpointsDir, cleanupDir string, from int64, deleteInstead bool) (int, error) { - if deleteInstead { - return deleteCheckpoints(checkpointsDir, from) + var usedNodes map[string][]string + if ms := GetMemoryStore(); ms != nil && ms.nodeProvider != nil { + un, err := ms.nodeProvider.GetUsedNodes(from) + if err != nil { + return 0, fmt.Errorf("aborting cleanup, could not determine nodes in use: %w", err) + } + usedNodes = un } - return archiveCheckpoints(checkpointsDir, cleanupDir, from) + if deleteInstead { + return deleteCheckpoints(checkpointsDir, from, usedNodes) + } + + return archiveCheckpoints(checkpointsDir, cleanupDir, from, usedNodes) } // deleteCheckpoints removes checkpoint files older than `from` across all clusters/hosts. -func deleteCheckpoints(checkpointsDir string, from int64) (int, error) { +// Hosts present in usedNodes are skipped. +func deleteCheckpoints(checkpointsDir string, from int64, usedNodes map[string][]string) (int, error) { entries1, err := os.ReadDir(checkpointsDir) if err != nil { return 0, err @@ -157,6 +173,10 @@ func deleteCheckpoints(checkpointsDir string, from int64) (int, error) { } for _, de2 := range entries2 { + if isNodeUsed(usedNodes, de1.Name(), de2.Name()) { + cclog.Debugf("[METRICSTORE]> skipping delete for %s/%s: node in use by running job", de1.Name(), de2.Name()) + continue + } work <- workItem{ dir: filepath.Join(checkpointsDir, de1.Name(), de2.Name()), cluster: de1.Name(), @@ -182,7 +202,7 @@ func deleteCheckpoints(checkpointsDir string, from int64) (int, error) { // Workers load checkpoint files from disk and send CheckpointFile trees on a // back-pressured channel. The main thread streams each tree directly to Parquet // rows without materializing all rows in memory. -func archiveCheckpoints(checkpointsDir, cleanupDir string, from int64) (int, error) { +func archiveCheckpoints(checkpointsDir, cleanupDir string, from int64, usedNodes map[string][]string) (int, error) { cclog.Info("[METRICSTORE]> start archiving checkpoints to parquet") startTime := time.Now() @@ -252,6 +272,10 @@ func archiveCheckpoints(checkpointsDir, cleanupDir string, from int64) (int, err if !hostEntry.IsDir() { continue } + if isNodeUsed(usedNodes, cluster, hostEntry.Name()) { + cclog.Debugf("[METRICSTORE]> skipping archive for %s/%s: node in use by running job", cluster, hostEntry.Name()) + continue + } dir := filepath.Join(checkpointsDir, cluster, hostEntry.Name()) work <- struct { dir, host string diff --git a/pkg/metricstore/nodeprovider_test.go b/pkg/metricstore/nodeprovider_test.go index 7071b167..fb834d01 100644 --- a/pkg/metricstore/nodeprovider_test.go +++ b/pkg/metricstore/nodeprovider_test.go @@ -161,3 +161,104 @@ func TestFromCheckpointProviderErrorFallsBack(t *testing.T) { t.Fatalf("expected fallback to cutoff load (4 files), got %d", n) } } + +func TestDeleteCheckpointsSkipsUsedNodes(t *testing.T) { + oldWorkers := Keys.NumWorkers + Keys.NumWorkers = 2 + t.Cleanup(func() { Keys.NumWorkers = oldWorkers }) + + dir := t.TempDir() + writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node001"), 1000) + writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node002"), 1000) + + used := map[string][]string{"fritz": {"node001"}} + n, err := deleteCheckpoints(dir, 3000, used) + if err != nil { + t.Fatal(err) + } + if n != 1 { + t.Fatalf("expected 1 file deleted, got %d", n) + } + if _, err := os.Stat(filepath.Join(dir, "fritz", "node001", "1000.json")); err != nil { + t.Errorf("used node's checkpoint must survive: %v", err) + } + if _, err := os.Stat(filepath.Join(dir, "fritz", "node002", "1000.json")); !os.IsNotExist(err) { + t.Error("unused node's checkpoint must be deleted") + } +} + +func TestArchiveCheckpointsSkipsUsedNodes(t *testing.T) { + oldWorkers := Keys.NumWorkers + Keys.NumWorkers = 2 + t.Cleanup(func() { Keys.NumWorkers = oldWorkers }) + + cpDir := t.TempDir() + archiveDir := t.TempDir() + writeTestCheckpoint(t, filepath.Join(cpDir, "fritz", "node001"), 1000) + writeTestCheckpoint(t, filepath.Join(cpDir, "fritz", "node002"), 1000) + + used := map[string][]string{"fritz": {"node001"}} + n, err := archiveCheckpoints(cpDir, archiveDir, 3000, used) + if err != nil { + t.Fatal(err) + } + if n != 1 { + t.Fatalf("expected 1 file archived, got %d", n) + } + if _, err := os.Stat(filepath.Join(cpDir, "fritz", "node001", "1000.json")); err != nil { + t.Errorf("used node's checkpoint must survive: %v", err) + } + if _, err := os.Stat(filepath.Join(cpDir, "fritz", "node002", "1000.json")); !os.IsNotExist(err) { + t.Error("unused node's checkpoint must be removed after archiving") + } + if _, err := os.Stat(filepath.Join(archiveDir, "fritz", "3000.parquet")); err != nil { + t.Errorf("parquet archive must exist: %v", err) + } +} + +func TestCleanupCheckpointsAbortsOnProviderError(t *testing.T) { + oldWorkers := Keys.NumWorkers + Keys.NumWorkers = 2 + t.Cleanup(func() { Keys.NumWorkers = oldWorkers }) + + oldMS := msInstance + msInstance = newTestStore() + msInstance.SetNodeProvider(&fakeNodeProvider{err: errTestProvider}) + t.Cleanup(func() { msInstance = oldMS }) + + dir := t.TempDir() + writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node001"), 1000) + + if _, err := CleanupCheckpoints(dir, "", 3000, true); err == nil { + t.Fatal("expected error when GetUsedNodes fails") + } + if _, err := os.Stat(filepath.Join(dir, "fritz", "node001", "1000.json")); err != nil { + t.Errorf("no files may be deleted when provider errors: %v", err) + } +} + +func TestCleanupCheckpointsUsedNodesSurvive(t *testing.T) { + oldWorkers := Keys.NumWorkers + Keys.NumWorkers = 2 + t.Cleanup(func() { Keys.NumWorkers = oldWorkers }) + + oldMS := msInstance + msInstance = newTestStore() + msInstance.SetNodeProvider(&fakeNodeProvider{nodes: map[string][]string{"fritz": {"node001"}}}) + t.Cleanup(func() { msInstance = oldMS }) + + dir := t.TempDir() + writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node001"), 1000) + writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node002"), 1000) + + n, err := CleanupCheckpoints(dir, "", 3000, true) + if err != nil { + t.Fatal(err) + } + if n != 1 { + t.Fatalf("expected 1 file deleted, got %d", n) + } + if _, err := os.Stat(filepath.Join(dir, "fritz", "node001", "1000.json")); err != nil { + t.Errorf("used node's checkpoint must survive: %v", err) + } +}