mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-20 13:30:46 +02:00
admin: honour persisted task configs when building the maintenance policy
buildPolicyFromTaskConfigs passed a literal nil to vacuum, erasure_coding
and balance LoadConfigFromPersistence. Those functions look for their
LoadXTaskPolicy() accessor via a type assertion, which a nil interface can
never satisfy, so every call fell through to NewDefaultConfig() and the
policy came back with the compiled-in defaults - Enabled: true among them.
A task disabled on disk was therefore still scheduled, and the only trace
was a glog.V(1) "Using default ... configuration" line.
Thread the real ConfigPersistence through instead. There are two copies of
this function: the one in weed/admin/dash builds config.Policy on the
normal admin startup path and can simply take cp as its receiver, and the
one in weed/admin/maintenance is the fallback used when the config carries
no policy yet, which now receives the store from NewMaintenanceManager.
weed/admin/dash already imports weed/admin/maintenance, so the maintenance
side has to keep the duck-typed interface{} parameter that the task
loaders already use rather than importing the concrete type back.
The store is only handed over when a data directory is configured: an
unconfigured one has nothing to read, and a typed nil pointer would pass
the loaders' type assertion and then panic on first use.
Fixes #10874
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user