diff --git a/weed/admin/dash/admin_server.go b/weed/admin/dash/admin_server.go index 37f89de83..505208923 100644 --- a/weed/admin/dash/admin_server.go +++ b/weed/admin/dash/admin_server.go @@ -1779,7 +1779,16 @@ func (s *AdminServer) ListPluginSchedulerStates() ([]adminplugin.SchedulerJobTyp // InitMaintenanceManager initializes the maintenance manager func (s *AdminServer) InitMaintenanceManager(config *maintenance.MaintenanceConfig) { - s.maintenanceManager = maintenance.NewMaintenanceManager(s, config) + // Hand the real config store to the manager so that, if it has to build the maintenance policy + // itself, it reads the persisted task configs instead of compiled-in defaults. Only pass it when + // a data directory is actually configured: an unconfigured store has nothing to read, and a typed + // nil pointer would satisfy the loaders' type assertion and then panic on use. + var configPersistence interface{} + if s.configPersistence != nil && s.configPersistence.IsConfigured() { + configPersistence = s.configPersistence + } + + s.maintenanceManager = maintenance.NewMaintenanceManager(s, config, configPersistence) // Set up task persistence if config persistence is available if s.configPersistence != nil { diff --git a/weed/admin/dash/config_persistence.go b/weed/admin/dash/config_persistence.go index fa3a0b43d..b36abc85b 100644 --- a/weed/admin/dash/config_persistence.go +++ b/weed/admin/dash/config_persistence.go @@ -156,7 +156,7 @@ func (cp *ConfigPersistence) LoadMaintenanceConfig() (*MaintenanceConfig, error) var config MaintenanceConfig if err := proto.Unmarshal(configData, &config); err == nil { // Always populate policy from separate task configuration files - config.Policy = buildPolicyFromTaskConfigs() + config.Policy = cp.buildPolicyFromTaskConfigs() return &config, nil } } @@ -687,8 +687,10 @@ func (cp *ConfigPersistence) GetConfigInfo() map[string]interface{} { return info } -// buildPolicyFromTaskConfigs loads task configurations from separate files and builds a MaintenancePolicy -func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { +// buildPolicyFromTaskConfigs loads task configurations from separate files and builds a MaintenancePolicy. +// cp is passed to each task loader so the persisted configs are actually honoured; the loaders fall +// back to compiled-in defaults for anything that has never been saved. +func (cp *ConfigPersistence) buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { policy := &worker_pb.MaintenancePolicy{ GlobalMaxConcurrent: 4, DefaultRepeatIntervalSeconds: 6 * 3600, // 6 hours in seconds @@ -697,7 +699,7 @@ func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { } // Load vacuum task configuration - if vacuumConfig := vacuum.LoadConfigFromPersistence(nil); vacuumConfig != nil { + if vacuumConfig := vacuum.LoadConfigFromPersistence(cp); vacuumConfig != nil { policy.TaskPolicies["vacuum"] = &worker_pb.TaskPolicy{ Enabled: vacuumConfig.Enabled, MaxConcurrent: int32(vacuumConfig.MaxConcurrent), @@ -713,7 +715,7 @@ func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { } // Load erasure coding task configuration - if ecConfig := erasure_coding.LoadConfigFromPersistence(nil); ecConfig != nil { + if ecConfig := erasure_coding.LoadConfigFromPersistence(cp); ecConfig != nil { policy.TaskPolicies["erasure_coding"] = &worker_pb.TaskPolicy{ Enabled: ecConfig.Enabled, MaxConcurrent: int32(ecConfig.MaxConcurrent), @@ -731,7 +733,7 @@ func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { } // Load balance task configuration - if balanceConfig := balance.LoadConfigFromPersistence(nil); balanceConfig != nil { + if balanceConfig := balance.LoadConfigFromPersistence(cp); balanceConfig != nil { policy.TaskPolicies["balance"] = &worker_pb.TaskPolicy{ Enabled: balanceConfig.Enabled, MaxConcurrent: int32(balanceConfig.MaxConcurrent), diff --git a/weed/admin/dash/maintenance_policy_persistence_test.go b/weed/admin/dash/maintenance_policy_persistence_test.go new file mode 100644 index 000000000..77e650fe5 --- /dev/null +++ b/weed/admin/dash/maintenance_policy_persistence_test.go @@ -0,0 +1,74 @@ +package dash + +import ( + "testing" + + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/balance" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" +) + +// TestLoadMaintenanceConfigHonoursPersistedTaskConfigs guards against the regression in +// https://github.com/seaweedfs/seaweedfs/issues/10874: buildPolicyFromTaskConfigs used to call +// LoadConfigFromPersistence(nil), which can never satisfy the loaders' type assertion, so every +// task silently fell back to its compiled-in defaults (Enabled: true) and a task disabled in the +// admin UI kept being scheduled. +func TestLoadMaintenanceConfigHonoursPersistedTaskConfigs(t *testing.T) { + dir := t.TempDir() + cp := NewConfigPersistence(dir) + + // A maintenance.pb must exist, otherwise LoadMaintenanceConfig returns early with defaults. + if err := cp.SaveMaintenanceConfig(DefaultMaintenanceConfig()); err != nil { + t.Fatalf("save maintenance config: %v", err) + } + + // Disable balance and vacuum the way the admin UI does, and change a value that is not a bool + // so a fallback to defaults cannot pass by coincidence. + disabledBalance := balance.NewDefaultConfig() + disabledBalance.Enabled = false + disabledBalance.MinServerCount = 7 + if err := cp.SaveBalanceTaskPolicy(disabledBalance.ToTaskPolicy()); err != nil { + t.Fatalf("save balance policy: %v", err) + } + + disabledVacuum := vacuum.NewDefaultConfig() + disabledVacuum.Enabled = false + if err := cp.SaveVacuumTaskPolicy(disabledVacuum.ToTaskPolicy()); err != nil { + t.Fatalf("save vacuum policy: %v", err) + } + + config, err := cp.LoadMaintenanceConfig() + if err != nil { + t.Fatalf("load maintenance config: %v", err) + } + if config.Policy == nil { + t.Fatal("policy is nil, want it populated from the persisted task configs") + } + + balancePolicy := config.Policy.TaskPolicies["balance"] + if balancePolicy == nil { + t.Fatal("no balance task policy in the built maintenance policy") + } + if balancePolicy.Enabled { + t.Error("balance enabled = true, want false from the persisted config") + } + if got := balancePolicy.GetBalanceConfig().GetMinServerCount(); got != 7 { + t.Errorf("balance min server count = %d, want persisted 7", got) + } + + vacuumPolicy := config.Policy.TaskPolicies["vacuum"] + if vacuumPolicy == nil { + t.Fatal("no vacuum task policy in the built maintenance policy") + } + if vacuumPolicy.Enabled { + t.Error("vacuum enabled = true, want false from the persisted config") + } + + // erasure_coding was never saved, so it keeps the loader's default of enabled. + ecPolicy := config.Policy.TaskPolicies["erasure_coding"] + if ecPolicy == nil { + t.Fatal("no erasure_coding task policy in the built maintenance policy") + } + if !ecPolicy.Enabled { + t.Error("erasure_coding enabled = false, want the default true for a config never saved") + } +} diff --git a/weed/admin/maintenance/maintenance_manager.go b/weed/admin/maintenance/maintenance_manager.go index 509cdc899..90e9a5d29 100644 --- a/weed/admin/maintenance/maintenance_manager.go +++ b/weed/admin/maintenance/maintenance_manager.go @@ -14,8 +14,13 @@ import ( "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" ) -// buildPolicyFromTaskConfigs loads task configurations from separate files and builds a MaintenancePolicy -func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { +// buildPolicyFromTaskConfigs loads task configurations from separate files and builds a MaintenancePolicy. +// +// configPersistence is duck-typed as interface{} because weed/admin/dash already imports this +// package, so importing *dash.ConfigPersistence back here would create an import cycle. It must be +// a value implementing the LoadXTaskPolicy() accessors the task loaders assert on; passing nil (or +// anything else) makes every task fall back to its compiled-in defaults. +func buildPolicyFromTaskConfigs(configPersistence interface{}) *worker_pb.MaintenancePolicy { policy := &worker_pb.MaintenancePolicy{ GlobalMaxConcurrent: 4, DefaultRepeatIntervalSeconds: 6 * 3600, // 6 hours in seconds @@ -24,7 +29,7 @@ func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { } // Load vacuum task configuration - if vacuumConfig := vacuum.LoadConfigFromPersistence(nil); vacuumConfig != nil { + if vacuumConfig := vacuum.LoadConfigFromPersistence(configPersistence); vacuumConfig != nil { policy.TaskPolicies["vacuum"] = &worker_pb.TaskPolicy{ Enabled: vacuumConfig.Enabled, MaxConcurrent: int32(vacuumConfig.MaxConcurrent), @@ -40,7 +45,7 @@ func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { } // Load erasure coding task configuration - if ecConfig := erasure_coding.LoadConfigFromPersistence(nil); ecConfig != nil { + if ecConfig := erasure_coding.LoadConfigFromPersistence(configPersistence); ecConfig != nil { policy.TaskPolicies["erasure_coding"] = &worker_pb.TaskPolicy{ Enabled: ecConfig.Enabled, MaxConcurrent: int32(ecConfig.MaxConcurrent), @@ -58,7 +63,7 @@ func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy { } // Load balance task configuration - if balanceConfig := balance.LoadConfigFromPersistence(nil); balanceConfig != nil { + if balanceConfig := balance.LoadConfigFromPersistence(configPersistence); balanceConfig != nil { policy.TaskPolicies["balance"] = &worker_pb.TaskPolicy{ Enabled: balanceConfig.Enabled, MaxConcurrent: int32(balanceConfig.MaxConcurrent), @@ -94,8 +99,12 @@ type MaintenanceManager struct { scanInProgress bool } -// NewMaintenanceManager creates a new maintenance manager -func NewMaintenanceManager(adminClient AdminClient, config *MaintenanceConfig) *MaintenanceManager { +// NewMaintenanceManager creates a new maintenance manager. +// +// configPersistence is the config store to read persisted task configs from when the policy has to +// be built here. See buildPolicyFromTaskConfigs for why it is duck-typed; pass nil when no config +// store is available. +func NewMaintenanceManager(adminClient AdminClient, config *MaintenanceConfig, configPersistence interface{}) *MaintenanceManager { if config == nil { config = DefaultMaintenanceConfig() } @@ -104,7 +113,7 @@ func NewMaintenanceManager(adminClient AdminClient, config *MaintenanceConfig) * policy := config.Policy if policy == nil { // Fallback: build policy from separate task configuration files if not already populated - policy = buildPolicyFromTaskConfigs() + policy = buildPolicyFromTaskConfigs(configPersistence) } queue := NewMaintenanceQueue(policy) diff --git a/weed/admin/maintenance/maintenance_manager_test.go b/weed/admin/maintenance/maintenance_manager_test.go index 243a88f5e..06b446baf 100644 --- a/weed/admin/maintenance/maintenance_manager_test.go +++ b/weed/admin/maintenance/maintenance_manager_test.go @@ -10,7 +10,7 @@ func TestMaintenanceManager_ErrorHandling(t *testing.T) { config := DefaultMaintenanceConfig() config.ScanIntervalSeconds = 1 // Short interval for testing (1 second) - manager := NewMaintenanceManager(nil, config) + manager := NewMaintenanceManager(nil, config, nil) // Test initial state if manager.errorCount != 0 { @@ -96,7 +96,7 @@ func TestIsConnectionError(t *testing.T) { func TestMaintenanceManager_GetErrorState(t *testing.T) { config := DefaultMaintenanceConfig() - manager := NewMaintenanceManager(nil, config) + manager := NewMaintenanceManager(nil, config, nil) // Test initial state errorCount, lastError, backoffDelay := manager.GetErrorState() @@ -118,7 +118,7 @@ func TestMaintenanceManager_GetErrorState(t *testing.T) { func TestMaintenanceManager_LogThrottling(t *testing.T) { config := DefaultMaintenanceConfig() - manager := NewMaintenanceManager(nil, config) + manager := NewMaintenanceManager(nil, config, nil) // This is a basic test to ensure the error handling doesn't panic // In practice, you'd want to capture log output to verify throttling diff --git a/weed/admin/maintenance/maintenance_policy_test.go b/weed/admin/maintenance/maintenance_policy_test.go new file mode 100644 index 000000000..a2ec6d4e0 --- /dev/null +++ b/weed/admin/maintenance/maintenance_policy_test.go @@ -0,0 +1,121 @@ +package maintenance + +import ( + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" +) + +// stubConfigPersistence implements the LoadXTaskPolicy accessors that the task config loaders +// type-assert on. The real implementation is *dash.ConfigPersistence, which this package cannot +// import: weed/admin/dash already imports weed/admin/maintenance, so the dependency only runs one +// way and the persistence argument has to stay duck-typed. +type stubConfigPersistence struct { + vacuum *worker_pb.TaskPolicy + ec *worker_pb.TaskPolicy + balance *worker_pb.TaskPolicy +} + +func (s *stubConfigPersistence) LoadVacuumTaskPolicy() (*worker_pb.TaskPolicy, error) { + return s.vacuum, nil +} + +func (s *stubConfigPersistence) LoadErasureCodingTaskPolicy() (*worker_pb.TaskPolicy, error) { + return s.ec, nil +} + +func (s *stubConfigPersistence) LoadBalanceTaskPolicy() (*worker_pb.TaskPolicy, error) { + return s.balance, nil +} + +func disabledStub() *stubConfigPersistence { + return &stubConfigPersistence{ + vacuum: &worker_pb.TaskPolicy{ + Enabled: false, + MaxConcurrent: 2, + RepeatIntervalSeconds: 2 * 3600, + TaskConfig: &worker_pb.TaskPolicy_VacuumConfig{ + VacuumConfig: &worker_pb.VacuumTaskConfig{GarbageThreshold: 0.3, MinVolumeAgeHours: 24}, + }, + }, + ec: &worker_pb.TaskPolicy{ + Enabled: false, + MaxConcurrent: 1, + RepeatIntervalSeconds: 3600, + TaskConfig: &worker_pb.TaskPolicy_ErasureCodingConfig{ + ErasureCodingConfig: &worker_pb.ErasureCodingTaskConfig{FullnessRatio: 0.95, QuietForSeconds: 3600, MinVolumeSizeMb: 30}, + }, + }, + balance: &worker_pb.TaskPolicy{ + Enabled: false, + MaxConcurrent: 1, + RepeatIntervalSeconds: 30 * 60, + TaskConfig: &worker_pb.TaskPolicy_BalanceConfig{ + BalanceConfig: &worker_pb.BalanceTaskConfig{ImbalanceThreshold: 0.2, MinServerCount: 7}, + }, + }, + } +} + +// TestBuildPolicyFromTaskConfigsUsesPersistence covers the bug reported in +// https://github.com/seaweedfs/seaweedfs/issues/10874: the persistence argument used to be a +// literal nil, which no type assertion can satisfy, so a task disabled on disk came back enabled. +func TestBuildPolicyFromTaskConfigsUsesPersistence(t *testing.T) { + policy := buildPolicyFromTaskConfigs(disabledStub()) + + for _, taskType := range []string{"vacuum", "erasure_coding", "balance"} { + taskPolicy := policy.TaskPolicies[taskType] + if taskPolicy == nil { + t.Fatalf("no %s task policy built", taskType) + } + if taskPolicy.Enabled { + t.Errorf("%s enabled = true, want false from the persisted config", taskType) + } + } + + if got := policy.TaskPolicies["balance"].GetBalanceConfig().GetMinServerCount(); got != 7 { + t.Errorf("balance min server count = %d, want persisted 7", got) + } +} + +// TestBuildPolicyFromTaskConfigsWithoutPersistence keeps the documented fallback: with no config +// store there is nothing to read, so the compiled-in defaults apply. +func TestBuildPolicyFromTaskConfigsWithoutPersistence(t *testing.T) { + policy := buildPolicyFromTaskConfigs(nil) + + for _, taskType := range []string{"vacuum", "erasure_coding", "balance"} { + taskPolicy := policy.TaskPolicies[taskType] + if taskPolicy == nil { + t.Fatalf("no %s task policy built", taskType) + } + if !taskPolicy.Enabled { + t.Errorf("%s enabled = false, want the compiled-in default true", taskType) + } + } +} + +// TestNewMaintenanceManagerUsesPersistenceForPolicyFallback covers the path the admin server takes +// when no maintenance.pb has been written yet: the config carries no policy, so the manager builds +// one itself and must read the persisted task configs to do it. +func TestNewMaintenanceManagerUsesPersistenceForPolicyFallback(t *testing.T) { + config := DefaultMaintenanceConfig() + if config.Policy != nil { + t.Fatal("default config unexpectedly carries a policy, test no longer covers the fallback") + } + + manager := NewMaintenanceManager(nil, config, disabledStub()) + + policy := manager.queue.policy + if policy == nil { + t.Fatal("queue policy is nil") + } + if IsTaskEnabled(policy, MaintenanceTaskType("balance")) { + t.Error("balance enabled = true, want false from the persisted config") + } + if IsTaskEnabled(policy, MaintenanceTaskType("vacuum")) { + t.Error("vacuum enabled = true, want false from the persisted config") + } + if manager.scanner.policy != policy { + t.Error("scanner and queue disagree on the policy") + } +} diff --git a/weed/admin/maintenance/maintenance_queue_test.go b/weed/admin/maintenance/maintenance_queue_test.go index 7a1328919..5da1de337 100644 --- a/weed/admin/maintenance/maintenance_queue_test.go +++ b/weed/admin/maintenance/maintenance_queue_test.go @@ -618,7 +618,7 @@ func TestMaintenanceQueue_StaleWorkerCapacityRelease(t *testing.T) { func TestMaintenanceManager_CancelTaskCapacityRelease(t *testing.T) { // Setup Manager config := DefaultMaintenanceConfig() - mm := NewMaintenanceManager(nil, config) + mm := NewMaintenanceManager(nil, config, nil) integration := mm.scanner.integration mq := mm.queue at := integration.GetActiveTopology()