From 7e809c9991cf65cc5672c441ca02912c93a1f33d Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Mon, 5 Oct 2026 12:02:53 +0800 Subject: [PATCH] admin: stop leaking cancelled maintenance task files in -dataDir/tasks (#11597) * admin: delete persisted state when scan cancels pending tasks Each detection cycle cancels every pending task of a type before re-detecting it, and the cancel path saved the cancelled task back to disk. Nothing ever removed those files, so -dataDir/tasks gained one orphaned .pb per candidate volume per scan cycle. Cancelled is terminal, so drop the file the same way CompleteTask does for completed/failed tasks. The cancelled entry stays in memory for the UI until the next purge. Refs #11595 * admin: delete persisted state when CancelTask cancels a pending task The manual cancel path only updated memory, leaving the pending .pb on disk where a restart would resurrect the cancelled task as pending and the file would linger until then. Delete it like the scan-cycle cancel path now does. * admin: count cancelled tasks toward task retention cleanup CleanupOldTasks and ConfigPersistence.CleanupCompletedTasks only filtered completed/failed tasks, so cancelled entries were exempt from retention in both memory and on disk. Treat all terminal states alike; nil CompletedAt entries also count and sort last, so they are pruned first. * admin: run task file retention in the periodic cleanup loop cleanupCompletedTasks had no callers, so the on-disk retention bound never ran during uptime. Invoke it from performCleanup alongside the in-memory CleanupOldTasks sweep. * admin: guard task state writes against stale saves and failed deletes saveTaskState runs after mq.mutex is released, so the task may have gone terminal in between; a delayed pending save could then recreate the file a cancel just deleted and resurrect the task on restart. Skip saving non-terminal snapshots once the live task is terminal or gone. If a cancel file removal fails, fall back to writing the cancelled snapshot so the file is terminal rather than pending. deleteTaskState now returns its error, and CancelTask captures task.Status while still holding the queue lock. * admin: serialize task file check+write against cancel deletes The saveTaskState guard still had a check-then-write window: a pending snapshot could pass the terminal check before a cancel deleted the file, then write it back after. A persistMu on the queue now covers the check+save and the cancel paths' delete (with its terminal-state fallback), so the two cannot interleave for the same task. --- weed/admin/dash/config_persistence.go | 6 +- .../dash/maintenance_task_persistence_test.go | 108 ++++++++++++++++++ weed/admin/maintenance/maintenance_manager.go | 60 ++++++---- weed/admin/maintenance/maintenance_queue.go | 71 +++++++++--- weed/admin/maintenance/maintenance_types.go | 4 + 5 files changed, 207 insertions(+), 42 deletions(-) create mode 100644 weed/admin/dash/maintenance_task_persistence_test.go diff --git a/weed/admin/dash/config_persistence.go b/weed/admin/dash/config_persistence.go index a4ef90c38..f0b3e0e85 100644 --- a/weed/admin/dash/config_persistence.go +++ b/weed/admin/dash/config_persistence.go @@ -895,7 +895,7 @@ func (cp *ConfigPersistence) ListTaskDetails() ([]string, error) { return taskIDs, nil } -// CleanupCompletedTasks removes old completed tasks beyond the retention limit +// CleanupCompletedTasks removes old terminal task files beyond the retention limit func (cp *ConfigPersistence) CleanupCompletedTasks() error { cp.tasksMu.Lock() defer cp.tasksMu.Unlock() @@ -914,10 +914,10 @@ func (cp *ConfigPersistence) CleanupCompletedTasks() error { return fmt.Errorf("failed to load tasks for cleanup: %w", err) } - // Filter completed and failed tasks, sort by completion time var completedTasks []*maintenance.MaintenanceTask for _, task := range allTasks { - if (task.Status == maintenance.TaskStatusCompleted || task.Status == maintenance.TaskStatusFailed) && task.CompletedAt != nil { + switch task.Status { + case maintenance.TaskStatusCompleted, maintenance.TaskStatusFailed, maintenance.TaskStatusCancelled: completedTasks = append(completedTasks, task) } } diff --git a/weed/admin/dash/maintenance_task_persistence_test.go b/weed/admin/dash/maintenance_task_persistence_test.go new file mode 100644 index 000000000..ea72dcfbe --- /dev/null +++ b/weed/admin/dash/maintenance_task_persistence_test.go @@ -0,0 +1,108 @@ +package dash + +import ( + "fmt" + "os" + "path/filepath" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/admin/maintenance" +) + +func countTaskStateFiles(t *testing.T, dir string) int { + t.Helper() + entries, err := os.ReadDir(filepath.Join(dir, TasksSubdir)) + if os.IsNotExist(err) { + return 0 + } + if err != nil { + t.Fatalf("read tasks dir: %v", err) + } + count := 0 + for _, entry := range entries { + if !entry.IsDir() && filepath.Ext(entry.Name()) == ".pb" { + count++ + } + } + return count +} + +// TestCancelledTaskFilesDoNotAccumulate reproduces issue #11595: every scan +// cycle cancels pending tasks of each detected type and re-detects them, so a +// cancelled task file per candidate volume accumulates on disk forever. +func TestCancelledTaskFilesDoNotAccumulate(t *testing.T) { + dir := t.TempDir() + cp := NewConfigPersistence(dir) + + queue := maintenance.NewMaintenanceQueue(nil) + queue.SetPersistence(cp) + + for cycle := 0; cycle < 3; cycle++ { + queue.AddTask(&maintenance.MaintenanceTask{ + ID: fmt.Sprintf("ec_vol_%d_cycle_%d", cycle, cycle), + Type: "erasure_coding", + VolumeID: uint32(cycle + 1), + Server: "server1", + }) + if cancelled := queue.CancelPendingTasksByType("erasure_coding"); cancelled != 1 { + t.Fatalf("cycle %d: cancelled %d tasks, want 1", cycle, cancelled) + } + } + + if n := countTaskStateFiles(t, dir); n != 0 { + t.Errorf("%d task files on disk after %d cancel cycles, want 0", n, 3) + } +} + +// TestManuallyCancelledTaskFileIsRemoved covers the CancelTask path used by +// the UI: a cancelled pending task must not leave its file behind, where a +// restart would resurrect it as pending. +func TestManuallyCancelledTaskFileIsRemoved(t *testing.T) { + dir := t.TempDir() + cp := NewConfigPersistence(dir) + + manager := maintenance.NewMaintenanceManager(nil, nil, cp) + queue := manager.GetQueue() + queue.SetPersistence(cp) + + queue.AddTask(&maintenance.MaintenanceTask{ + ID: "manual_1", + Type: "vacuum", + VolumeID: 7, + Server: "server1", + }) + if n := countTaskStateFiles(t, dir); n != 1 { + t.Fatalf("%d task files after AddTask, want 1", n) + } + + if err := manager.CancelTask("manual_1"); err != nil { + t.Fatalf("CancelTask: %v", err) + } + if n := countTaskStateFiles(t, dir); n != 0 { + t.Errorf("%d task files on disk after CancelTask, want 0", n) + } +} + +// TestCleanupCompletedTasksBoundsCancelledFiles checks retention covers +// cancelled files, e.g. ones written by older versions. +func TestCleanupCompletedTasksBoundsCancelledFiles(t *testing.T) { + dir := t.TempDir() + cp := NewConfigPersistence(dir) + + for i := 0; i < MaxCompletedTasks+5; i++ { + if err := cp.SaveTaskState(&maintenance.MaintenanceTask{ + ID: fmt.Sprintf("old_cancelled_%02d", i), + Type: "erasure_coding", + Status: maintenance.TaskStatusCancelled, + }); err != nil { + t.Fatalf("save task state: %v", err) + } + } + + if err := cp.CleanupCompletedTasks(); err != nil { + t.Fatalf("CleanupCompletedTasks: %v", err) + } + if n := countTaskStateFiles(t, dir); n > MaxCompletedTasks { + t.Errorf("%d task files after cleanup, want at most %d", n, MaxCompletedTasks) + } +} diff --git a/weed/admin/maintenance/maintenance_manager.go b/weed/admin/maintenance/maintenance_manager.go index 32618fb1b..4d2a518e1 100644 --- a/weed/admin/maintenance/maintenance_manager.go +++ b/weed/admin/maintenance/maintenance_manager.go @@ -452,6 +452,7 @@ func (mm *MaintenanceManager) performCleanup() { removedTasks := mm.queue.CleanupOldTasks(taskRetention) removedWorkers := mm.queue.RemoveStaleWorkers(workerTimeout) + mm.queue.cleanupCompletedTasks() // Clean up stale pending operations (operations running for more than 4 hours) staleOperationTimeout := 4 * time.Hour @@ -616,37 +617,48 @@ func (mm *MaintenanceManager) saveTaskConfigsFromPolicy(policy *worker_pb.Mainte // CancelTask cancels a pending task func (mm *MaintenanceManager) CancelTask(taskID string) error { mm.queue.mutex.Lock() - defer mm.queue.mutex.Unlock() task, exists := mm.queue.tasks[taskID] if !exists { + mm.queue.mutex.Unlock() return fmt.Errorf("task %s not found", taskID) } - - if task.Status == TaskStatusPending { - task.Status = TaskStatusCancelled - task.CompletedAt = &[]time.Time{time.Now()}[0] - - // Remove from pending tasks - for i, pendingTask := range mm.queue.pendingTasks { - if pendingTask.ID == taskID { - mm.queue.pendingTasks = append(mm.queue.pendingTasks[:i], mm.queue.pendingTasks[i+1:]...) - break - } - } - - // Notify ActiveTopology to release capacity - if mm.scanner != nil && mm.scanner.integration != nil { - if at := mm.scanner.integration.GetActiveTopology(); at != nil { - _ = at.CompleteTask(taskID) - } - } - - glog.V(2).Infof("Cancelled task %s", taskID) - return nil + if task.Status != TaskStatusPending { + status := task.Status + mm.queue.mutex.Unlock() + return fmt.Errorf("task %s cannot be cancelled (status: %s)", taskID, status) } - return fmt.Errorf("task %s cannot be cancelled (status: %s)", taskID, task.Status) + task.Status = TaskStatusCancelled + completedTime := time.Now() + task.CompletedAt = &completedTime + cancelledSnapshot := snapshotTask(task) + + // Remove from pending tasks + for i, pendingTask := range mm.queue.pendingTasks { + if pendingTask.ID == taskID { + mm.queue.pendingTasks = append(mm.queue.pendingTasks[:i], mm.queue.pendingTasks[i+1:]...) + break + } + } + + // Notify ActiveTopology to release capacity + if mm.scanner != nil && mm.scanner.integration != nil { + if at := mm.scanner.integration.GetActiveTopology(); at != nil { + _ = at.CompleteTask(taskID) + } + } + mm.queue.mutex.Unlock() + + if mm.queue.persistence != nil { + mm.queue.persistMu.Lock() + if mm.queue.deleteTaskStateLocked(taskID) != nil { + mm.queue.saveTaskStateLocked(cancelledSnapshot) + } + mm.queue.persistMu.Unlock() + } + glog.V(2).Infof("Cancelled task %s", taskID) + return nil } // RegisterWorker registers a new worker diff --git a/weed/admin/maintenance/maintenance_queue.go b/weed/admin/maintenance/maintenance_queue.go index b71063106..de4e2643b 100644 --- a/weed/admin/maintenance/maintenance_queue.go +++ b/weed/admin/maintenance/maintenance_queue.go @@ -97,21 +97,53 @@ func (mq *MaintenanceQueue) LoadTasksFromPersistence() error { return nil } -// saveTaskState saves a task to persistent storage +// isTerminalStatus reports whether the status is a terminal task state. +func isTerminalStatus(status MaintenanceTaskStatus) bool { + return status == TaskStatusCompleted || status == TaskStatusFailed || status == TaskStatusCancelled +} + +// saveTaskState saves a task to persistent storage. persistMu orders the +// status check and the write against the cancel paths' deletes, so a stale +// non-terminal file cannot resurrect the task on the next restart. func (mq *MaintenanceQueue) saveTaskState(task *MaintenanceTask) { - if mq.persistence != nil { - if err := mq.persistence.SaveTaskState(task); err != nil { - glog.Errorf("Failed to save task state for %s: %v", task.ID, err) - } + if mq.persistence == nil { + return + } + mq.persistMu.Lock() + defer mq.persistMu.Unlock() + mq.saveTaskStateLocked(task) +} + +// saveTaskStateLocked must be called with persistMu held. +func (mq *MaintenanceQueue) saveTaskStateLocked(task *MaintenanceTask) { + mq.mutex.RLock() + live, ok := mq.tasks[task.ID] + stale := !isTerminalStatus(task.Status) && (!ok || isTerminalStatus(live.Status)) + mq.mutex.RUnlock() + if stale { + return + } + if err := mq.persistence.SaveTaskState(task); err != nil { + glog.Errorf("Failed to save task state for %s: %v", task.ID, err) } } -func (mq *MaintenanceQueue) deleteTaskState(taskID string) { - if mq.persistence != nil { - if err := mq.persistence.DeleteTaskState(taskID); err != nil { - glog.V(2).Infof("Failed to delete task state for %s: %v", taskID, err) - } +func (mq *MaintenanceQueue) deleteTaskState(taskID string) error { + if mq.persistence == nil { + return nil } + mq.persistMu.Lock() + defer mq.persistMu.Unlock() + return mq.deleteTaskStateLocked(taskID) +} + +// deleteTaskStateLocked must be called with persistMu held. +func (mq *MaintenanceQueue) deleteTaskStateLocked(taskID string) error { + if err := mq.persistence.DeleteTaskState(taskID); err != nil { + glog.Warningf("Failed to delete task state for %s: %v", taskID, err) + return err + } + return nil } // cleanupCompletedTasks removes old completed tasks beyond the retention limit @@ -297,9 +329,18 @@ func (mq *MaintenanceQueue) CancelPendingTasksByType(taskType MaintenanceTaskTyp mq.pendingTasks = remaining mq.mutex.Unlock() - // Persist cancelled state outside the lock to avoid blocking - for _, snapshot := range cancelledSnapshots { - mq.saveTaskState(snapshot) + // Cancelled is terminal: drop the file like CompleteTask does instead of + // leaving one orphaned .pb per cancelled task per scan cycle. If removal + // fails, persist the cancelled state so a restart drops it instead of + // re-queueing it. + if mq.persistence != nil { + mq.persistMu.Lock() + for _, snapshot := range cancelledSnapshots { + if mq.deleteTaskStateLocked(snapshot.ID) != nil { + mq.saveTaskStateLocked(snapshot) + } + } + mq.persistMu.Unlock() } return cancelled } @@ -922,7 +963,7 @@ func generateTaskID() string { return fmt.Sprintf("%s-%04d", string(b), timestamp) } -// CleanupOldTasks removes old completed and failed tasks +// CleanupOldTasks removes old terminal tasks from memory func (mq *MaintenanceQueue) CleanupOldTasks(retention time.Duration) int { mq.mutex.Lock() defer mq.mutex.Unlock() @@ -931,7 +972,7 @@ func (mq *MaintenanceQueue) CleanupOldTasks(retention time.Duration) int { removed := 0 for id, task := range mq.tasks { - if (task.Status == TaskStatusCompleted || task.Status == TaskStatusFailed) && + if (task.Status == TaskStatusCompleted || task.Status == TaskStatusFailed || task.Status == TaskStatusCancelled) && task.CompletedAt != nil && task.CompletedAt.Before(cutoff) { delete(mq.tasks, id) diff --git a/weed/admin/maintenance/maintenance_types.go b/weed/admin/maintenance/maintenance_types.go index 1a20d764d..f7120a3c5 100644 --- a/weed/admin/maintenance/maintenance_types.go +++ b/weed/admin/maintenance/maintenance_types.go @@ -218,6 +218,10 @@ type MaintenanceQueue struct { policy *MaintenancePolicy integration *MaintenanceIntegration persistence TaskPersistence // Interface for task persistence + // persistMu serializes the check+write in saveTaskState against the + // deletes in the cancel paths, so a stale pending save cannot land on + // disk after the task's file was removed. + persistMu sync.Mutex } // MaintenanceScanner analyzes the cluster and generates maintenance tasks