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