From c73f38b2f48903f33e6184730e6741a702837deb Mon Sep 17 00:00:00 2001 From: JasonHu520 Date: Tue, 3 Mar 2026 11:45:34 +0800 Subject: [PATCH] fix(scheduler): manual migrate task also should fresh disk info with: #1000744290 Signed-off-by: JasonHu520 --- blobstore/scheduler/balancer.go | 8 +++----- blobstore/scheduler/disk_droper.go | 3 +-- blobstore/scheduler/manual_migrater.go | 4 ++-- blobstore/scheduler/manual_migrater_test.go | 2 +- blobstore/scheduler/migrate.go | 13 ++++++------ blobstore/scheduler/migrate_test.go | 22 ++++++++++++++++----- blobstore/scheduler/startup.go | 2 +- 7 files changed, 32 insertions(+), 22 deletions(-) diff --git a/blobstore/scheduler/balancer.go b/blobstore/scheduler/balancer.go index 278d290fc..3cd8d538a 100644 --- a/blobstore/scheduler/balancer.go +++ b/blobstore/scheduler/balancer.go @@ -68,11 +68,9 @@ func NewBalanceMgr(clusterMgrCli client.ClusterMgrAPI, volumeUpdater client.Task clusterMgrCli: clusterMgrCli, cfg: conf, } - baseMgr := NewMigrateMgr(clusterMgrCli, volumeUpdater, taskSwitch, taskLogger, - &conf.MigrateConfig, proto.TaskTypeBalance) - baseMgr.diskGetter = clusterTopology - baseMgr.isBalanceAlloc = true - mgr.IMigrator = baseMgr + conf.MigrateConfig.IsBalanceAlloc = true + mgr.IMigrator = NewMigrateMgr(clusterMgrCli, volumeUpdater, taskSwitch, taskLogger, + &conf.MigrateConfig, proto.TaskTypeBalance, clusterTopology) return mgr } diff --git a/blobstore/scheduler/disk_droper.go b/blobstore/scheduler/disk_droper.go index 726df7bb9..265e90dd0 100644 --- a/blobstore/scheduler/disk_droper.go +++ b/blobstore/scheduler/disk_droper.go @@ -179,8 +179,7 @@ func NewDiskDropMgr(clusterMgrCli client.ClusterMgrAPI, volumeUpdater client.Tas return ErrHandleLockVolFail } baseMgr := NewMigrateMgr(clusterMgrCli, volumeUpdater, taskSwitch, taskLogger, - &conf.MigrateConfig, proto.TaskTypeDiskDrop) - baseMgr.diskGetter = topologyMgr + &conf.MigrateConfig, proto.TaskTypeDiskDrop, topologyMgr) mgr.IMigrator = baseMgr return mgr diff --git a/blobstore/scheduler/manual_migrater.go b/blobstore/scheduler/manual_migrater.go index b25353ed5..96e4195d8 100644 --- a/blobstore/scheduler/manual_migrater.go +++ b/blobstore/scheduler/manual_migrater.go @@ -35,13 +35,13 @@ type ManualMigrateMgr struct { // NewManualMigrateMgr returns manual migrate manager func NewManualMigrateMgr(clusterMgrCli client.ClusterMgrAPI, volumeUpdater client.TaskAPI, - taskLogger recordlog.Encoder, conf *MigrateConfig, + taskLogger recordlog.Encoder, conf *MigrateConfig, diskGetter DiskGetter, ) *ManualMigrateMgr { mgr := &ManualMigrateMgr{ clusterMgrCli: clusterMgrCli, } mgr.IMigrator = NewMigrateMgr(clusterMgrCli, volumeUpdater, taskswitch.NewEnabledTaskSwitch(), taskLogger, - conf, proto.TaskTypeManualMigrate) + conf, proto.TaskTypeManualMigrate, diskGetter) mgr.abnormalReporter = base.NewAbnormalReporter(conf.ClusterID, ShardRepair, base.ChunkMissMigrateAbnormal) conf.reportTaskCallback = mgr.reportMissChuckMigrated return mgr diff --git a/blobstore/scheduler/manual_migrater_test.go b/blobstore/scheduler/manual_migrater_test.go index 7aebd1be9..ea80edcd8 100644 --- a/blobstore/scheduler/manual_migrater_test.go +++ b/blobstore/scheduler/manual_migrater_test.go @@ -35,7 +35,7 @@ func newManualMigrater(t *testing.T) *ManualMigrateMgr { volumeUpdater := NewMockTaskAPI(ctr) taskLogger := mocks.NewMockRecordLogEncoder(ctr) migrater := NewMockMigrater(ctr) - mgr := NewManualMigrateMgr(clusterMgr, volumeUpdater, taskLogger, &MigrateConfig{ClusterID: 1}) + mgr := NewManualMigrateMgr(clusterMgr, volumeUpdater, taskLogger, &MigrateConfig{ClusterID: 1}, nil) mgr.IMigrator = migrater return mgr } diff --git a/blobstore/scheduler/migrate.go b/blobstore/scheduler/migrate.go index 7a8adb391..d98879335 100644 --- a/blobstore/scheduler/migrate.go +++ b/blobstore/scheduler/migrate.go @@ -336,6 +336,8 @@ type MigrateConfig struct { ClusterID proto.ClusterID `json:"-"` // fill in config.go base.TaskCommonConfig + IsBalanceAlloc bool `json:"-"` + lockFailHandleFunc lockFailFunc // clear junk tasks clearJunkTasksWhenLoadingFunc clearJunkTasksFunc @@ -399,8 +401,6 @@ type MigrateMgr struct { finishTaskCallback, loadTaskCallback taskLimitFunc // report task reportTaskCallback taskReportFunc - - isBalanceAlloc bool } // NewMigrateMgr returns migrate manager @@ -411,6 +411,7 @@ func NewMigrateMgr( taskLogger recordlog.Encoder, conf *MigrateConfig, taskType proto.TaskType, + diskGetter DiskGetter, ) *MigrateMgr { checkMigrateConf(conf) mgr := &MigrateMgr{ @@ -421,6 +422,7 @@ func NewMigrateMgr( clusterMgrCli: clusterMgrCli, volumeUpdater: volumeUpdater, + diskGetter: diskGetter, prepareQueue: base.NewTaskQueue(time.Duration(conf.PrepareQueueRetryDelayS) * time.Second), workQueue: base.NewWorkerTaskQueue(time.Duration(conf.CancelPunishDurationS) * time.Second), @@ -441,7 +443,6 @@ func NewMigrateMgr( if mgr.lockVolFailHandleFunc == nil { mgr.lockVolFailHandleFunc = mgr.handleLockVolFail } - mgr.isBalanceAlloc = false mgr.taskStatsMgr = base.NewTaskStatsMgrAndRun(conf.ClusterID, taskType, mgr) return mgr } @@ -615,7 +616,7 @@ func (mgr *MigrateMgr) prepareTask() (err error) { } // alloc volume unit - ret, err := base.AllocVunitSafe(ctx, mgr.clusterMgrCli, migTask.SourceVuid, migTask.Sources, nil, mgr.isBalanceAlloc) + ret, err := base.AllocVunitSafe(ctx, mgr.clusterMgrCli, migTask.SourceVuid, migTask.Sources, nil, mgr.cfg.IsBalanceAlloc) if err != nil { span.Errorf("alloc volume unit failed: err[%+v]", err) return @@ -829,7 +830,7 @@ func (mgr *MigrateMgr) handleUpdateVolMappingFail(ctx context.Context, task *pro if base.ShouldAllocAndRedo(code) { span.Infof("realloc vunit and redo: task_id[%s]", task.TaskID) - newVunit, err := base.AllocVunitSafe(ctx, mgr.clusterMgrCli, task.SourceVuid, task.Sources, nil, mgr.isBalanceAlloc) + newVunit, err := base.AllocVunitSafe(ctx, mgr.clusterMgrCli, task.SourceVuid, task.Sources, nil, mgr.cfg.IsBalanceAlloc) if err != nil { span.Errorf("realloc failed: vuid[%d], err[%+v]", task.SourceVuid, err) return err @@ -1007,7 +1008,7 @@ func (mgr *MigrateMgr) ReclaimTask(ctx context.Context, args *api.TaskArgs) (err } newDst, err := base.AllocVunitSafe(ctx, mgr.clusterMgrCli, arg.Src[arg.Dest.Vuid.Index()].Vuid, - arg.Src, []proto.DiskID{arg.Dest.DiskID}, mgr.isBalanceAlloc) + arg.Src, []proto.DiskID{arg.Dest.DiskID}, mgr.cfg.IsBalanceAlloc) if err != nil { span.Errorf("alloc volume unit from clustermgr failed, err: %s", err) return err diff --git a/blobstore/scheduler/migrate_test.go b/blobstore/scheduler/migrate_test.go index 0ea7aafa3..879fff830 100644 --- a/blobstore/scheduler/migrate_test.go +++ b/blobstore/scheduler/migrate_test.go @@ -72,7 +72,7 @@ func newMigrateMgr(t *testing.T) *MigrateMgr { }, } - mgr := NewMigrateMgr(clusterMgr, volumeUpdater, taskSwitch, taskLogger, conf, proto.TaskTypeBalance) + mgr := NewMigrateMgr(clusterMgr, volumeUpdater, taskSwitch, taskLogger, conf, proto.TaskTypeBalance, nil) return mgr } @@ -375,17 +375,29 @@ func TestAcquireMigrateTask(t *testing.T) { require.True(t, errors.Is(err, proto.ErrTaskEmpty)) } { - mgr := newMigrateMgr(t) ctr := gomock.NewController(t) + clusterMgr := NewMockClusterMgrAPI(ctr) + taskSwitch := mocks.NewMockSwitcher(ctr) + taskLogger := mocks.NewMockRecordLogEncoder(ctr) + volumeUpdater := NewMockTaskAPI(ctr) diskCache := NewMockDiskGetter(ctr) - mgr.diskGetter = diskCache + conf := &MigrateConfig{ + ClusterID: 0, + TaskCommonConfig: base.TaskCommonConfig{ + PrepareQueueRetryDelayS: 0, + FinishQueueRetryDelayS: 0, + CancelPunishDurationS: 0, + WorkQueueSize: 3, + }, + } + mgr := NewMigrateMgr(clusterMgr, volumeUpdater, taskSwitch, taskLogger, conf, proto.TaskTypeBalance, diskCache) t1 := mockGenMigrateTask(proto.TaskTypeDiskRepair, idc, 1, 1, proto.MigrateStatePrepared, newMockVolInfoMap()) mgr.workQueue.AddPreparedTask(idc, t1.TaskID, t1) - mgr.diskGetter.(*MockDiskGetter).EXPECT().GetDisk(any).AnyTimes().Return(&client.DiskInfoSimple{ + diskCache.EXPECT().GetDisk(any).AnyTimes().Return(&client.DiskInfoSimple{ DiskID: 12121, Host: "http://127.2.1.2:8889", }, true) - mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + clusterMgr.EXPECT().UpdateMigrateTask(any, any).Return(nil) task, err := mgr.AcquireTask(ctx, idc) require.NoError(t, err) diff --git a/blobstore/scheduler/startup.go b/blobstore/scheduler/startup.go index 6b21eeae3..5317d2b80 100644 --- a/blobstore/scheduler/startup.go +++ b/blobstore/scheduler/startup.go @@ -186,7 +186,7 @@ func NewService(conf *Config) (svr *Service, err error) { diskRepairMgr := NewDiskRepairMgr(clusterMgrCli, diskRepairTaskSwitch, taskLogger, &conf.DiskRepair, topologyMgr) - manualMigMgr := NewManualMigrateMgr(clusterMgrCli, taskCli, taskLogger, &conf.ManualMigrate) + manualMigMgr := NewManualMigrateMgr(clusterMgrCli, taskCli, taskLogger, &conf.ManualMigrate, topologyMgr) mqProxy := client.NewProxyClient(&conf.Proxy, cmapi.New(&conf.ClusterMgr), conf.ClusterID) inspectorTaskSwitch, err := switchMgr.AddSwitch(proto.TaskTypeVolumeInspect.String())