fix(scheduler): manual migrate task also should fresh disk info

with: #1000744290

Signed-off-by: JasonHu520 <huzongchao@oppo.com>
This commit is contained in:
JasonHu520 2026-03-03 11:45:34 +08:00 committed by 梁曟風
parent 0c7f1fbb42
commit c73f38b2f4
7 changed files with 32 additions and 22 deletions

View File

@ -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
}

View File

@ -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

View File

@ -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

View File

@ -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
}

View File

@ -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

View File

@ -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)

View File

@ -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())