diff --git a/blobstore/api/scheduler/task.go b/blobstore/api/scheduler/task.go index c1efb0f50..bef72d1d3 100644 --- a/blobstore/api/scheduler/task.go +++ b/blobstore/api/scheduler/task.go @@ -187,6 +187,10 @@ type MigrateTaskDetail struct { Stat proto.TaskStatistics `json:"stat"` } +type ShardTaskDetail struct { + Task proto.ShardMigrateTask `json:"task"` +} + type PerMinStats struct { FinishedCnt string `json:"finished_cnt"` ShardCnt string `json:"shard_cnt"` @@ -335,14 +339,14 @@ func hostWithScheme(host string) string { // ShardTaskArgs for shard node task action. type ShardTaskArgs struct { - IDC string `json:"idc"` - TaskID string `json:"task_id"` - TaskType proto.TaskType `json:"task_type"` - Source proto.SunitLocation `json:"source"` - Dest proto.SunitLocation `json:"dest"` - Leader proto.SunitLocation `json:"leader"` - Learner bool `json:"learner"` - Reason string `json:"reason"` + IDC string `json:"idc"` + TaskID string `json:"task_id"` + TaskType proto.TaskType `json:"task_type"` + Source proto.ShardUnitInfoSimple `json:"source"` + Dest proto.ShardUnitInfoSimple `json:"dest"` + Leader proto.ShardUnitInfoSimple `json:"leader"` + Learner bool `json:"learner"` + Reason string `json:"reason"` } func (t *ShardTaskArgs) Unmarshal(data []byte) error { diff --git a/blobstore/blobnode/task_shardnode_migrate.go b/blobstore/blobnode/task_shardnode_migrate.go index 959221445..61cd4dc22 100644 --- a/blobstore/blobnode/task_shardnode_migrate.go +++ b/blobstore/blobnode/task_shardnode_migrate.go @@ -42,7 +42,6 @@ func (s *ShardWorker) OperateArgs(reason string) *scheduler.TaskArgs { Source: s.t.Source, Dest: s.t.Destination, Leader: s.t.Leader, - Learner: s.t.Learner, Reason: reason, } data, _ := shardTaskArgs.Marshal() diff --git a/blobstore/common/proto/scheduler.go b/blobstore/common/proto/scheduler.go index 5c561d158..982ad32d1 100644 --- a/blobstore/common/proto/scheduler.go +++ b/blobstore/common/proto/scheduler.go @@ -91,17 +91,17 @@ func CheckVunitLocations(locations []VunitLocation) bool { return true } -// SunitLocation shard location -type SunitLocation struct { - // todo for proto suid - Suid Suid `json:"suid"` - Host string `json:"host"` - DiskID DiskID `json:"disk_id"` +// ShardUnitInfoSimple shard location +type ShardUnitInfoSimple struct { + Suid Suid `json:"suid"` + Host string `json:"host"` + DiskID DiskID `json:"disk_id"` + Learner bool `json:"learner"` } -// CheckSunitLocation CheckSunitLocations for shard task check -func CheckSunitLocation(location SunitLocation) bool { - if location.Suid == InvalidSuid || location.Host == "" || location.DiskID == InvalidDiskID { +// CheckShardUnitInfo for shard task check +func CheckShardUnitInfo(info ShardUnitInfoSimple) bool { + if info.Suid == InvalidSuid || info.Host == "" || info.DiskID == InvalidDiskID { return false } return true @@ -227,10 +227,9 @@ type ShardMigrateTask struct { Ctime string `json:"ctime"` // create time MTime string `json:"mtime"` // modify time - Source SunitLocation `json:"source"` // old shard location - Leader SunitLocation `json:"leader"` // shard leader location - Destination SunitLocation `json:"destination"` // new shard location - Learner bool `json:"learner"` + Source ShardUnitInfoSimple `json:"source"` // old shard location + Leader ShardUnitInfoSimple `json:"leader"` // shard leader location + Destination ShardUnitInfoSimple `json:"destination"` // new shard location } func (s *ShardMigrateTask) Unmarshal(data []byte) error { @@ -254,24 +253,27 @@ func (s *ShardMigrateTask) Task() (*Task, error) { return ret, err } -func (s *ShardMigrateTask) GetSource() SunitLocation { +func (s *ShardMigrateTask) GetSource() ShardUnitInfoSimple { return s.Source } -func (s *ShardMigrateTask) GetLeader() SunitLocation { +func (s *ShardMigrateTask) GetLeader() ShardUnitInfoSimple { return s.Leader } -func (s *ShardMigrateTask) GetDestination() SunitLocation { +func (s *ShardMigrateTask) GetDestination() ShardUnitInfoSimple { return s.Destination } -func (s *ShardMigrateTask) SetDestination(dest SunitLocation) { +func (s *ShardMigrateTask) SetDestination(dest ShardUnitInfoSimple) { s.Destination = dest } +func (s *ShardMigrateTask) Running() bool { + return s.State == ShardTaskStatePrepared || s.State == ShardTaskStateWorkCompleted +} func (s *ShardMigrateTask) IsValid() bool { - return CheckSunitLocation(s.Source) && CheckSunitLocation(s.Destination) + return CheckShardUnitInfo(s.Source) && CheckShardUnitInfo(s.Destination) } type VolumeInspectCheckPoint struct { diff --git a/blobstore/scheduler/balancer.go b/blobstore/scheduler/balancer.go index d5b5fd0db..b950ed9bc 100644 --- a/blobstore/scheduler/balancer.go +++ b/blobstore/scheduler/balancer.go @@ -177,7 +177,7 @@ func (mgr *BalanceMgr) genOneBalanceTask(ctx context.Context, diskInfo *client.D span.Debugf("select balance volume unit; vuid[%d], volume_id[%v]", vuid, vuid.Vid()) task := &proto.MigrateTask{ - TaskID: client.GenMigrateTaskID(proto.TaskTypeBalance, diskInfo.DiskID, vuid.Vid()), + TaskID: client.GenMigrateTaskID(proto.TaskTypeBalance, diskInfo.DiskID, uint32(vuid.Vid())), TaskType: proto.TaskTypeBalance, State: proto.MigrateStateInited, SourceIDC: diskInfo.Idc, diff --git a/blobstore/scheduler/base/queue.go b/blobstore/scheduler/base/queue.go index 4f45f46a8..573ab9fae 100644 --- a/blobstore/scheduler/base/queue.go +++ b/blobstore/scheduler/base/queue.go @@ -467,10 +467,10 @@ func vunitSliceEqual(a, b []proto.VunitLocation) bool { // ShardTask define shard task interface type ShardTask interface { - GetSource() proto.SunitLocation - GetLeader() proto.SunitLocation - GetDestination() proto.SunitLocation - SetDestination(dest proto.SunitLocation) + GetSource() proto.ShardUnitInfoSimple + GetLeader() proto.ShardUnitInfoSimple + GetDestination() proto.ShardUnitInfoSimple + SetDestination(dest proto.ShardUnitInfoSimple) Task() (*proto.Task, error) } @@ -483,6 +483,17 @@ type ShardTaskQueue struct { leaseExpiredS time.Duration } +func NewShardTaskQueue(cancelPunishDuration time.Duration) *ShardTaskQueue { + // extended lock duration of task leasing + leaseExpiredS := proto.TaskLeaseExpiredS * time.Second + + return &ShardTaskQueue{ + idcQueues: make(map[string]*Queue), + cancelPunishDuration: cancelPunishDuration, + leaseExpiredS: leaseExpiredS, + } +} + // AddPreparedTask add prepared task func (q *ShardTaskQueue) AddPreparedTask(idc, taskID string, wtask ShardTask) { q.mu.Lock() @@ -517,7 +528,7 @@ func (q *ShardTaskQueue) Acquire(idc string) (taskID string, wtask ShardTask, ex } // Cancel cancel task -func (q *ShardTaskQueue) Cancel(idc, taskID string, src, dst proto.SunitLocation) error { +func (q *ShardTaskQueue) Cancel(idc, taskID string, src, dst proto.ShardUnitInfoSimple) error { q.mu.Lock() defer q.mu.Unlock() @@ -539,7 +550,7 @@ func (q *ShardTaskQueue) Cancel(idc, taskID string, src, dst proto.SunitLocation } // Reclaim reclaim task -func (q *ShardTaskQueue) Reclaim(idc, taskID string, src, oldDest, newDest proto.SunitLocation, newDiskID proto.DiskID) error { +func (q *ShardTaskQueue) Reclaim(idc, taskID string, src, oldDest, newDest proto.ShardUnitInfoSimple, newDiskID proto.DiskID) error { q.mu.Lock() defer q.mu.Unlock() @@ -561,6 +572,23 @@ func (q *ShardTaskQueue) Reclaim(idc, taskID string, src, oldDest, newDest proto return idcQueue.Requeue(taskID, 0) } +// Query find task by idc and taskID +func (q *ShardTaskQueue) Query(idc, taskID string) (ShardTask, error) { + q.mu.Lock() + defer q.mu.Unlock() + + idcQueue, ok := q.idcQueues[idc] + if !ok { + return nil, errNoSuchIDCQueue + } + + wt, err := idcQueue.Get(taskID) + if err != nil { + return nil, err + } + return wt.(ShardTask), nil +} + // Renewal renewal task func (q *ShardTaskQueue) Renewal(idc, taskID string) error { q.mu.Lock() @@ -573,7 +601,7 @@ func (q *ShardTaskQueue) Renewal(idc, taskID string) error { } // Complete task -func (q *ShardTaskQueue) Complete(idc, taskID string, src, dst proto.SunitLocation) (ShardTask, error) { +func (q *ShardTaskQueue) Complete(idc, taskID string, src, dst proto.ShardUnitInfoSimple) (ShardTask, error) { q.mu.Lock() defer q.mu.Unlock() @@ -613,7 +641,7 @@ func (q *ShardTaskQueue) StatsTasks() (todo int, doing int) { return todo, doing } -func checkShardTaskValid(task ShardTask, src proto.SunitLocation, dst proto.SunitLocation) error { +func checkShardTaskValid(task ShardTask, src proto.ShardUnitInfoSimple, dst proto.ShardUnitInfoSimple) error { if task.GetSource() != src || task.GetDestination() != dst { return ErrUnmatchedVuids } diff --git a/blobstore/scheduler/base/task_types.go b/blobstore/scheduler/base/task_types.go index 197f3807b..07b3a4614 100644 --- a/blobstore/scheduler/base/task_types.go +++ b/blobstore/scheduler/base/task_types.go @@ -39,9 +39,10 @@ const ( // err use for task var ( - ErrNoTaskInQueue = errors.New("no task in queue") - ErrVolNotOnlyOneTask = errors.New("vol not only one task running") - ErrUpdateVolumeCache = errors.New("update volume cache failed") + ErrNoTaskInQueue = errors.New("no task in queue") + ErrVolNotOnlyOneTask = errors.New("vol not only one task running") + ErrUpdateVolumeCache = errors.New("update volume cache failed") + ErrShardNotOnlyOneTask = errors.New("shard not only one task running") ) // TaskCommonConfig task common config diff --git a/blobstore/scheduler/base/utils.go b/blobstore/scheduler/base/utils.go index 4ab51e3c3..b553fde59 100644 --- a/blobstore/scheduler/base/utils.go +++ b/blobstore/scheduler/base/utils.go @@ -63,6 +63,30 @@ func AllocVunitSafe( return allocVunit, nil } +type IAllocShardUnit interface { + AllocShardUnit(ctx context.Context, suid comproto.Suid) (ret *client.AllocShardUnitInfo, err error) +} + +// AllocShardUnitSafe alloc volume unit safe +func AllocShardUnitSafe( + ctx context.Context, + cli IAllocShardUnit, + src, dest comproto.ShardUnitInfoSimple, +) (ret *client.AllocShardUnitInfo, err error) { + span := trace.SpanFromContextSafe(ctx) + + allocShardUnit, err := cli.AllocShardUnit(ctx, src.Suid) + if err != nil { + return nil, err + } + + if allocShardUnit.DiskID == src.DiskID || allocShardUnit.DiskID == dest.DiskID { + span.Panic("alloc chunk and others chunks are on same disk") + } + + return allocShardUnit, nil +} + // Subtraction c = a - b func Subtraction(a, b []comproto.Vuid) (c []comproto.Vuid) { m := make(map[comproto.Vuid]struct{}) @@ -78,6 +102,21 @@ func Subtraction(a, b []comproto.Vuid) (c []comproto.Vuid) { return c } +// Sub c = a - b +func Sub(a, b []comproto.Suid) (c []comproto.Suid) { + m := make(map[comproto.Suid]struct{}) + for _, suid := range b { + m[suid] = struct{}{} + } + + for _, suid := range a { + if _, ok := m[suid]; !ok { + c = append(c, suid) + } + } + return c +} + // GenTaskID return task id func GenTaskID(prefix string, vid comproto.Vid) string { return fmt.Sprintf("%s-%d-%v", prefix, vid, xid.New().String()) @@ -120,6 +159,14 @@ func ShouldAllocAndRedo(errCode int) bool { return false } +// ShouldAllocShardUnitAndRedo return true if should alloc and redo task +func ShouldAllocShardUnitAndRedo(errCode int) bool { + if errCode == errors.CodeNewSuidNotMatch { + return true + } + return false +} + func InsistOn(ctx context.Context, errMsg string, on func() error) { span := trace.SpanFromContextSafe(ctx) attempt := 0 diff --git a/blobstore/scheduler/base/volume_task_locker.go b/blobstore/scheduler/base/volume_task_locker.go index 5937a015e..6eacb4305 100644 --- a/blobstore/scheduler/base/volume_task_locker.go +++ b/blobstore/scheduler/base/volume_task_locker.go @@ -19,29 +19,28 @@ import ( "errors" "sync" - "github.com/cubefs/cubefs/blobstore/common/proto" "github.com/cubefs/cubefs/blobstore/common/trace" ) // make sure only one task in same volume to run in cluster var ( // ErrVidTaskConflict vid task conflict - ErrVidTaskConflict = errors.New("vid task conflict") + ErrVidTaskConflict = errors.New("id task conflict") ) -// VolTaskLocker volume task locker -type VolTaskLocker struct { - taskMap map[proto.Vid]struct{} +// TaskLocker task locker +type TaskLocker struct { + taskMap map[uint32]struct{} mu sync.Mutex } -// TryLock try lock task volume and return error if there is task doing -func (m *VolTaskLocker) TryLock(ctx context.Context, vid proto.Vid) error { +// TryLock try lock task and return error if there is task doing +func (m *TaskLocker) TryLock(ctx context.Context, vid uint32) error { m.mu.Lock() defer m.mu.Unlock() span := trace.SpanFromContextSafe(ctx) - span.Infof("vid %d mutex try lock", vid) + span.Infof("id %d mutex try lock", vid) if _, ok := m.taskMap[vid]; ok { return ErrVidTaskConflict @@ -51,27 +50,41 @@ func (m *VolTaskLocker) TryLock(ctx context.Context, vid proto.Vid) error { } // Unlock unlock task volume -func (m *VolTaskLocker) Unlock(ctx context.Context, vid proto.Vid) { +func (m *TaskLocker) Unlock(ctx context.Context, vid uint32) { m.mu.Lock() defer m.mu.Unlock() span := trace.SpanFromContextSafe(ctx) - span.Infof("vid %d mutex unlock", vid) + span.Infof("id %d mutex unlock", vid) delete(m.taskMap, vid) } -var volTaskLocker *VolTaskLocker +var volTaskLocker *TaskLocker // NewVolTaskLockerOnce singleton mode:make sure only one instance in global var NewVolTaskLockerOnce sync.Once // VolTaskLockerInst ensure that only one background task is executing on the same volume -func VolTaskLockerInst() *VolTaskLocker { +func VolTaskLockerInst() *TaskLocker { NewVolTaskLockerOnce.Do(func() { - volTaskLocker = &VolTaskLocker{ - taskMap: make(map[proto.Vid]struct{}), + volTaskLocker = &TaskLocker{ + taskMap: make(map[uint32]struct{}), } }) return volTaskLocker } + +var ( + shardTaskLocker *TaskLocker + NewShardTaskLockerOnce sync.Once +) + +func ShardTaskLockerInst() *TaskLocker { + NewShardTaskLockerOnce.Do(func() { + shardTaskLocker = &TaskLocker{ + taskMap: make(map[uint32]struct{}), + } + }) + return shardTaskLocker +} diff --git a/blobstore/scheduler/base/volume_task_locker_test.go b/blobstore/scheduler/base/volume_task_locker_test.go index 7a9a59798..a0f667cda 100644 --- a/blobstore/scheduler/base/volume_task_locker_test.go +++ b/blobstore/scheduler/base/volume_task_locker_test.go @@ -19,14 +19,12 @@ import ( "testing" "github.com/stretchr/testify/require" - - "github.com/cubefs/cubefs/blobstore/common/proto" ) func MockEmptyVolTaskLocker() { VolTaskLockerInst().mu.Lock() defer VolTaskLockerInst().mu.Unlock() - VolTaskLockerInst().taskMap = make(map[proto.Vid]struct{}) + VolTaskLockerInst().taskMap = make(map[uint32]struct{}) } func TestVolTaskLocker(t *testing.T) { diff --git a/blobstore/scheduler/client/clustermgr.go b/blobstore/scheduler/client/clustermgr.go index bfc7e1087..f133d7652 100644 --- a/blobstore/scheduler/client/clustermgr.go +++ b/blobstore/scheduler/client/clustermgr.go @@ -50,12 +50,20 @@ type ClusterMgrVolumeAPI interface { } type ClusterMgrShardAPI interface { - GetShardInfo(ctx context.Context, Sid proto.Vid) (ret *ShardInfoSimple, err error) - UpdateShard(ctx context.Context, newVuid, oldVuid proto.Suid, newDiskID proto.DiskID) (err error) - AllocShardUnit(ctx context.Context, vuid proto.Vuid) (ret *AllocVunitInfo, err error) - ReleaseShardUnit(ctx context.Context, vuid proto.Vuid, diskID proto.DiskID) (err error) - ListDiskShardUnits(ctx context.Context, diskID proto.DiskID) (ret []*VunitInfoSimple, err error) - ListShard(ctx context.Context, marker proto.Vid, count int) (volInfo []*VolumeInfoSimple, retVid proto.Vid, err error) + GetShardInfo(ctx context.Context, ShardID proto.ShardID) (ret *ShardInfoSimple, err error) + UpdateShard(ctx context.Context, newSuid, oldSuid proto.Suid, newDiskID proto.DiskID) (err error) + AllocShardUnit(ctx context.Context, suid proto.Suid) (ret *AllocShardUnitInfo, err error) + ListDiskShardUnits(ctx context.Context, diskID proto.DiskID) (ret []*ShardUnitInfoSimple, err error) + ListShard(ctx context.Context, marker proto.ShardID, count int) (volInfo []*ShardInfoSimple, retShardID proto.ShardID, err error) +} + +type ClusterMgrShardDiskAPI interface { + ListShardDisk(ctx context.Context) (ret []*ShardNodeDiskInfo, err error) + ListBrokenShardDisk(ctx context.Context) (ret []*ShardNodeDiskInfo, err error) + ListRepairingShardDisk(ctx context.Context) (ret []*ShardNodeDiskInfo, err error) + SetShardDiskRepairing(ctx context.Context, id proto.DiskID) (err error) + SetShardDiskRepaired(ctx context.Context, id proto.DiskID) (err error) + GetShardDiskInfo(ctx context.Context, id proto.DiskID) (ret *ShardNodeDiskInfo, err error) } type ClusterMgrDiskAPI interface { @@ -69,17 +77,6 @@ type ClusterMgrDiskAPI interface { GetDiskInfo(ctx context.Context, diskID proto.DiskID) (ret *DiskInfoSimple, err error) } -type ClusterShardNodeAPI interface { - ListClusterDisks(ctx context.Context) (disks []*ShardNodeDiskInfo, err error) - ListBrokenDisks(ctx context.Context) (disks []*ShardNodeDiskInfo, err error) - ListRepairingDisks(ctx context.Context) (disks []*ShardNodeDiskInfo, err error) - ListDropDisks(ctx context.Context) (disks []*ShardNodeDiskInfo, err error) - SetDiskRepairing(ctx context.Context, diskID proto.DiskID) (err error) - SetDiskRepaired(ctx context.Context, diskID proto.DiskID) (err error) - SetDiskDropped(ctx context.Context, diskID proto.DiskID) (err error) - GetDiskInfo(ctx context.Context, diskID proto.DiskID) (ret *ShardNodeDiskInfo, err error) -} - type ClusterMgrServiceAPI interface { Register(ctx context.Context, info RegisterInfo) error GetService(ctx context.Context, name string, clusterID proto.ClusterID) (hosts []string, err error) @@ -110,6 +107,8 @@ type ClusterMgrAPI interface { ClusterMgrDiskAPI ClusterMgrServiceAPI ClusterMgrTaskAPI + ClusterMgrShardAPI + ClusterMgrShardDiskAPI } // migrate task key @@ -178,7 +177,7 @@ func genMigratingDiskPrefix(taskType proto.TaskType) string { } // GenMigrateTaskID return uniq task id -func GenMigrateTaskID(taskType proto.TaskType, diskID proto.DiskID, volumeID proto.Vid) string { +func GenMigrateTaskID(taskType proto.TaskType, diskID proto.DiskID, volumeID uint32) string { return fmt.Sprintf("%s%d%s%s", GenMigrateTaskPrefixByDiskID(taskType, diskID), volumeID, _delimiter, xid.New().String()) } @@ -279,6 +278,10 @@ type AllocVunitInfo struct { proto.VunitLocation } +type AllocShardUnitInfo struct { + proto.ShardUnitInfoSimple +} + // Location returns volume unit location func (vunit *AllocVunitInfo) Location() proto.VunitLocation { return vunit.VunitLocation @@ -319,19 +322,6 @@ type DiskInfoSimple struct { FreeChunkCnt int64 `json:"free_chunk_cnt"` } -// ShardNodeDiskInfo diskInfo for shard node -type ShardNodeDiskInfo struct { - ClusterID proto.ClusterID `json:"cluster_id"` - DiskID proto.DiskID `json:"disk_id"` - Idc string `json:"idc"` - Rack string `json:"rack"` - Host string `json:"host"` - Status proto.DiskStatus `json:"status"` - Readonly bool `json:"readonly"` - UsedShardCnt int64 `json:"used_shard_cnt"` - FreeShardCnt int64 `json:"free_shard_cnt"` -} - // IsHealth return true if disk is health func (disk *DiskInfoSimple) IsHealth() bool { return disk.Status == proto.DiskStatusNormal @@ -377,6 +367,63 @@ func (disk *DiskInfoSimple) set(info *clustermgr.BlobNodeDiskInfo) { disk.FreeChunkCnt = info.FreeChunkCnt } +// ShardNodeDiskInfo diskInfo for shard node +type ShardNodeDiskInfo struct { + ClusterID proto.ClusterID `json:"cluster_id"` + DiskID proto.DiskID `json:"disk_id"` + Idc string `json:"idc"` + Rack string `json:"rack"` + Host string `json:"host"` + Status proto.DiskStatus `json:"status"` + Readonly bool `json:"readonly"` + UsedShardCnt int32 `json:"used_shard_cnt"` + FreeShardCnt int32 `json:"free_shard_cnt"` +} + +// IsHealth return true if disk is health +func (disk *ShardNodeDiskInfo) IsHealth() bool { + return disk.Status == proto.DiskStatusNormal +} + +// IsBroken return true if disk is broken +func (disk *ShardNodeDiskInfo) IsBroken() bool { + return disk.Status == proto.DiskStatusBroken +} + +// IsDropped return true if disk is dropped +func (disk *ShardNodeDiskInfo) IsDropped() bool { + return disk.Status == proto.DiskStatusDropped +} + +// IsRepaired return true if disk is repaired +func (disk *ShardNodeDiskInfo) IsRepaired() bool { + return disk.Status == proto.DiskStatusRepaired +} + +// CanDropped disk can drop when disk is normal or has repaired or has dropped +// for simplicity we not allow to set disk status dropped +// when disk is repairing +func (disk *ShardNodeDiskInfo) CanDropped() bool { + if disk.Status == proto.DiskStatusNormal || + disk.Status == proto.DiskStatusRepaired || + disk.Status == proto.DiskStatusDropped { + return true + } + return false +} + +func (disk *ShardNodeDiskInfo) set(info *clustermgr.ShardNodeDiskInfo) { + disk.ClusterID = info.ClusterID + disk.Idc = info.Idc + disk.Rack = info.Rack + disk.Host = info.Host + disk.DiskID = info.DiskID + disk.Status = info.Status + disk.Readonly = info.Readonly + disk.FreeShardCnt = info.FreeShardCnt + disk.UsedShardCnt = info.UsedShardCnt +} + // RegisterInfo register info use for clustermgr type RegisterInfo struct { ClusterID uint64 `json:"cluster_id"` @@ -1024,11 +1071,11 @@ func (c *clustermgrClient) SetConsumeOffset(taskType proto.TaskType, topic strin // ShardInfoSimple shard info used by scheduler type ShardInfoSimple struct { - Sid proto.ShardID `json:"sid"` - ApplyIndex uint64 `json:"apply_index"` - Leader proto.NodeID `json:"leader"` - Status proto.ShardStatus `json:"status"` - SunitLocations []proto.SunitLocation `json:"sunit_locations"` + ShardID proto.ShardID `json:"sid"` + ApplyIndex uint64 `json:"apply_index"` + Leader uint8 `json:"leader"` + Status proto.ShardStatus `json:"status"` + ShardUnitInfoSimples []proto.ShardUnitInfoSimple `json:"sunit_locations"` } type UpdateShardArgs struct { @@ -1038,3 +1085,55 @@ type UpdateShardArgs struct { NewIsLearner bool `json:"new_is_learner"` NewDiskID bool `json:"new_disk_id"` } + +type ShardUnitInfoSimple struct { + Suid proto.Suid `json:"suid"` + DiskID proto.DiskID `json:"disk_id"` + Learner bool `json:"learner"` + Host string `json:"host"` + Status proto.ShardStatus `json:"status"` +} + +func (c *clustermgrClient) GetShardInfo(ctx context.Context, ShardID proto.ShardID) (ret *ShardInfoSimple, err error) { + return nil, nil +} + +func (c *clustermgrClient) UpdateShard(ctx context.Context, newSuid, oldSuid proto.Suid, newDiskID proto.DiskID) (err error) { + return nil +} + +func (c *clustermgrClient) AllocShardUnit(ctx context.Context, suid proto.Suid) (ret *AllocShardUnitInfo, err error) { + return nil, nil +} + +func (c *clustermgrClient) ListDiskShardUnits(ctx context.Context, diskID proto.DiskID) (ret []*ShardUnitInfoSimple, err error) { + return nil, nil +} + +func (c *clustermgrClient) ListShard(ctx context.Context, marker proto.ShardID, count int) (volInfo []*ShardInfoSimple, retVid proto.ShardID, err error) { + return nil, 0, err +} + +func (c *clustermgrClient) ListShardDisk(ctx context.Context) (ret []*ShardNodeDiskInfo, err error) { + return +} + +func (c *clustermgrClient) ListBrokenShardDisk(ctx context.Context) (ret []*ShardNodeDiskInfo, err error) { + return +} + +func (c *clustermgrClient) ListRepairingShardDisk(ctx context.Context) (ret []*ShardNodeDiskInfo, err error) { + return +} + +func (c *clustermgrClient) SetShardDiskRepairing(ctx context.Context, id proto.DiskID) (err error) { + return +} + +func (c *clustermgrClient) SetShardDiskRepaired(ctx context.Context, id proto.DiskID) (err error) { + return +} + +func (c *clustermgrClient) GetShardDiskInfo(ctx context.Context, id proto.DiskID) (ret *ShardNodeDiskInfo, err error) { + return +} diff --git a/blobstore/scheduler/client/clustermgr_test.go b/blobstore/scheduler/client/clustermgr_test.go index b61191f54..6e018daaa 100644 --- a/blobstore/scheduler/client/clustermgr_test.go +++ b/blobstore/scheduler/client/clustermgr_test.go @@ -83,16 +83,16 @@ func TestClustermgrClient(t *testing.T) { { // get volume info cli.client.(*MockClusterManager).EXPECT().GetVolumeInfo(any, any).Return(nil, errMock) - _, err := cli.GetVolumeInfo(ctx, proto.Vid(1)) + _, err := cli.GetVolumeInfo(ctx, 1) require.True(t, errors.Is(err, errMock)) volume := MockGenVolInfo(10, codemode.EC6P6, proto.VolumeStatusIdle) volume2 := MockGenVolInfo(10, codemode.EC6P6, proto.VolumeStatusActive) cli.client.(*MockClusterManager).EXPECT().GetVolumeInfo(any, any).Return(volume, nil) cli.client.(*MockClusterManager).EXPECT().GetVolumeInfo(any, any).Return(volume2, nil) - vol, err := cli.GetVolumeInfo(ctx, proto.Vid(1)) + vol, err := cli.GetVolumeInfo(ctx, 1) require.NoError(t, err) - vol2, err := cli.GetVolumeInfo(ctx, proto.Vid(2)) + vol2, err := cli.GetVolumeInfo(ctx, 1) require.NoError(t, err) require.Equal(t, vol.Vid, volume.Vid) require.Equal(t, vol.CodeMode, volume.CodeMode) @@ -106,21 +106,21 @@ func TestClustermgrClient(t *testing.T) { { // lock volume cli.client.(*MockClusterManager).EXPECT().LockVolume(any, any).Return(nil) - err := cli.LockVolume(ctx, proto.Vid(1)) + err := cli.LockVolume(ctx, 1) require.NoError(t, err) } { // unlock volume cli.client.(*MockClusterManager).EXPECT().UnlockVolume(any, any).Return(nil) - err := cli.UnlockVolume(ctx, proto.Vid(1)) + err := cli.UnlockVolume(ctx, 1) require.NoError(t, err) cli.client.(*MockClusterManager).EXPECT().UnlockVolume(any, any).Return(errcode.ErrUnlockNotAllow) - err = cli.UnlockVolume(ctx, proto.Vid(1)) + err = cli.UnlockVolume(ctx, 1) require.ErrorIs(t, errcode.ErrUnlockNotAllow, err) cli.client.(*MockClusterManager).EXPECT().UnlockVolume(any, any).Return(errMock) - err = cli.UnlockVolume(ctx, proto.Vid(1)) + err = cli.UnlockVolume(ctx, 1) require.True(t, errors.Is(err, errMock)) } { @@ -307,20 +307,20 @@ func TestClustermgrClient(t *testing.T) { { // add migrate task cli.client.(*MockClusterManager).EXPECT().SetKV(any, any, any).Return(nil) - task, _ := (&proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), proto.Vid(1))}).Task() + task, _ := (&proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), 1)}).Task() err := cli.AddMigrateTask(ctx, task) require.NoError(t, err) } { // update migrate task cli.client.(*MockClusterManager).EXPECT().SetKV(any, any, any).Return(nil) - task, _ := (&proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), proto.Vid(1))}).Task() + task, _ := (&proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), 1)}).Task() err := cli.UpdateMigrateTask(ctx, task) require.NoError(t, err) } { // get migrate task - task1, _ := (&proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), proto.Vid(1))}).Task() + task1, _ := (&proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), 1)}).Task() taskBytes, _, _ := task1.Marshal() cli.client.(*MockClusterManager).EXPECT().SetKV(any, any, any).Return(nil) cli.client.(*MockClusterManager).EXPECT().GetKV(any, any).Return(cmapi.GetKvRet{Value: taskBytes}, nil) @@ -337,7 +337,7 @@ func TestClustermgrClient(t *testing.T) { } { // delete migrate task - task1 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), proto.Vid(1))} + task1 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), 1)} cli.client.(*MockClusterManager).EXPECT().DeleteKV(any, any).Return(nil) err := cli.DeleteMigrateTask(ctx, task1.TaskID) require.NoError(t, err) @@ -357,9 +357,9 @@ func TestClustermgrClient(t *testing.T) { { // list all migrate tasks by disk_id diskID := proto.DiskID(100) - task1 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeBalance, diskID, proto.Vid(1)), TaskType: proto.TaskTypeBalance} + task1 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeBalance, diskID, 1), TaskType: proto.TaskTypeBalance} task1Bytes, _ := json.Marshal(task1) - task2 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeBalance, diskID, proto.Vid(2)), TaskType: proto.TaskTypeBalance} + task2 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeBalance, diskID, 1), TaskType: proto.TaskTypeBalance} task2Bytes, _ := json.Marshal(task2) cli.client.(*MockClusterManager).EXPECT().ListKV(any, any).Return(cmapi.ListKvRet{Kvs: []*cmapi.KeyValue{{Key: task1.TaskID, Value: task1Bytes}}, Marker: task1.TaskID}, nil) cli.client.(*MockClusterManager).EXPECT().ListKV(any, any).Return(cmapi.ListKvRet{Kvs: []*cmapi.KeyValue{{Key: task2.TaskID, Value: task2Bytes}}, Marker: task2.TaskID}, nil) @@ -381,7 +381,7 @@ func TestClustermgrClient(t *testing.T) { require.True(t, errors.Is(err, errMock)) // list all migrate task - task3 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeBalance, proto.DiskID(200), proto.Vid(2)), TaskType: proto.TaskTypeBalance} + task3 := &proto.MigrateTask{TaskID: GenMigrateTaskID(proto.TaskTypeBalance, proto.DiskID(200), 1), TaskType: proto.TaskTypeBalance} task3Bytes, _ := json.Marshal(task3) cli.client.(*MockClusterManager).EXPECT().ListKV(any, any).Return(cmapi.ListKvRet{Kvs: []*cmapi.KeyValue{ {Key: task1.TaskID, Value: task1Bytes}, diff --git a/blobstore/scheduler/client_mock_test.go b/blobstore/scheduler/client_mock_test.go index 50943d1e0..a7c88f380 100644 --- a/blobstore/scheduler/client_mock_test.go +++ b/blobstore/scheduler/client_mock_test.go @@ -65,6 +65,21 @@ func (mr *MockClusterMgrAPIMockRecorder) AddMigratingDisk(arg0, arg1 interface{} return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AddMigratingDisk", reflect.TypeOf((*MockClusterMgrAPI)(nil).AddMigratingDisk), arg0, arg1) } +// AllocShardUnit mocks base method. +func (m *MockClusterMgrAPI) AllocShardUnit(arg0 context.Context, arg1 proto.Suid) (*client.AllocShardUnitInfo, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "AllocShardUnit", arg0, arg1) + ret0, _ := ret[0].(*client.AllocShardUnitInfo) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// AllocShardUnit indicates an expected call of AllocShardUnit. +func (mr *MockClusterMgrAPIMockRecorder) AllocShardUnit(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AllocShardUnit", reflect.TypeOf((*MockClusterMgrAPI)(nil).AllocShardUnit), arg0, arg1) +} + // AllocVolumeUnit mocks base method. func (m *MockClusterMgrAPI) AllocVolumeUnit(arg0 context.Context, arg1 proto.Vuid) (*client.AllocVunitInfo, error) { m.ctrl.T.Helper() @@ -198,6 +213,36 @@ func (mr *MockClusterMgrAPIMockRecorder) GetService(arg0, arg1, arg2 interface{} return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetService", reflect.TypeOf((*MockClusterMgrAPI)(nil).GetService), arg0, arg1, arg2) } +// GetShardDiskInfo mocks base method. +func (m *MockClusterMgrAPI) GetShardDiskInfo(arg0 context.Context, arg1 proto.DiskID) (*client.ShardNodeDiskInfo, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "GetShardDiskInfo", arg0, arg1) + ret0, _ := ret[0].(*client.ShardNodeDiskInfo) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// GetShardDiskInfo indicates an expected call of GetShardDiskInfo. +func (mr *MockClusterMgrAPIMockRecorder) GetShardDiskInfo(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetShardDiskInfo", reflect.TypeOf((*MockClusterMgrAPI)(nil).GetShardDiskInfo), arg0, arg1) +} + +// GetShardInfo mocks base method. +func (m *MockClusterMgrAPI) GetShardInfo(arg0 context.Context, arg1 proto.ShardID) (*client.ShardInfoSimple, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "GetShardInfo", arg0, arg1) + ret0, _ := ret[0].(*client.ShardInfoSimple) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// GetShardInfo indicates an expected call of GetShardInfo. +func (mr *MockClusterMgrAPIMockRecorder) GetShardInfo(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetShardInfo", reflect.TypeOf((*MockClusterMgrAPI)(nil).GetShardInfo), arg0, arg1) +} + // GetVolumeInfo mocks base method. func (m *MockClusterMgrAPI) GetVolumeInfo(arg0 context.Context, arg1 proto.Vid) (*client.VolumeInfoSimple, error) { m.ctrl.T.Helper() @@ -273,6 +318,21 @@ func (mr *MockClusterMgrAPIMockRecorder) ListBrokenDisks(arg0 interface{}) *gomo return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListBrokenDisks", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListBrokenDisks), arg0) } +// ListBrokenShardDisk mocks base method. +func (m *MockClusterMgrAPI) ListBrokenShardDisk(arg0 context.Context) ([]*client.ShardNodeDiskInfo, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListBrokenShardDisk", arg0) + ret0, _ := ret[0].([]*client.ShardNodeDiskInfo) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListBrokenShardDisk indicates an expected call of ListBrokenShardDisk. +func (mr *MockClusterMgrAPIMockRecorder) ListBrokenShardDisk(arg0 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListBrokenShardDisk", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListBrokenShardDisk), arg0) +} + // ListClusterDisks mocks base method. func (m *MockClusterMgrAPI) ListClusterDisks(arg0 context.Context) ([]*client.DiskInfoSimple, error) { m.ctrl.T.Helper() @@ -288,6 +348,21 @@ func (mr *MockClusterMgrAPIMockRecorder) ListClusterDisks(arg0 interface{}) *gom return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListClusterDisks", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListClusterDisks), arg0) } +// ListDiskShardUnits mocks base method. +func (m *MockClusterMgrAPI) ListDiskShardUnits(arg0 context.Context, arg1 proto.DiskID) ([]*client.ShardUnitInfoSimple, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListDiskShardUnits", arg0, arg1) + ret0, _ := ret[0].([]*client.ShardUnitInfoSimple) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListDiskShardUnits indicates an expected call of ListDiskShardUnits. +func (mr *MockClusterMgrAPIMockRecorder) ListDiskShardUnits(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListDiskShardUnits", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListDiskShardUnits), arg0, arg1) +} + // ListDiskVolumeUnits mocks base method. func (m *MockClusterMgrAPI) ListDiskVolumeUnits(arg0 context.Context, arg1 proto.DiskID) ([]*client.VunitInfoSimple, error) { m.ctrl.T.Helper() @@ -364,6 +439,52 @@ func (mr *MockClusterMgrAPIMockRecorder) ListRepairingDisks(arg0 interface{}) *g return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListRepairingDisks", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListRepairingDisks), arg0) } +// ListRepairingShardDisk mocks base method. +func (m *MockClusterMgrAPI) ListRepairingShardDisk(arg0 context.Context) ([]*client.ShardNodeDiskInfo, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListRepairingShardDisk", arg0) + ret0, _ := ret[0].([]*client.ShardNodeDiskInfo) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListRepairingShardDisk indicates an expected call of ListRepairingShardDisk. +func (mr *MockClusterMgrAPIMockRecorder) ListRepairingShardDisk(arg0 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListRepairingShardDisk", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListRepairingShardDisk), arg0) +} + +// ListShard mocks base method. +func (m *MockClusterMgrAPI) ListShard(arg0 context.Context, arg1 proto.ShardID, arg2 int) ([]*client.ShardInfoSimple, proto.ShardID, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListShard", arg0, arg1, arg2) + ret0, _ := ret[0].([]*client.ShardInfoSimple) + ret1, _ := ret[1].(proto.ShardID) + ret2, _ := ret[2].(error) + return ret0, ret1, ret2 +} + +// ListShard indicates an expected call of ListShard. +func (mr *MockClusterMgrAPIMockRecorder) ListShard(arg0, arg1, arg2 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListShard", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListShard), arg0, arg1, arg2) +} + +// ListShardDisk mocks base method. +func (m *MockClusterMgrAPI) ListShardDisk(arg0 context.Context) ([]*client.ShardNodeDiskInfo, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListShardDisk", arg0) + ret0, _ := ret[0].([]*client.ShardNodeDiskInfo) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListShardDisk indicates an expected call of ListShardDisk. +func (mr *MockClusterMgrAPIMockRecorder) ListShardDisk(arg0 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListShardDisk", reflect.TypeOf((*MockClusterMgrAPI)(nil).ListShardDisk), arg0) +} + // ListVolume mocks base method. func (m *MockClusterMgrAPI) ListVolume(arg0 context.Context, arg1 proto.Vid, arg2 int) ([]*client.VolumeInfoSimple, proto.Vid, error) { m.ctrl.T.Helper() @@ -492,6 +613,34 @@ func (mr *MockClusterMgrAPIMockRecorder) SetDiskRepairing(arg0, arg1 interface{} return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SetDiskRepairing", reflect.TypeOf((*MockClusterMgrAPI)(nil).SetDiskRepairing), arg0, arg1) } +// SetShardDiskRepaired mocks base method. +func (m *MockClusterMgrAPI) SetShardDiskRepaired(arg0 context.Context, arg1 proto.DiskID) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "SetShardDiskRepaired", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// SetShardDiskRepaired indicates an expected call of SetShardDiskRepaired. +func (mr *MockClusterMgrAPIMockRecorder) SetShardDiskRepaired(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SetShardDiskRepaired", reflect.TypeOf((*MockClusterMgrAPI)(nil).SetShardDiskRepaired), arg0, arg1) +} + +// SetShardDiskRepairing mocks base method. +func (m *MockClusterMgrAPI) SetShardDiskRepairing(arg0 context.Context, arg1 proto.DiskID) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "SetShardDiskRepairing", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// SetShardDiskRepairing indicates an expected call of SetShardDiskRepairing. +func (mr *MockClusterMgrAPIMockRecorder) SetShardDiskRepairing(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SetShardDiskRepairing", reflect.TypeOf((*MockClusterMgrAPI)(nil).SetShardDiskRepairing), arg0, arg1) +} + // SetVolumeInspectCheckPoint mocks base method. func (m *MockClusterMgrAPI) SetVolumeInspectCheckPoint(arg0 context.Context, arg1 proto.Vid) error { m.ctrl.T.Helper() @@ -534,6 +683,20 @@ func (mr *MockClusterMgrAPIMockRecorder) UpdateMigrateTask(arg0, arg1 interface{ return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpdateMigrateTask", reflect.TypeOf((*MockClusterMgrAPI)(nil).UpdateMigrateTask), arg0, arg1) } +// UpdateShard mocks base method. +func (m *MockClusterMgrAPI) UpdateShard(arg0 context.Context, arg1, arg2 proto.Suid, arg3 proto.DiskID) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "UpdateShard", arg0, arg1, arg2, arg3) + ret0, _ := ret[0].(error) + return ret0 +} + +// UpdateShard indicates an expected call of UpdateShard. +func (mr *MockClusterMgrAPIMockRecorder) UpdateShard(arg0, arg1, arg2, arg3 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpdateShard", reflect.TypeOf((*MockClusterMgrAPI)(nil).UpdateShard), arg0, arg1, arg2, arg3) +} + // UpdateVolume mocks base method. func (m *MockClusterMgrAPI) UpdateVolume(arg0 context.Context, arg1, arg2 proto.Vuid, arg3 proto.DiskID) error { m.ctrl.T.Helper() diff --git a/blobstore/scheduler/config.go b/blobstore/scheduler/config.go index 76f16b55e..fd799ad6d 100644 --- a/blobstore/scheduler/config.go +++ b/blobstore/scheduler/config.go @@ -93,6 +93,8 @@ type Config struct { VolumeInspect VolumeInspectMgrCfg `json:"volume_inspect"` TaskLog recordlog.Config `json:"task_log"` + ShardDiskRepair ShardMigrateConfig `json:"shard_disk_repair"` + Kafka KafkaConfig `json:"kafka"` ShardRepair ShardRepairConfig `json:"shard_repair"` BlobDelete BlobDeleteConfig `json:"blob_delete"` diff --git a/blobstore/scheduler/disk_droper.go b/blobstore/scheduler/disk_droper.go index 19cee0278..20006f88c 100644 --- a/blobstore/scheduler/disk_droper.go +++ b/blobstore/scheduler/disk_droper.go @@ -423,7 +423,7 @@ func (mgr *DiskDropMgr) listUnMigratedVuid(ctx context.Context, diskID proto.Dis func (mgr *DiskDropMgr) initOneTask(ctx context.Context, src proto.Vuid, dropDiskID proto.DiskID, diskIDC string) { t := proto.MigrateTask{ - TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskDrop, dropDiskID, src.Vid()), + TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskDrop, dropDiskID, uint32(src.Vid())), TaskType: proto.TaskTypeDiskDrop, State: proto.MigrateStateInited, SourceDiskID: dropDiskID, diff --git a/blobstore/scheduler/disk_repairer.go b/blobstore/scheduler/disk_repairer.go index c52b39de4..4d095f306 100644 --- a/blobstore/scheduler/disk_repairer.go +++ b/blobstore/scheduler/disk_repairer.go @@ -113,7 +113,7 @@ func (mgr *DiskRepairMgr) Load() error { continue } if task.Running() { - err = base.VolTaskLockerInst().TryLock(ctx, task.Vid()) + err = base.VolTaskLockerInst().TryLock(ctx, uint32(task.Vid())) if err != nil { return fmt.Errorf("repair task conflict: task[%+v], err[%+v]", t, err.Error()) @@ -354,7 +354,7 @@ func (mgr *DiskRepairMgr) initOneTask(ctx context.Context, badVuid proto.Vuid, b span := trace.SpanFromContextSafe(ctx) t := proto.MigrateTask{ - TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, brokenDiskID, badVuid.Vid()), + TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, brokenDiskID, uint32(badVuid.Vid())), TaskType: proto.TaskTypeDiskRepair, State: proto.MigrateStateInited, SourceDiskID: brokenDiskID, @@ -414,7 +414,7 @@ func (mgr *DiskRepairMgr) popTaskAndPrepare() error { t := task.(*proto.MigrateTask).Copy() span.Infof("pop task: task_id[%s], task[%+v]", t.TaskID, t) // whether vid has another running task - err = base.VolTaskLockerInst().TryLock(ctx, t.Vid()) + err = base.VolTaskLockerInst().TryLock(ctx, uint32(t.Vid())) if err != nil { span.Warnf("tryLock failed: vid[%d]", t.Vid()) return base.ErrVolNotOnlyOneTask @@ -422,7 +422,7 @@ func (mgr *DiskRepairMgr) popTaskAndPrepare() error { defer func() { if err != nil { span.Errorf("prepare task failed: task_id[%s], err[%+v]", t.TaskID, err) - base.VolTaskLockerInst().Unlock(ctx, t.Vid()) + base.VolTaskLockerInst().Unlock(ctx, uint32(t.Vid())) } }() @@ -500,7 +500,7 @@ func (mgr *DiskRepairMgr) finishTaskInAdvance(ctx context.Context, task *proto.M mgr.finishTaskCounter.Add() mgr.prepareQueue.RemoveTask(task.TaskID) mgr.deletedTasks.add(task.SourceDiskID, task.TaskID) - base.VolTaskLockerInst().Unlock(ctx, task.Vid()) + base.VolTaskLockerInst().Unlock(ctx, uint32(task.Vid())) } func (mgr *DiskRepairMgr) finishTaskLoop() { @@ -582,7 +582,7 @@ func (mgr *DiskRepairMgr) finishTask(ctx context.Context, task *proto.MigrateTas // add delete task and check it again mgr.deletedTasks.add(task.SourceDiskID, task.TaskID) - base.VolTaskLockerInst().Unlock(ctx, task.Vid()) + base.VolTaskLockerInst().Unlock(ctx, uint32(task.Vid())) return nil } diff --git a/blobstore/scheduler/disk_repairer_test.go b/blobstore/scheduler/disk_repairer_test.go index ef39a9d4d..f2a0be141 100644 --- a/blobstore/scheduler/disk_repairer_test.go +++ b/blobstore/scheduler/disk_repairer_test.go @@ -77,10 +77,29 @@ func generateTaskArgs(task *proto.MigrateTask, reason string) *api.TaskArgs { return ret } +func genShardTaskArgs(task *proto.ShardMigrateTask, reason string) *api.TaskArgs { + ret := new(api.TaskArgs) + args := &api.ShardTaskArgs{ + IDC: task.SourceIDC, + TaskType: task.TaskType, + Source: task.Source, + TaskID: task.TaskID, + Reason: reason, + Dest: task.Destination, + Learner: task.Source.Learner, + Leader: task.Leader, + } + data, _ := args.Marshal() + ret.Data = data + ret.TaskType = task.TaskType + ret.ModuleType = proto.TypeBlobNode + return ret +} + func TestDiskRepairerLoad(t *testing.T) { task, err := (&proto.MigrateTask{ SourceDiskID: testDisk1.DiskID, - TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, testDisk1.DiskID, proto.Vid(1)), + TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, testDisk1.DiskID, 1), }).Task() require.NoError(t, err) { @@ -324,7 +343,7 @@ func TestDiskRepairerCollectTask(t *testing.T) { units = append(units, &ele) } t1, _ := (&proto.MigrateTask{ - TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), volume.Vid), + TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), uint32(volume.Vid)), TaskType: proto.TaskTypeDiskRepair, SourceVuid: units[0].Vuid, }).Task() @@ -582,7 +601,7 @@ func TestDiskRepairerCheckRepairedAndClear(t *testing.T) { units = append(units, &ele) } task := &proto.MigrateTask{ - TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), volume.Vid), + TaskID: client.GenMigrateTaskID(proto.TaskTypeDiskRepair, proto.DiskID(1), uint32(volume.Vid)), TaskType: proto.TaskTypeDiskRepair, SourceVuid: units[0].Vuid, } diff --git a/blobstore/scheduler/manual_migrater.go b/blobstore/scheduler/manual_migrater.go index c018db132..f2c435a72 100644 --- a/blobstore/scheduler/manual_migrater.go +++ b/blobstore/scheduler/manual_migrater.go @@ -60,7 +60,7 @@ func (mgr *ManualMigrateMgr) AddManualTask(ctx context.Context, vuid proto.Vuid, } task := &proto.MigrateTask{ - TaskID: client.GenMigrateTaskID(proto.TaskTypeManualMigrate, disk.DiskID, vuid.Vid()), + TaskID: client.GenMigrateTaskID(proto.TaskTypeManualMigrate, disk.DiskID, uint32(vuid.Vid())), TaskType: proto.TaskTypeManualMigrate, State: proto.MigrateStateInited, SourceIDC: disk.Idc, diff --git a/blobstore/scheduler/migrate.go b/blobstore/scheduler/migrate.go index bf0ffa962..b2d896524 100644 --- a/blobstore/scheduler/migrate.go +++ b/blobstore/scheduler/migrate.go @@ -53,6 +53,7 @@ type MMigrator interface { var ( _ BaseMigrator = (*DiskRepairMgr)(nil) _ BaseMigrator = (*MigrateMgr)(nil) + _ BaseMigrator = (*ShardMigrateMgr)(nil) ) // BaseMigrator base interface for shard and blobnode task. @@ -448,7 +449,7 @@ func (mgr *MigrateMgr) Load() (err error) { continue } if task.Running() { - err = base.VolTaskLockerInst().TryLock(ctx, task.SourceVuid.Vid()) + err = base.VolTaskLockerInst().TryLock(ctx, uint32(task.SourceVuid.Vid())) if err != nil { return fmt.Errorf("migrate task conflict: vid[%d], task[%+v], err[%+v]", task.SourceVuid.Vid(), tasks[i], err) @@ -544,14 +545,14 @@ func (mgr *MigrateMgr) prepareTask() (err error) { span.Infof("prepare task phase: task_id[%s], state[%+v]", migTask.TaskID, migTask.State) - err = base.VolTaskLockerInst().TryLock(ctx, migTask.SourceVuid.Vid()) + err = base.VolTaskLockerInst().TryLock(ctx, uint32(migTask.SourceVuid.Vid())) if err != nil { span.Warnf("lock volume failed: volume_id[%v], err[%+v]", migTask.SourceVuid.Vid(), err) return base.ErrVolNotOnlyOneTask } defer func() { if err != nil { - base.VolTaskLockerInst().Unlock(ctx, task.(*proto.MigrateTask).SourceVuid.Vid()) + base.VolTaskLockerInst().Unlock(ctx, uint32(task.(*proto.MigrateTask).SourceVuid.Vid())) } }() @@ -721,7 +722,7 @@ func (mgr *MigrateMgr) finishTask() (err error) { _ = mgr.finishQueue.RemoveTask(migrateTask.TaskID) _ = mgr.updateVolumeCache(ctx, migrateTask) - base.VolTaskLockerInst().Unlock(ctx, migrateTask.SourceVuid.Vid()) + base.VolTaskLockerInst().Unlock(ctx, uint32(migrateTask.SourceVuid.Vid())) mgr.deleteMigratingVuid(migrateTask.SourceDiskID, migrateTask.SourceVuid) mgr.finishTaskCounter.Add() @@ -782,7 +783,7 @@ func (mgr *MigrateMgr) finishTaskInAdvance(ctx context.Context, task *proto.Migr mgr.finishTaskCallback(task.SourceDiskID) - base.VolTaskLockerInst().Unlock(ctx, task.SourceVuid.Vid()) + base.VolTaskLockerInst().Unlock(ctx, uint32(task.SourceVuid.Vid())) } func (mgr *MigrateMgr) handleUpdateVolMappingFail(ctx context.Context, task *proto.MigrateTask, err error) error { diff --git a/blobstore/scheduler/scheduler_mock_test.go b/blobstore/scheduler/scheduler_mock_test.go index 1e54f6cdc..bc720e9c2 100644 --- a/blobstore/scheduler/scheduler_mock_test.go +++ b/blobstore/scheduler/scheduler_mock_test.go @@ -1,5 +1,5 @@ // Code generated by MockGen. DO NOT EDIT. -// Source: github.com/cubefs/cubefs/blobstore/scheduler (interfaces: ITaskRunner,IVolumeCache,MMigrator,IVolumeInspector,IClusterTopology) +// Source: github.com/cubefs/cubefs/blobstore/scheduler (interfaces: ITaskRunner,IVolumeCache,MMigrator,IVolumeInspector,IClusterTopology,ShardDiskMigrator) // Package scheduler is a generated GoMock package. package scheduler @@ -862,3 +862,306 @@ func (mr *MockClusterTopologyMockRecorder) UpdateVolume(arg0 interface{}) *gomoc mr.mock.ctrl.T.Helper() return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpdateVolume", reflect.TypeOf((*MockClusterTopology)(nil).UpdateVolume), arg0) } + +// MockShardDisMigrator is a mock of ShardDiskMigrator interface. +type MockShardDisMigrator struct { + ctrl *gomock.Controller + recorder *MockShardDisMigratorMockRecorder +} + +// MockShardDisMigratorMockRecorder is the mock recorder for MockShardDisMigrator. +type MockShardDisMigratorMockRecorder struct { + mock *MockShardDisMigrator +} + +// NewMockShardDisMigrator creates a new mock instance. +func NewMockShardDisMigrator(ctrl *gomock.Controller) *MockShardDisMigrator { + mock := &MockShardDisMigrator{ctrl: ctrl} + mock.recorder = &MockShardDisMigratorMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockShardDisMigrator) EXPECT() *MockShardDisMigratorMockRecorder { + return m.recorder +} + +// AcquireTask mocks base method. +func (m *MockShardDisMigrator) AcquireTask(arg0 context.Context, arg1 string) (*proto.Task, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "AcquireTask", arg0, arg1) + ret0, _ := ret[0].(*proto.Task) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// AcquireTask indicates an expected call of AcquireTask. +func (mr *MockShardDisMigratorMockRecorder) AcquireTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AcquireTask", reflect.TypeOf((*MockShardDisMigrator)(nil).AcquireTask), arg0, arg1) +} + +// AddTask mocks base method. +func (m *MockShardDisMigrator) AddTask(arg0 context.Context, arg1 *proto.ShardMigrateTask) { + m.ctrl.T.Helper() + m.ctrl.Call(m, "AddTask", arg0, arg1) +} + +// AddTask indicates an expected call of AddTask. +func (mr *MockShardDisMigratorMockRecorder) AddTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AddTask", reflect.TypeOf((*MockShardDisMigrator)(nil).AddTask), arg0, arg1) +} + +// CancelTask mocks base method. +func (m *MockShardDisMigrator) CancelTask(arg0 context.Context, arg1 *scheduler.TaskArgs) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "CancelTask", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// CancelTask indicates an expected call of CancelTask. +func (mr *MockShardDisMigratorMockRecorder) CancelTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "CancelTask", reflect.TypeOf((*MockShardDisMigrator)(nil).CancelTask), arg0, arg1) +} + +// Close mocks base method. +func (m *MockShardDisMigrator) Close() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "Close") +} + +// Close indicates an expected call of Close. +func (mr *MockShardDisMigratorMockRecorder) Close() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Close", reflect.TypeOf((*MockShardDisMigrator)(nil).Close)) +} + +// CompleteTask mocks base method. +func (m *MockShardDisMigrator) CompleteTask(arg0 context.Context, arg1 *scheduler.TaskArgs) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "CompleteTask", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// CompleteTask indicates an expected call of CompleteTask. +func (mr *MockShardDisMigratorMockRecorder) CompleteTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "CompleteTask", reflect.TypeOf((*MockShardDisMigrator)(nil).CompleteTask), arg0, arg1) +} + +// DiskProgress mocks base method. +func (m *MockShardDisMigrator) DiskProgress(arg0 context.Context, arg1 proto.DiskID) (*scheduler.DiskMigratingStats, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "DiskProgress", arg0, arg1) + ret0, _ := ret[0].(*scheduler.DiskMigratingStats) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// DiskProgress indicates an expected call of DiskProgress. +func (mr *MockShardDisMigratorMockRecorder) DiskProgress(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "DiskProgress", reflect.TypeOf((*MockShardDisMigrator)(nil).DiskProgress), arg0, arg1) +} + +// Done mocks base method. +func (m *MockShardDisMigrator) Done() <-chan struct{} { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Done") + ret0, _ := ret[0].(<-chan struct{}) + return ret0 +} + +// Done indicates an expected call of Done. +func (mr *MockShardDisMigratorMockRecorder) Done() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Done", reflect.TypeOf((*MockShardDisMigrator)(nil).Done)) +} + +// Enabled mocks base method. +func (m *MockShardDisMigrator) Enabled() bool { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Enabled") + ret0, _ := ret[0].(bool) + return ret0 +} + +// Enabled indicates an expected call of Enabled. +func (mr *MockShardDisMigratorMockRecorder) Enabled() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Enabled", reflect.TypeOf((*MockShardDisMigrator)(nil).Enabled)) +} + +// GetTask mocks base method. +func (m *MockShardDisMigrator) GetTask(arg0 context.Context, arg1 string) (*proto.ShardMigrateTask, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "GetTask", arg0, arg1) + ret0, _ := ret[0].(*proto.ShardMigrateTask) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// GetTask indicates an expected call of GetTask. +func (mr *MockShardDisMigratorMockRecorder) GetTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetTask", reflect.TypeOf((*MockShardDisMigrator)(nil).GetTask), arg0, arg1) +} + +// ListImmigratedSuid mocks base method. +func (m *MockShardDisMigrator) ListImmigratedSuid(arg0 context.Context, arg1 proto.DiskID) ([]proto.Suid, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListImmigratedSuid", arg0, arg1) + ret0, _ := ret[0].([]proto.Suid) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListImmigratedSuid indicates an expected call of ListImmigratedSuid. +func (mr *MockShardDisMigratorMockRecorder) ListImmigratedSuid(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListImmigratedSuid", reflect.TypeOf((*MockShardDisMigrator)(nil).ListImmigratedSuid), arg0, arg1) +} + +// ListMigratingSuid mocks base method. +func (m *MockShardDisMigrator) ListMigratingSuid(arg0 context.Context, arg1 proto.DiskID) ([]proto.Suid, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListMigratingSuid", arg0, arg1) + ret0, _ := ret[0].([]proto.Suid) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListMigratingSuid indicates an expected call of ListMigratingSuid. +func (mr *MockShardDisMigratorMockRecorder) ListMigratingSuid(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListMigratingSuid", reflect.TypeOf((*MockShardDisMigrator)(nil).ListMigratingSuid), arg0, arg1) +} + +// Load mocks base method. +func (m *MockShardDisMigrator) Load() error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Load") + ret0, _ := ret[0].(error) + return ret0 +} + +// Load indicates an expected call of Load. +func (mr *MockShardDisMigratorMockRecorder) Load() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Load", reflect.TypeOf((*MockShardDisMigrator)(nil).Load)) +} + +// Progress mocks base method. +func (m *MockShardDisMigrator) Progress(arg0 context.Context) ([]proto.DiskID, int, int) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Progress", arg0) + ret0, _ := ret[0].([]proto.DiskID) + ret1, _ := ret[1].(int) + ret2, _ := ret[2].(int) + return ret0, ret1, ret2 +} + +// Progress indicates an expected call of Progress. +func (mr *MockShardDisMigratorMockRecorder) Progress(arg0 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Progress", reflect.TypeOf((*MockShardDisMigrator)(nil).Progress), arg0) +} + +// QueryTask mocks base method. +func (m *MockShardDisMigrator) QueryTask(arg0 context.Context, arg1 string) (*scheduler.TaskRet, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "QueryTask", arg0, arg1) + ret0, _ := ret[0].(*scheduler.TaskRet) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// QueryTask indicates an expected call of QueryTask. +func (mr *MockShardDisMigratorMockRecorder) QueryTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "QueryTask", reflect.TypeOf((*MockShardDisMigrator)(nil).QueryTask), arg0, arg1) +} + +// ReclaimTask mocks base method. +func (m *MockShardDisMigrator) ReclaimTask(arg0 context.Context, arg1 *scheduler.TaskArgs) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ReclaimTask", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// ReclaimTask indicates an expected call of ReclaimTask. +func (mr *MockShardDisMigratorMockRecorder) ReclaimTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ReclaimTask", reflect.TypeOf((*MockShardDisMigrator)(nil).ReclaimTask), arg0, arg1) +} + +// RenewalTask mocks base method. +func (m *MockShardDisMigrator) RenewalTask(arg0 context.Context, arg1, arg2 string) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "RenewalTask", arg0, arg1, arg2) + ret0, _ := ret[0].(error) + return ret0 +} + +// RenewalTask indicates an expected call of RenewalTask. +func (mr *MockShardDisMigratorMockRecorder) RenewalTask(arg0, arg1, arg2 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RenewalTask", reflect.TypeOf((*MockShardDisMigrator)(nil).RenewalTask), arg0, arg1, arg2) +} + +// ReportTask mocks base method. +func (m *MockShardDisMigrator) ReportTask(arg0 context.Context, arg1 *scheduler.TaskArgs) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ReportTask", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// ReportTask indicates an expected call of ReportTask. +func (mr *MockShardDisMigratorMockRecorder) ReportTask(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ReportTask", reflect.TypeOf((*MockShardDisMigrator)(nil).ReportTask), arg0, arg1) +} + +// Run mocks base method. +func (m *MockShardDisMigrator) Run() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "Run") +} + +// Run indicates an expected call of Run. +func (mr *MockShardDisMigratorMockRecorder) Run() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Run", reflect.TypeOf((*MockShardDisMigrator)(nil).Run)) +} + +// Stats mocks base method. +func (m *MockShardDisMigrator) Stats() scheduler.ShardTaskStat { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Stats") + ret0, _ := ret[0].(scheduler.ShardTaskStat) + return ret0 +} + +// Stats indicates an expected call of Stats. +func (mr *MockShardDisMigratorMockRecorder) Stats() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Stats", reflect.TypeOf((*MockShardDisMigrator)(nil).Stats)) +} + +// WaitEnable mocks base method. +func (m *MockShardDisMigrator) WaitEnable() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "WaitEnable") +} + +// WaitEnable indicates an expected call of WaitEnable. +func (mr *MockShardDisMigratorMockRecorder) WaitEnable() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "WaitEnable", reflect.TypeOf((*MockShardDisMigrator)(nil).WaitEnable)) +} diff --git a/blobstore/scheduler/scheduler_test.go b/blobstore/scheduler/scheduler_test.go index 354240db5..2bd8d64b3 100644 --- a/blobstore/scheduler/scheduler_test.go +++ b/blobstore/scheduler/scheduler_test.go @@ -31,7 +31,7 @@ import ( // github.com/cubefs/cubefs/blobstore/scheduler/... module scheduler interfaces //go:generate mockgen -destination=./client_mock_test.go -package=scheduler -mock_names ClusterMgrAPI=MockClusterMgrAPI,BlobnodeAPI=MockBlobnodeAPI,IVolumeUpdater=MockVolumeUpdater,ProxyAPI=MockMqProxyAPI github.com/cubefs/cubefs/blobstore/scheduler/client ClusterMgrAPI,BlobnodeAPI,IVolumeUpdater,ProxyAPI //go:generate mockgen -destination=./base_mock_test.go -package=scheduler -mock_names KafkaConsumer=MockKafkaConsumer,GroupConsumer=MockGroupConsumer,IProducer=MockProducer github.com/cubefs/cubefs/blobstore/scheduler/base KafkaConsumer,GroupConsumer,IProducer -//go:generate mockgen -destination=./scheduler_mock_test.go -package=scheduler -mock_names ITaskRunner=MockTaskRunner,IVolumeCache=MockVolumeCache,MMigrator=MockMigrater,IVolumeInspector=MockVolumeInspector,IClusterTopology=MockClusterTopology github.com/cubefs/cubefs/blobstore/scheduler ITaskRunner,IVolumeCache,MMigrator,IVolumeInspector,IClusterTopology +//go:generate mockgen -destination=./scheduler_mock_test.go -package=scheduler -mock_names ITaskRunner=MockTaskRunner,IVolumeCache=MockVolumeCache,MMigrator=MockMigrater,IVolumeInspector=MockVolumeInspector,IClusterTopology=MockClusterTopology,ShardDiskMigrator=MockShardDisMigrator github.com/cubefs/cubefs/blobstore/scheduler ITaskRunner,IVolumeCache,MMigrator,IVolumeInspector,IClusterTopology,ShardDiskMigrator const ( testTopic = "test_topic" @@ -90,7 +90,7 @@ func mockGenMigrateTask(taskType proto.TaskType, idc string, diskID proto.DiskID codeMode := volInfoMap[vid].CodeMode vunitInfo := MockAlloc(volInfoMap[vid].VunitLocations[0].Vuid) task = &proto.MigrateTask{ - TaskID: client.GenMigrateTaskID(taskType, diskID, vid), + TaskID: client.GenMigrateTaskID(taskType, diskID, uint32(vid)), TaskType: taskType, State: state, SourceIDC: idc, @@ -106,6 +106,20 @@ func mockGenMigrateTask(taskType proto.TaskType, idc string, diskID proto.DiskID return task } +func mockGenShardMigrateTask(shardID proto.ShardID, taskType proto.TaskType, idc string, diskID proto.DiskID, + state proto.ShardTaskState, shardInfoMap map[proto.ShardID]*client.ShardInfoSimple) (task *proto.ShardMigrateTask) { + task = &proto.ShardMigrateTask{ + TaskID: client.GenMigrateTaskID(taskType, diskID, uint32(shardID)), + TaskType: taskType, + Ctime: time.Now().String(), + SourceIDC: idc, + State: state, + Source: shardInfoMap[shardID].ShardUnitInfoSimples[2], + Leader: shardInfoMap[shardID].ShardUnitInfoSimples[shardInfoMap[shardID].Leader], + } + return task +} + func MockGenVolInfo(vid proto.Vid, cm codemode.CodeMode, status proto.VolumeStatus) *client.VolumeInfoSimple { vol := client.VolumeInfoSimple{} cmInfo := cm.Tactic() @@ -126,6 +140,27 @@ func MockGenVolInfo(vid proto.Vid, cm codemode.CodeMode, status proto.VolumeStat return &vol } +func MockGenShardInfo(shardID proto.ShardID, leader uint8, status proto.ShardStatus) *client.ShardInfoSimple { + shard := new(client.ShardInfoSimple) + shard.Leader = leader + shard.ShardID = shardID + shard.Status = status + shard.ApplyIndex = 0 + host := "127.0.0.0:xxx" + sunits := make([]proto.ShardUnitInfoSimple, 0, 3) + for i := 0; i < 3; i++ { + sunits = append(sunits, proto.ShardUnitInfoSimple{ + DiskID: proto.DiskID(i + 1), + Suid: proto.EncodeSuid(shardID, uint8(i), 0), + Host: host, + Learner: false, + }) + } + shard.ShardUnitInfoSimples = sunits + + return shard +} + func MockAlloc(vuid proto.Vuid) *client.AllocVunitInfo { vid := vuid.Vid() idx := vuid.Index() diff --git a/blobstore/scheduler/service.go b/blobstore/scheduler/service.go index d1c244966..6871bc850 100644 --- a/blobstore/scheduler/service.go +++ b/blobstore/scheduler/service.go @@ -44,6 +44,8 @@ type Service struct { manualMigMgr IManualMigrator inspectMgr IVolumeInspector + shardDiskRepairMgr ShardDiskMigrator + shardRepairMgr ITaskRunner blobDeleteMgr ITaskRunner clusterTopology IClusterTopology @@ -64,8 +66,7 @@ func (svr *Service) mgrByType(typ proto.TaskType) (BaseMigrator, error) { case proto.TaskTypeManualMigrate: return svr.manualMigMgr, nil case proto.TaskTypeShardDiskRepair: - // todo - return nil, nil + return svr.shardDiskRepairMgr, nil case proto.TaskTypeShardInspect: return nil, errIllegalTaskType case proto.TaskTypeShardMigrate: diff --git a/blobstore/scheduler/service_test.go b/blobstore/scheduler/service_test.go index a0805978e..f0586075f 100644 --- a/blobstore/scheduler/service_test.go +++ b/blobstore/scheduler/service_test.go @@ -196,13 +196,13 @@ func TestServiceAPI(t *testing.T) { for _, taskType := range taskTypes { taskArgs, err := (&api.OperateTaskArgs{ IDC: idc, TaskType: taskType, - TaskID: client.GenMigrateTaskID(taskType, diskID, volumeID), + TaskID: client.GenMigrateTaskID(taskType, diskID, uint32(volumeID)), }).TaskArgs() require.NoError(t, err) require.NoError(t, cli.ReclaimTask(ctx, taskArgs)) require.NoError(t, cli.CancelTask(ctx, taskArgs)) require.NoError(t, cli.CompleteTask(ctx, taskArgs)) - args, err := (&api.TaskReportArgs{TaskType: taskType, TaskID: client.GenMigrateTaskID(taskType, diskID, volumeID)}).TaskArgs() + args, err := (&api.TaskReportArgs{TaskType: taskType, TaskID: client.GenMigrateTaskID(taskType, diskID, uint32(volumeID))}).TaskArgs() require.NoError(t, err) require.NoError(t, cli.ReportTask(ctx, args)) } @@ -274,9 +274,9 @@ func TestServiceAPI(t *testing.T) { require.Error(t, err) } for _, taskType := range taskTypes { - _, err = cli.DetailMigrateTask(ctx, &api.MigrateTaskDetailArgs{Type: taskType, ID: client.GenMigrateTaskID(taskType, diskID, volumeID)}) + _, err = cli.DetailMigrateTask(ctx, &api.MigrateTaskDetailArgs{Type: taskType, ID: client.GenMigrateTaskID(taskType, diskID, uint32(volumeID))}) require.NoError(t, err) - _, err = cli.DetailMigrateTask(ctx, &api.MigrateTaskDetailArgs{Type: taskType, ID: client.GenMigrateTaskID(taskType, diskID, volumeID)}) + _, err = cli.DetailMigrateTask(ctx, &api.MigrateTaskDetailArgs{Type: taskType, ID: client.GenMigrateTaskID(taskType, diskID, uint32(volumeID))}) require.Error(t, err) } // disk migrating stats diff --git a/blobstore/scheduler/shard_disk_repairer.go b/blobstore/scheduler/shard_disk_repairer.go new file mode 100644 index 000000000..5a18680d8 --- /dev/null +++ b/blobstore/scheduler/shard_disk_repairer.go @@ -0,0 +1,396 @@ +// Copyright 2024 The CubeFS Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or +// implied. See the License for the specific language governing +// permissions and limitations under the License. + +package scheduler + +import ( + "context" + "errors" + "fmt" + "time" + + api "github.com/cubefs/cubefs/blobstore/api/scheduler" + "github.com/cubefs/cubefs/blobstore/common/proto" + "github.com/cubefs/cubefs/blobstore/common/taskswitch" + "github.com/cubefs/cubefs/blobstore/common/trace" + "github.com/cubefs/cubefs/blobstore/scheduler/base" + "github.com/cubefs/cubefs/blobstore/scheduler/client" +) + +type ShardDiskMigrator interface { + ShardMigrator + + Progress(ctx context.Context) (migratingDisks []proto.DiskID, total, migrated int) + DiskProgress(ctx context.Context, diskID proto.DiskID) (stats *api.DiskMigratingStats, err error) +} + +type ShardDiskRepairMgr struct { + ShardMigrator + + repairedDisks *migratedDisks + repairingDisks *migratingShardDisks + + hasRevised bool + + clusterMgrCli client.ClusterMgrAPI + + cfg *ShardMigrateConfig +} + +func NewShardDiskRepairMgr( + cfg *ShardMigrateConfig, + clusterMgrCli client.ClusterMgrAPI, + taskSwitch taskswitch.ISwitcher) *ShardDiskRepairMgr { + mgr := &ShardDiskRepairMgr{ + clusterMgrCli: clusterMgrCli, + repairingDisks: newMigratingShardDisks(), + repairedDisks: newMigratedDisks(), + } + cfg.loadTaskCallback = mgr.loadTaskCallback + mgr.ShardMigrator = NewShardMigrateMgr(clusterMgrCli, taskSwitch, cfg, proto.TaskTypeShardDiskRepair) + return mgr +} + +// Load load repair task from database +func (mgr *ShardDiskRepairMgr) Load() error { + span, ctx := trace.StartSpanFromContext(context.Background(), "shard_disk_repair.Load") + diskInfos, err := mgr.clusterMgrCli.ListRepairingShardDisk(ctx) + if err != nil { + span.Errorf("list repairing shard disk failed, err: %s", err) + return err + } + for _, disk := range diskInfos { + mgr.repairingDisks.add(disk.DiskID, disk) + } + err = mgr.ShardMigrator.Load() + if err != nil { + return err + } + + return nil +} + +// Run shard disk repair task manager +func (mgr *ShardDiskRepairMgr) Run() { + go mgr.collectTaskLoop() + mgr.ShardMigrator.Run() + go mgr.checkRepairedAndClearLoop() + go mgr.checkAndClearJunkTasksLoop() +} + +// Close shard disk repair task manager +func (mgr *ShardDiskRepairMgr) Close() { + mgr.ShardMigrator.Close() +} + +func (mgr *ShardDiskRepairMgr) Progress(ctx context.Context) (repairingDisk []proto.DiskID, total, migrated int) { + span := trace.SpanFromContextSafe(ctx) + repairingDisk = make([]proto.DiskID, 0) + + for _, disk := range mgr.repairingDisks.list() { + total += int(disk.UsedShardCnt) + remainTasks, err := mgr.clusterMgrCli.ListAllMigrateTasksByDiskID(ctx, proto.TaskTypeShardDiskRepair, disk.DiskID) + if err != nil { + span.Errorf("find all task failed: err[%+v]", err) + return repairingDisk, 0, 0 + } + migrated += int(disk.UsedShardCnt) - len(remainTasks) + repairingDisk = append(repairingDisk, disk.DiskID) + } + return +} + +func (mgr *ShardDiskRepairMgr) DiskProgress(ctx context.Context, diskID proto.DiskID) (stats *api.DiskMigratingStats, err error) { + span := trace.SpanFromContextSafe(ctx) + + migratingDisk, ok := mgr.repairingDisks.get(diskID) + if !ok { + err = errors.New("not repairing disk") + return + } + remainTasks, err := mgr.clusterMgrCli.ListAllMigrateTasksByDiskID(ctx, proto.TaskTypeDiskRepair, diskID) + if err != nil { + span.Errorf("find all task failed: err[%+v]", err) + return + } + stats = &api.DiskMigratingStats{} + stats.TotalTasksCnt = int(migratingDisk.UsedShardCnt) + stats.MigratedTasksCnt = stats.TotalTasksCnt - len(remainTasks) + return +} + +func (mgr *ShardDiskRepairMgr) collectTaskLoop() { + t := time.NewTicker(time.Duration(mgr.cfg.CollectTaskIntervalS) * time.Second) + defer t.Stop() + + for { + select { + case <-t.C: + mgr.ShardMigrator.WaitEnable() + mgr.collectionTask() + case <-mgr.ShardMigrator.Done(): + return + } + } +} + +func (mgr *ShardDiskRepairMgr) collectionTask() { + span, ctx := trace.StartSpanFromContext(context.Background(), "disk_repair.collectTask") + defer span.Finish() + + // revise repair tasks to make sure data consistency when services start + if !mgr.hasRevised { + if err := mgr.reviseRepairDisks(ctx); err != nil { + return + } + mgr.hasRevised = true + } + + if mgr.repairingDisks.size() >= mgr.cfg.DiskConcurrency { + return + } + disk, err := mgr.acquireBrokenDisk(ctx) + if err != nil { + span.Warnf("acquire broken shard disk from clustermgr failed, err[%s]", err) + return + } + if disk == nil { + return + } + if err = mgr.generateTask(ctx, disk); err != nil { + span.Errorf("generate shard disk repair tasks failed: err[%+v]", err) + return + } + + execMsg := fmt.Sprintf("set shard disk diskId %d repairing", disk.DiskID) + base.InsistOn(ctx, execMsg, func() error { + return mgr.clusterMgrCli.SetShardDiskRepairing(ctx, disk.DiskID) + }) + + mgr.repairingDisks.add(disk.DiskID, disk) + span.Infof("init repair task for shard disk[%d] success", disk.DiskID) +} + +func (mgr *ShardDiskRepairMgr) acquireBrokenDisk(ctx context.Context) (*client.ShardNodeDiskInfo, error) { + diskInfos, err := mgr.clusterMgrCli.ListBrokenShardDisk(ctx) + if err != nil { + return nil, err + } + for _, disk := range diskInfos { + if _, exist := mgr.repairingDisks.get(disk.DiskID); !exist { + return disk, nil + } + } + return nil, nil +} + +// checkAndClearJunkTasksLoop due to network timeout, it may still have some junk migrate tasks in clustermgr, +// and we need to clear those tasks later +func (mgr *ShardDiskRepairMgr) checkAndClearJunkTasksLoop() { + t := time.NewTicker(clearJunkMigrationTaskInterval) + defer t.Stop() + + for { + select { + case <-t.C: + mgr.checkAndClearJunkTasks() + case <-mgr.ShardMigrator.Done(): + return + } + } +} + +func (mgr *ShardDiskRepairMgr) checkAndClearJunkTasks() { + span, ctx := trace.StartSpanFromContext(context.Background(), "shard_disk_repair.clearJunkTasks") + + for _, disk := range mgr.repairedDisks.list() { + if time.Since(disk.finishedTime) < junkMigrationTaskProtectionWindow { + continue + } + span.Debugf("check repaired shard disk: disk_id[%d], repaired time[%v]", disk.diskID, disk.finishedTime) + diskInfo, err := mgr.clusterMgrCli.GetShardDiskInfo(ctx, disk.diskID) + if err != nil { + span.Errorf("get disk info failed: disk_id[%d], err[%+v]", disk.diskID, err) + continue + } + if !diskInfo.IsRepaired() { + continue + } + tasks, err := mgr.clusterMgrCli.ListAllMigrateTasksByDiskID(ctx, proto.TaskTypeShardDiskRepair, disk.diskID) + if err != nil { + continue + } + if len(tasks) != 0 { + span.Warnf("clear junk tasks of repaired disk: disk_id[%d], tasks size[%d]", disk.diskID, len(tasks)) + for _, task := range tasks { + span.Warnf("check and delete junk task: task_id[%s]", task.TaskID) + base.InsistOn(ctx, "chek and delete junk task", func() error { + return mgr.clusterMgrCli.DeleteMigrateTask(ctx, task.TaskID) + }) + } + } + mgr.repairedDisks.delete(disk.diskID) + } +} + +func (mgr *ShardDiskRepairMgr) checkRepairedAndClearLoop() { + t := time.NewTicker(time.Duration(mgr.cfg.CheckTaskIntervalS) * time.Second) + defer t.Stop() + + for { + select { + case <-t.C: + mgr.WaitEnable() + mgr.checkRepairedAndClear() + case <-mgr.Done(): + return + } + } +} + +func (mgr *ShardDiskRepairMgr) checkRepairedAndClear() { + span, ctx := trace.StartSpanFromContext(context.Background(), "disk_repair.checkRepairedAndClear") + + for _, disk := range mgr.repairingDisks.list() { + if !mgr.checkDiskRepaired(ctx, disk.DiskID) { + continue + } + err := mgr.clusterMgrCli.SetShardDiskRepaired(ctx, disk.DiskID) + if err != nil { + return + } + span.Infof("disk repaired will start clear: disk_id[%d]", disk.DiskID) + // mgr.clearTasksByDiskID(ctx, disk.DiskID) + } +} + +func (mgr *ShardDiskRepairMgr) checkDiskRepaired(ctx context.Context, diskID proto.DiskID) bool { + span := trace.SpanFromContextSafe(ctx) + span.Infof("check repaired: disk_id[%d]", diskID) + + tasks, err := mgr.clusterMgrCli.ListAllMigrateTasksByDiskID(ctx, proto.TaskTypeShardDiskRepair, diskID) + if err != nil { + span.Errorf("check repaired and find task failed: disk_iD[%d], err[%+v]", diskID, err) + return false + } + sunitInfos, err := mgr.clusterMgrCli.ListDiskShardUnits(ctx, diskID) + if err != nil { + span.Errorf("check repaired list disk shard units failed: disk_id[%s], err[%+v]", diskID, err) + return false + } + if len(sunitInfos) == 0 && len(tasks) != 0 { + // due to network timeout, it may lead to repeated insertion of deleted tasks, and need to delete it again + // mgr.clearJunkTasks(ctx, diskID, tasks) + return false + } + if len(sunitInfos) != 0 && len(tasks) == 0 { + // it may be occur when migration done and repair tasks generate concurrent, list volume units may not return the migrate unit + span.Warnf("clustermgr has some shard unit not repair and revise again: disk_id[%d], shard units len[%d]", diskID, len(sunitInfos)) + if err = mgr.repairDisk(ctx, diskID); err != nil { + span.Errorf("revise repair task failed: err[%+v]", err) + } + return false + } + return len(tasks) == 0 && len(sunitInfos) == 0 +} + +func (mgr *ShardDiskRepairMgr) loadTaskCallback(ctx context.Context, diskId proto.DiskID) { + if _, exist := mgr.repairingDisks.get(diskId); exist { + return + } + info, err := mgr.clusterMgrCli.GetShardDiskInfo(ctx, diskId) + if err != nil { + return + } + mgr.repairingDisks.add(diskId, &client.ShardNodeDiskInfo{ + ClusterID: info.ClusterID, + DiskID: info.DiskID, + Idc: info.Idc, + Rack: info.Rack, + Host: info.Host, + Status: info.Status, + Readonly: info.Readonly, + UsedShardCnt: info.UsedShardCnt, + FreeShardCnt: info.FreeShardCnt, + }) +} + +func (mgr *ShardDiskRepairMgr) reviseRepairDisks(ctx context.Context) error { + span := trace.SpanFromContextSafe(ctx) + for _, disk := range mgr.repairingDisks.list() { + if err := mgr.repairDisk(ctx, disk.DiskID); err != nil { + span.Errorf("revise repair shard tasks failed: disk_id[%d]", disk.DiskID) + return err + } + } + return nil +} + +func (mgr *ShardDiskRepairMgr) repairDisk(ctx context.Context, diskID proto.DiskID) error { + span := trace.SpanFromContextSafe(ctx) + + diskInfo, err := mgr.clusterMgrCli.GetShardDiskInfo(ctx, diskID) + if err != nil { + span.Errorf("get shard disk info failed: err[%+v]", err) + return err + } + + if err = mgr.generateTask(ctx, diskInfo); err != nil { + span.Errorf("generate shard disk repair tasks failed: err[%+v]", err) + return err + } + + if diskInfo.IsBroken() { + execMsg := fmt.Sprintf("set shard disk diskId %d repairing", diskID) + base.InsistOn(ctx, execMsg, func() error { + return mgr.clusterMgrCli.SetShardDiskRepairing(ctx, diskID) + }) + } + return nil +} + +func (mgr *ShardDiskRepairMgr) generateTask(ctx context.Context, disk *client.ShardNodeDiskInfo) error { + span := trace.SpanFromContextSafe(ctx) + span.Infof("start generate shard disk repair tasks: disk_id[%d], disk_idc[%s]", disk.DiskID, disk.Idc) + + migratingSuids, err := mgr.ListMigratingSuid(ctx, disk.DiskID) + if err != nil { + span.Errorf("list repairing suids failed: err[%s]", err) + return err + } + + immigratedSuids, err := mgr.ListImmigratedSuid(ctx, disk.DiskID) + if err != nil { + span.Errorf("list un repaired suids failed: err[%s]", err) + return err + } + + remain := base.Sub(immigratedSuids, migratingSuids) + span.Infof("should gen shard tasks remain: len[%d]", len(remain)) + for _, suid := range remain { + mgr.initOneTask(ctx, suid, disk.DiskID, disk.Idc) + } + return nil +} + +func (mgr *ShardDiskRepairMgr) initOneTask(ctx context.Context, suid proto.Suid, diskID proto.DiskID, idc string) { + task := &proto.ShardMigrateTask{ + TaskID: client.GenMigrateTaskID(proto.TaskTypeShardDiskRepair, diskID, uint32(suid.ShardID())), + TaskType: proto.TaskTypeShardDiskRepair, + State: proto.ShardTaskStateInited, + SourceIDC: idc, + Source: proto.ShardUnitInfoSimple{Suid: suid, DiskID: diskID}, + } + mgr.ShardMigrator.AddTask(ctx, task) +} diff --git a/blobstore/scheduler/shard_migrate.go b/blobstore/scheduler/shard_migrate.go index ab9d46859..98066e744 100644 --- a/blobstore/scheduler/shard_migrate.go +++ b/blobstore/scheduler/shard_migrate.go @@ -16,13 +16,17 @@ package scheduler import ( "context" + "encoding/json" "errors" + "fmt" + "sync" "time" api "github.com/cubefs/cubefs/blobstore/api/scheduler" "github.com/cubefs/cubefs/blobstore/common/counter" errcode "github.com/cubefs/cubefs/blobstore/common/errors" "github.com/cubefs/cubefs/blobstore/common/proto" + "github.com/cubefs/cubefs/blobstore/common/rpc" "github.com/cubefs/cubefs/blobstore/common/taskswitch" "github.com/cubefs/cubefs/blobstore/common/trace" "github.com/cubefs/cubefs/blobstore/scheduler/base" @@ -34,12 +38,18 @@ import ( type ShardMigrateConfig struct { ClusterID proto.ClusterID `json:"-"` // fill in config.go base.TaskCommonConfig + + loadTaskCallback loadTaskCallback } type ShardMigrator interface { BaseMigrator Stats() api.ShardTaskStat + AddTask(ctx context.Context, task *proto.ShardMigrateTask) + GetTask(ctx context.Context, taskID string) (*proto.ShardMigrateTask, error) + ListMigratingSuid(ctx context.Context, diskID proto.DiskID) (suids []proto.Suid, err error) + ListImmigratedSuid(ctx context.Context, diskID proto.DiskID) (suids []proto.Suid, err error) taskswitch.ISwitcher closer.Closer @@ -65,6 +75,31 @@ type ShardMigrateMgr struct { cfg *ShardMigrateConfig } +type loadTaskCallback func(ctx context.Context, diskId proto.DiskID) + +// NewShardMigrateMgr returns migrate manager +func NewShardMigrateMgr( + clusterMgrCli client.ClusterMgrAPI, + taskSwitch taskswitch.ISwitcher, + conf *ShardMigrateConfig, + taskType proto.TaskType, +) ShardMigrator { + mgr := &ShardMigrateMgr{ + taskType: taskType, + taskSwitch: taskSwitch, + clusterMgrCli: clusterMgrCli, + + prepareQueue: base.NewTaskQueue(time.Duration(conf.PrepareQueueRetryDelayS) * time.Second), + workQueue: base.NewShardTaskQueue(time.Duration(conf.CancelPunishDurationS) * time.Second), + finishQueue: base.NewTaskQueue(time.Duration(conf.FinishQueueRetryDelayS) * time.Second), + + cfg: conf, + Closer: closer.New(), + } + mgr.taskStatsMgr = base.NewTaskStatsMgrAndRun(conf.ClusterID, taskType, mgr) + return mgr +} + func (mgr *ShardMigrateMgr) AcquireTask(ctx context.Context, idc string) (task *proto.Task, err error) { span := trace.SpanFromContextSafe(ctx) task = new(proto.Task) @@ -80,6 +115,7 @@ func (mgr *ShardMigrateMgr) AcquireTask(ctx context.Context, idc string) (task * return task, err } task.Data = data + task.TaskID = t.TaskID span.Infof("acquire %s taskId: %s", mgr.taskType, t.TaskID) return task, nil } @@ -92,7 +128,7 @@ func (mgr *ShardMigrateMgr) CancelTask(ctx context.Context, args *api.TaskArgs) arg := &api.ShardTaskArgs{} err := arg.Unmarshal(args.Data) if err != nil { - return err + return errcode.ErrIllegalArguments } if !client.ValidMigrateTask(args.TaskType, arg.TaskID) { return errcode.ErrIllegalArguments @@ -112,7 +148,7 @@ func (mgr *ShardMigrateMgr) CompleteTask(ctx context.Context, args *api.TaskArgs arg := &api.ShardTaskArgs{} err := arg.Unmarshal(args.Data) if err != nil { - return err + return errcode.ErrIllegalArguments } if !client.ValidMigrateTask(args.TaskType, arg.TaskID) { return errcode.ErrIllegalArguments @@ -141,16 +177,76 @@ func (mgr *ShardMigrateMgr) CompleteTask(ctx context.Context, args *api.TaskArgs return nil } -func (mgr *ShardMigrateMgr) ReclaimTask(ctx context.Context, args *api.TaskArgs) error { - return nil +func (mgr *ShardMigrateMgr) ReclaimTask(ctx context.Context, args *api.TaskArgs) (err error) { + mgr.taskStatsMgr.ReclaimTask() + span := trace.SpanFromContextSafe(ctx) + + arg := &api.ShardTaskArgs{} + err = arg.Unmarshal(args.Data) + if err != nil { + return err + } + + if !client.ValidMigrateTask(args.TaskType, arg.TaskID) { + return errcode.ErrIllegalArguments + } + + newDst, err := base.AllocShardUnitSafe(ctx, mgr.clusterMgrCli, arg.Source, arg.Dest) + if err != nil { + span.Errorf("alloc volume unit from clustermgr failed, err: %s", err) + return err + } + + err = mgr.workQueue.Reclaim(arg.IDC, arg.TaskID, arg.Source, arg.Dest, newDst.ShardUnitInfoSimple, newDst.DiskID) + if err != nil { + span.Errorf("reclaim migrate task failed: task_type:[%s],task_id[%s], err[%+v]", mgr.taskType, arg.TaskID, err) + return err + } + + task, err := mgr.workQueue.Query(arg.IDC, arg.TaskID) + if err != nil { + span.Errorf("found task in workQueue failed: idc[%s], task_id[%s], err[%+v]", arg.IDC, arg.TaskID, err) + return err + } + t, err := task.Task() + if err != nil { + return err + } + if err = mgr.clusterMgrCli.UpdateMigrateTask(ctx, t); err != nil { + span.Errorf("update reclaim task failed: task_id[%s], err[%+v]", arg.TaskID, err) + } + return } -func (mgr *ShardMigrateMgr) RenewalTask(ctx context.Context, idc, taskID string) error { - return nil +func (mgr *ShardMigrateMgr) RenewalTask(ctx context.Context, idc, taskID string) (err error) { + if !mgr.taskSwitch.Enabled() { + return proto.ErrTaskPaused + } + + err = mgr.workQueue.Renewal(idc, taskID) + if err != nil { + span := trace.SpanFromContextSafe(ctx) + span.Warnf("renewal migrate task failed: task_type[%s], task_id[%s], err[%+v]", mgr.taskType, taskID, err) + } + return } func (mgr *ShardMigrateMgr) QueryTask(ctx context.Context, taskID string) (*api.TaskRet, error) { - return nil, nil + detail := &api.ShardTaskDetail{} + task, err := mgr.GetTask(ctx, taskID) + if err != nil { + return nil, err + } + detail.Task = *task + + // todo add statics + + data, err := json.Marshal(detail) + if err != nil { + return nil, err + } + + return &api.TaskRet{TaskType: task.TaskType, Data: data}, nil } func (mgr *ShardMigrateMgr) ReportTask(ctx context.Context, args *api.TaskArgs) (err error) { @@ -161,6 +257,10 @@ func (mgr *ShardMigrateMgr) Stats() api.ShardTaskStat { return api.ShardTaskStat{} } +func (mgr *ShardMigrateMgr) StatQueueTaskCnt() (preparing, workerDoing, finishing int) { + return +} + // Enabled task switch func (mgr *ShardMigrateMgr) Enabled() bool { return mgr.taskSwitch.Enabled() @@ -179,8 +279,62 @@ func (mgr *ShardMigrateMgr) Done() <-chan struct{} { return mgr.Closer.Done() } -func (mgr *ShardMigrateMgr) Load() error { - return nil +func (mgr *ShardMigrateMgr) Load() (err error) { + ctx := context.Background() + span := trace.SpanFromContextSafe(ctx) + + span.Infof("start load shard node migrate task: task_type[%s]", mgr.taskType) + // load task + tasks, err := mgr.clusterMgrCli.ListAllMigrateTasks(ctx, mgr.taskType) + if err != nil { + span.Errorf("list tasks from clustermgr failed: err[%+v]", err) + return + } + span.Infof("load task success: task_type[%s], tasks len[%d]", mgr.taskType, len(tasks)) + + if len(tasks) == 0 { + return + } + + var junkTasks []*proto.ShardMigrateTask + for i := range tasks { + task := &proto.ShardMigrateTask{} + err = task.Unmarshal(tasks[i].Data) + if err != nil { + return err + } + if mgr.isJunkTask(ctx, task) { + junkTasks = append(junkTasks, task) + continue + } + + if mgr.cfg.loadTaskCallback != nil { + mgr.cfg.loadTaskCallback(ctx, task.Source.DiskID) + } + + if task.Running() { + err = base.ShardTaskLockerInst().TryLock(ctx, uint32(task.Source.Suid.ShardID())) + if err != nil { + return fmt.Errorf("migrate task conflict: vid[%d], task[%+v], err[%+v]", + task.Source.Suid.ShardID(), tasks[i], err) + } + } + + span.Infof("load task success: task_type[%s], task_id[%s], state[%d]", mgr.taskType, task.TaskID, task.State) + switch task.State { + case proto.ShardTaskStateInited: + mgr.prepareQueue.PushTask(task.TaskID, task) + case proto.ShardTaskStatePrepared: + mgr.workQueue.AddPreparedTask(task.SourceIDC, task.TaskID, task) + case proto.ShardTaskStateWorkCompleted: + mgr.finishQueue.PushTask(task.TaskID, task) + case proto.ShardTaskStateFinished, proto.ShardTaskStateFinishedInAdvance: + return fmt.Errorf("task should be deleted from db: task[%+v]", task) + default: + return fmt.Errorf("unexpect migrate state: task[%+v]", task) + } + } + return mgr.clearJunkTasks(ctx, junkTasks) } func (mgr *ShardMigrateMgr) Run() { @@ -214,10 +368,322 @@ func (mgr *ShardMigrateMgr) finishTaskLoop() { } } -func (mgr *ShardMigrateMgr) prepareTask() error { +func (mgr *ShardMigrateMgr) prepareTask() (err error) { + _, task, exist := mgr.prepareQueue.PopTask() + if !exist { + return base.ErrNoTaskInQueue + } + + span, ctx := trace.StartSpanFromContext(context.Background(), "shard.migrate.prepareTask") + defer span.Finish() + + defer func() { + if err != nil { + mgr.prepareQueue.RetryTask(task.(*proto.ShardMigrateTask).TaskID) + } + }() + + taskC := *(task.(*proto.ShardMigrateTask)) + + err = base.ShardTaskLockerInst().TryLock(ctx, uint32(taskC.Source.Suid.ShardID())) + if err != nil { + span.Warnf("lock shard failed: shard_id[%v], err[%+v]", taskC.Source.Suid.ShardID(), err) + return base.ErrShardNotOnlyOneTask + } + defer func() { + if err != nil { + base.ShardTaskLockerInst().Unlock(ctx, uint32(taskC.Source.Suid.ShardID())) + } + }() + + shardInfo, err := mgr.clusterMgrCli.GetShardInfo(ctx, taskC.Source.Suid.ShardID()) + if err != nil { + span.Errorf("prepare task failed: err[%v]", err) + return err + } + + if taskC.Source.Suid != shardInfo.ShardUnitInfoSimples[taskC.Source.Suid.Index()].Suid { + span.Infof("the source unit has been moved and finish task immediately: task_id[%s], task source suid[%v], current suid[%v]", + taskC.TaskID, taskC.Source.Suid, shardInfo.ShardUnitInfoSimples[taskC.Source.Suid.Index()].Suid) + mgr.finishTaskInAdvance(ctx, &taskC) + return + } + taskC.Source = shardInfo.ShardUnitInfoSimples[taskC.Source.Suid.Index()] + + // alloc shard unit + ret, err := base.AllocShardUnitSafe(ctx, mgr.clusterMgrCli, taskC.Source, taskC.Source) + if err != nil { + span.Errorf("alloc shard unit failed: err[%+v]", err) + return + } + + taskC.Leader = shardInfo.ShardUnitInfoSimples[shardInfo.Leader] + taskC.Destination = ret.ShardUnitInfoSimple + taskC.State = proto.ShardTaskStatePrepared + taskC.Ctime = time.Now().String() + + // update db + base.InsistOn(ctx, "shard migrate prepare task update task tbl", func() error { + task, err := taskC.Task() + if err != nil { + return err + } + return mgr.clusterMgrCli.UpdateMigrateTask(ctx, task) + }) + + mgr.workQueue.AddPreparedTask(taskC.SourceIDC, taskC.TaskID, &taskC) + _ = mgr.prepareQueue.RemoveTask(taskC.TaskID) + + span.Infof("prepare task success: task_id[%s], state[%v]", taskC.TaskID, taskC.State) + return nil } -func (mgr *ShardMigrateMgr) finishTask() error { +func (mgr *ShardMigrateMgr) finishTask() (err error) { + _, task, exist := mgr.finishQueue.PopTask() + if !exist { + return base.ErrNoTaskInQueue + } + + span, ctx := trace.StartSpanFromContext(context.Background(), "migrate.finishTask") + defer span.Finish() + + defer func() { + if err != nil { + mgr.finishQueue.RetryTask(task.(*proto.ShardMigrateTask).TaskID) + } + }() + + migrateTask := *(task.(*proto.ShardMigrateTask)) + span.Infof("finish task phase: task_id[%s], state[%v]", migrateTask.TaskID, migrateTask.State) + + if migrateTask.State != proto.ShardTaskStateWorkCompleted { + span.Panicf("unexpect task state: task_id[%s], expect state[%d], actual state[%d]", + migrateTask.TaskID, proto.ShardTaskStateWorkCompleted, migrateTask.State) + } + + // because competed task did not persisted to the database, so in finish phase need to do it + // the task maybe update more than once, which is allowed + base.InsistOn(ctx, "migrate finish task update task tbl to state completed ", func() error { + task, err := migrateTask.Task() + if err != nil { + return err + } + return mgr.clusterMgrCli.UpdateMigrateTask(ctx, task) + }) + + // update shard mapping relationship + err = mgr.clusterMgrCli.UpdateShard(ctx, migrateTask.Destination.Suid, migrateTask.Source.Suid, + migrateTask.Destination.DiskID) + if err != nil { + info, err_ := mgr.clusterMgrCli.GetShardInfo(ctx, migrateTask.Source.Suid.ShardID()) + if err_ != nil { + span.Errorf("task[%s] get shard[%d] info from clustermgr failed, err: %s", + migrateTask.TaskID, migrateTask.Source.Suid.ShardID(), err_) + return err_ + } + idx := migrateTask.Source.Suid.Index() + if info.ShardUnitInfoSimples[idx] != migrateTask.Destination { + span.Errorf("change shard unit relationship failed: old suid[%d], new suid[%d], new diskId[%d], err[%+v]", + migrateTask.Source.Suid, + migrateTask.Destination.Suid, + migrateTask.Destination.DiskID, + err) + return mgr.handleUpdateShardMappingFail(ctx, &migrateTask, err) + } + } + + // remove task from clustermgr + migrateTask.State = proto.ShardTaskStateFinished + base.InsistOn(ctx, "migrate finish task update task tbl", func() error { + return mgr.clusterMgrCli.DeleteMigrateTask(ctx, migrateTask.TaskID) + }) + + _ = mgr.finishQueue.RemoveTask(migrateTask.TaskID) + //_ = mgr.updateVolumeCache(ctx, migrateTask) + + base.ShardTaskLockerInst().Unlock(ctx, uint32(migrateTask.Source.Suid.ShardID())) + // mgr.deleteMigratingVuid(migrateTask.SourceDiskID, migrateTask.SourceVuid) + + mgr.finishTaskCounter.Add() + + // add delete task and check it again + // mgr.addDeletedTask(migrateTask) + // mgr.finishTaskCallback(migrateTask.SourceDiskID) + + span.Infof("finish task phase success: task_id[%s], state[%v]", migrateTask.TaskID, migrateTask.State) + return +} + +func (mgr *ShardMigrateMgr) isJunkTask(ctx context.Context, task *proto.ShardMigrateTask) bool { + return false +} + +func (mgr *ShardMigrateMgr) clearJunkTasks(ctx context.Context, tasks []*proto.ShardMigrateTask) (err error) { return nil } + +func (mgr *ShardMigrateMgr) finishTaskInAdvance(ctx context.Context, task *proto.ShardMigrateTask) { + span := trace.SpanFromContextSafe(ctx) + span.Infof("finish task in advance: task_id[%s], task[%+v]", task.TaskID, task) + + task.State = proto.ShardTaskStateFinishedInAdvance + + base.InsistOn(ctx, "migrate finish task in advance update tbl", func() error { + return mgr.clusterMgrCli.DeleteMigrateTask(ctx, task.TaskID) + }) + + mgr.finishTaskCounter.Add() + _ = mgr.prepareQueue.RemoveTask(task.TaskID) + + base.ShardTaskLockerInst().Unlock(ctx, uint32(task.Source.Suid.ShardID())) +} + +func (mgr *ShardMigrateMgr) handleUpdateShardMappingFail(ctx context.Context, task *proto.ShardMigrateTask, err error) error { + span := trace.SpanFromContextSafe(ctx) + span.Infof("handle update shard mapping failed: task_id[%s], state[%d], dest suid[%d]", + task.TaskID, task.State, task.Destination.Suid) + + code := rpc.DetectStatusCode(err) + if code == errcode.CodeOldSuidNotMatch { + span.Panicf("change shard unit relationship failed: old vuid not match") + } + + if base.ShouldAllocShardUnitAndRedo(code) { + span.Infof("realloc shard unit and redo: task_id[%s]", task.TaskID) + newSunit, err := base.AllocShardUnitSafe(ctx, mgr.clusterMgrCli, task.Source, task.Destination) + if err != nil { + span.Errorf("realloc failed: suid[%d], err[%+v]", task.Source.Suid, err) + return err + } + task.SetDestination(newSunit.ShardUnitInfoSimple) + task.State = proto.ShardTaskStatePrepared + task.MTime = time.Now().String() + + base.InsistOn(ctx, "migrate redo task update task tbl", func() error { + t, err := task.Task() + if err != nil { + return err + } + return mgr.clusterMgrCli.UpdateMigrateTask(ctx, t) + }) + + _ = mgr.finishQueue.RemoveTask(task.TaskID) + mgr.workQueue.AddPreparedTask(task.SourceIDC, task.TaskID, task) + span.Infof("task %+v redo again", task) + + return nil + } + + return err +} + +func (mgr *ShardMigrateMgr) GetTask(ctx context.Context, id string) (task *proto.ShardMigrateTask, err error) { + span := trace.SpanFromContextSafe(ctx) + t, err := mgr.clusterMgrCli.GetMigrateTask(ctx, mgr.taskType, id) + if err != nil { + span.Errorf("get task[%s] from clustermgr failed, err: %s", id, err) + return nil, err + } + task = new(proto.ShardMigrateTask) + err = task.Unmarshal(t.Data) + if err != nil { + span.Errorf("unmarshal task[%s] from clustermgr failed, err: %s", id, err) + return nil, err + } + return task, nil +} + +// AddTask add shard migrate task, such as balance manualMigrate and drop task +func (mgr *ShardMigrateMgr) AddTask(ctx context.Context, task *proto.ShardMigrateTask) { + // add task to db + base.InsistOn(ctx, "migrate add task insert task to tbl", func() error { + t, err := task.Task() + if err != nil { + return err + } + return mgr.clusterMgrCli.AddMigrateTask(ctx, t) + }) + + // add task to prepare queue + mgr.prepareQueue.PushTask(task.TaskID, task) + + // mgr.addMigratingVuid(task.SourceDiskID, task.SourceVuid, task.TaskID) +} + +// ListMigratingSuid for disk drop and repair +func (mgr *ShardMigrateMgr) ListMigratingSuid(ctx context.Context, diskID proto.DiskID) (bads []proto.Suid, err error) { + tasks, err := mgr.clusterMgrCli.ListAllMigrateTasksByDiskID(ctx, mgr.taskType, diskID) + if err != nil { + return nil, err + } + + task := &proto.ShardMigrateTask{} + for _, t := range tasks { + err = task.Unmarshal(t.Data) + if err != nil { + return nil, err + } + bads = append(bads, task.Source.Suid) + } + return bads, nil +} + +// ListImmigratedSuid for disk drop and repair +func (mgr *ShardMigrateMgr) ListImmigratedSuid(ctx context.Context, diskID proto.DiskID) (bads []proto.Suid, err error) { + shardUnits, err := mgr.clusterMgrCli.ListDiskShardUnits(ctx, diskID) + if err != nil { + return nil, err + } + + for _, sunit := range shardUnits { + bads = append(bads, sunit.Suid) + } + return bads, nil +} + +type migratingShardDisks struct { + disks map[proto.DiskID]*client.ShardNodeDiskInfo + sync.RWMutex +} + +func newMigratingShardDisks() *migratingShardDisks { + return &migratingShardDisks{ + disks: make(map[proto.DiskID]*client.ShardNodeDiskInfo), + } +} + +func (m *migratingShardDisks) add(diskID proto.DiskID, disk *client.ShardNodeDiskInfo) { + m.Lock() + m.disks[diskID] = disk + m.Unlock() +} + +func (m *migratingShardDisks) delete(diskID proto.DiskID) { + m.Lock() + delete(m.disks, diskID) + m.Unlock() +} + +func (m *migratingShardDisks) get(diskID proto.DiskID) (disk *client.ShardNodeDiskInfo, exist bool) { + m.RLock() + disk, exist = m.disks[diskID] + m.RUnlock() + return +} + +func (m *migratingShardDisks) list() (disks []*client.ShardNodeDiskInfo) { + m.RLock() + for _, disk := range m.disks { + disks = append(disks, disk) + } + m.RUnlock() + return +} + +func (m *migratingShardDisks) size() (size int) { + m.RLock() + size = len(m.disks) + m.RUnlock() + return +} diff --git a/blobstore/scheduler/shard_migrate_test.go b/blobstore/scheduler/shard_migrate_test.go new file mode 100644 index 000000000..e172cc976 --- /dev/null +++ b/blobstore/scheduler/shard_migrate_test.go @@ -0,0 +1,385 @@ +// Copyright 2024 The CubeFS Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or +// implied. See the License for the specific language governing +// permissions and limitations under the License. + +package scheduler + +import ( + "context" + "errors" + "testing" + + "github.com/golang/mock/gomock" + "github.com/stretchr/testify/require" + + api "github.com/cubefs/cubefs/blobstore/api/scheduler" + errcode "github.com/cubefs/cubefs/blobstore/common/errors" + "github.com/cubefs/cubefs/blobstore/common/proto" + "github.com/cubefs/cubefs/blobstore/scheduler/base" + "github.com/cubefs/cubefs/blobstore/scheduler/client" + "github.com/cubefs/cubefs/blobstore/testing/mocks" +) + +var MockMigrateShardInfoMap = map[proto.ShardID]*client.ShardInfoSimple{ + 100: MockGenShardInfo(100, 0, proto.ShardStatusActive), + 101: MockGenShardInfo(101, 0, proto.ShardStatusActive), + 102: MockGenShardInfo(102, 0, proto.ShardStatusInactive), + 103: MockGenShardInfo(103, 0, proto.ShardStatusInactive), +} + +func newShardMigrateMgr(t *testing.T) *ShardMigrateMgr { + ctr := gomock.NewController(t) + clusterMgr := NewMockClusterMgrAPI(ctr) + taskSwitch := mocks.NewMockSwitcher(ctr) + + conf := &ShardMigrateConfig{ + ClusterID: 0, + TaskCommonConfig: base.TaskCommonConfig{ + PrepareQueueRetryDelayS: 0, + FinishQueueRetryDelayS: 0, + CancelPunishDurationS: 0, + WorkQueueSize: 3, + }, + } + + mgr := NewShardMigrateMgr(clusterMgr, taskSwitch, conf, proto.TaskTypeShardDiskRepair) + + shardMigrateMgr, ok := mgr.(*ShardMigrateMgr) + require.True(t, ok) + + return shardMigrateMgr +} + +func TestShardMigrateLoad(t *testing.T) { + mgr := newShardMigrateMgr(t) + + { + t1, _ := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStateInited, MockMigrateShardInfoMap).Task() + t2, _ := mockGenShardMigrateTask(101, proto.TaskTypeShardDiskRepair, "z1", 4, proto.ShardTaskStateInited, MockMigrateShardInfoMap).Task() + t3, _ := mockGenShardMigrateTask(102, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStateInited, MockMigrateShardInfoMap).Task() + t4, _ := mockGenShardMigrateTask(103, proto.TaskTypeShardDiskRepair, "z1", 4, proto.ShardTaskStateInited, MockMigrateShardInfoMap).Task() + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().ListAllMigrateTasks(any, any).Return([]*proto.Task{t1, t2, t3, t4}, nil) + err := mgr.Load() + require.NoError(t, err) + } +} + +func TestPrepareShardMigrateTask(t *testing.T) { + ctx := context.Background() + { + // no task + mgr := newShardMigrateMgr(t) + err := mgr.prepareTask() + require.True(t, errors.Is(err, base.ErrNoTaskInQueue)) + } + { + // one task and finish in advance + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStateInited, MockMigrateShardInfoMap) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AddMigrateTask(any, any).Return(nil) + mgr.AddTask(ctx, t1) + + // lock failed and send task to queue + err := base.ShardTaskLockerInst().TryLock(ctx, 100) + require.NoError(t, err) + err = mgr.prepareTask() + require.True(t, errors.Is(err, base.ErrShardNotOnlyOneTask)) + base.ShardTaskLockerInst().Unlock(ctx, 100) + + // get shard info failed + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(nil, errMock) + err = mgr.prepareTask() + require.True(t, errors.Is(err, errMock)) + + // finish task in advance because source shard unit has moved + shard := MockMigrateShardInfoMap[100] + shard.ShardUnitInfoSimples[int(t1.Source.Suid.Index())].Suid = shard.ShardUnitInfoSimples[int(t1.Source.Suid.Index())].Suid + 1 + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(shard, nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().DeleteMigrateTask(any, any).Return(nil) + err = mgr.prepareTask() + require.NoError(t, err) + } + { + // one task and normal finish + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStateInited, MockMigrateShardInfoMap) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AddMigrateTask(any, any).Return(nil) + mgr.AddTask(ctx, t1) + + // alloc shard unit failed + shard := MockMigrateShardInfoMap[100] + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(shard, nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AllocShardUnit(any, any).Return(nil, errMock) + err := mgr.prepareTask() + require.True(t, errors.Is(err, errMock)) + + // alloc success + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(shard, nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AllocShardUnit(any, any).DoAndReturn( + func(ctx context.Context, vuid proto.Suid) (*client.AllocShardUnitInfo, error) { + shardID := vuid.ShardID() + idx := vuid.Index() + epoch := vuid.Epoch() + epoch++ + newSuid := proto.EncodeSuid(shardID, idx, epoch) + return &client.AllocShardUnitInfo{ + ShardUnitInfoSimple: proto.ShardUnitInfoSimple{ + Suid: newSuid, + DiskID: shard.ShardUnitInfoSimples[idx].DiskID + 3, + Host: shard.ShardUnitInfoSimples[idx].Host, + }, + }, nil + }) + err = mgr.prepareTask() + require.NoError(t, err) + } +} + +func TestFinishShardMigrateTask(t *testing.T) { + { + // no task + mgr := newShardMigrateMgr(t) + err := mgr.finishTask() + require.True(t, errors.Is(err, base.ErrNoTaskInQueue)) + } + { + // panic :status not eql proto.MigrateStateWorkCompleted + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStateInited, MockMigrateShardInfoMap) + mgr.finishQueue.PushTask(t1.TaskID, t1) + require.Panics(t, func() { + _ = mgr.finishTask() + }) + } + { + { + // one task and redo success finally + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStateWorkCompleted, MockMigrateShardInfoMap) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + mgr.finishQueue.PushTask(t1.TaskID, t1) + + // update relationship failed + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(MockMigrateShardInfoMap[100], nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateShard(any, any, any, any).Return(errMock) + err := mgr.finishTask() + require.True(t, errors.Is(err, errMock)) + + // update relationship failed and need redo + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(MockMigrateShardInfoMap[100], nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateShard(any, any, any, any).Return(errcode.ErrNewSuidNotMatch) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AllocShardUnit(any, any).Return(nil, errMock) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + // alloc failed + err = mgr.finishTask() + require.True(t, errors.Is(err, errMock)) + + // panic + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(MockMigrateShardInfoMap[100], nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateShard(any, any, any, any).Return(errcode.ErrOldSuidNotMatch) + require.Panics(t, func() { + _ = mgr.finishTask() + }) + + // redo success + shard := MockMigrateShardInfoMap[100] + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().GetShardInfo(any, any).Return(shard, nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateShard(any, any, any, any).Return(errcode.ErrNewSuidNotMatch) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AllocShardUnit(any, any).DoAndReturn( + func(ctx context.Context, vuid proto.Suid) (*client.AllocShardUnitInfo, error) { + shardID := vuid.ShardID() + idx := vuid.Index() + epoch := vuid.Epoch() + epoch++ + newSuid := proto.EncodeSuid(shardID, idx, epoch) + return &client.AllocShardUnitInfo{ + ShardUnitInfoSimple: proto.ShardUnitInfoSimple{ + Suid: newSuid, + DiskID: shard.ShardUnitInfoSimples[idx].DiskID + 3, + Host: shard.ShardUnitInfoSimples[idx].Host, + }, + }, nil + }) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Times(2).Return(nil) + err = mgr.finishTask() + require.NoError(t, err) + } + { + // one task and success normal + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStateWorkCompleted, MockMigrateShardInfoMap) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().DeleteMigrateTask(any, any).Return(nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateShard(any, any, any, any).Return(nil) + mgr.finishQueue.PushTask(t1.TaskID, t1) + err := mgr.finishTask() + require.NoError(t, err) + } + } +} + +func TestAcquireShardMigrateTask(t *testing.T) { + ctx := context.Background() + idc := "z0" + { + // task switch is close + mgr := newShardMigrateMgr(t) + mgr.taskSwitch.(*mocks.MockSwitcher).EXPECT().Enabled().Return(false) + _, err := mgr.AcquireTask(ctx, idc) + require.True(t, errors.Is(err, proto.ErrTaskPaused)) + } + { + // no task in queue + mgr := newShardMigrateMgr(t) + mgr.taskSwitch.(*mocks.MockSwitcher).EXPECT().Enabled().Return(true) + _, err := mgr.AcquireTask(ctx, idc) + require.True(t, errors.Is(err, proto.ErrTaskEmpty)) + } + { + // one task in queue + mgr := newShardMigrateMgr(t) + mgr.taskSwitch.(*mocks.MockSwitcher).EXPECT().Enabled().Return(true) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStatePrepared, MockMigrateShardInfoMap) + mgr.workQueue.AddPreparedTask(idc, t1.TaskID, t1) + task, err := mgr.AcquireTask(ctx, idc) + require.NoError(t, err) + require.Equal(t, t1.TaskID, task.TaskID) + } +} + +func TestCancelShardMigrateTask(t *testing.T) { + ctx := context.Background() + idc := "z0" + { + mgr := newShardMigrateMgr(t) + + err := mgr.CancelTask(ctx, &api.TaskArgs{}) + require.Error(t, err) + } + { + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStatePrepared, MockMigrateShardInfoMap) + mgr.workQueue.AddPreparedTask(idc, t1.TaskID, t1) + + // no such task + err := mgr.CancelTask(ctx, &api.TaskArgs{}) + require.Error(t, err) + taskArgs := genShardTaskArgs(t1, "") + err = mgr.CancelTask(ctx, taskArgs) + require.NoError(t, err) + } +} + +func TestReclaimShardMigrateTask(t *testing.T) { + ctx := context.Background() + idc := "z0" + { + // no task + mgr := newShardMigrateMgr(t) + err := mgr.ReclaimTask(ctx, &api.TaskArgs{}) + require.Error(t, err) + } + { + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStatePrepared, MockMigrateShardInfoMap) + location := t1.Destination + location.Suid += 1 + location.DiskID += 1 + mgr.workQueue.AddPreparedTask(idc, t1.TaskID, t1) + + // update failed + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AllocShardUnit(gomock.Any(), gomock.Any()).Return( + &client.AllocShardUnitInfo{ShardUnitInfoSimple: location}, nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(errMock) + taskArgs := genShardTaskArgs(t1, "") + err := mgr.ReclaimTask(ctx, taskArgs) + require.True(t, errors.Is(err, errMock)) + + // update success + task, err := mgr.workQueue.Query(t1.SourceIDC, t1.TaskID) + require.NoError(t, err) + t1 = task.(*proto.ShardMigrateTask) + taskArgs = genShardTaskArgs(t1, "") + location = t1.Source + location.Suid += 2 + location.DiskID += 2 + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().AllocShardUnit(gomock.Any(), gomock.Any()).Return(&client.AllocShardUnitInfo{ShardUnitInfoSimple: location}, nil) + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + err = mgr.ReclaimTask(ctx, taskArgs) + require.NoError(t, err) + } +} + +func TestCompleteShardMigrateTask(t *testing.T) { + ctx := context.Background() + idc := "z0" + { + // no task + mgr := newShardMigrateMgr(t) + err := mgr.CompleteTask(ctx, &api.TaskArgs{}) + require.Error(t, err) + } + { + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", 4, proto.ShardTaskStatePrepared, MockMigrateShardInfoMap) + mgr.workQueue.AddPreparedTask(idc, t1.TaskID, t1) + + // update failed + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(errMock) + taskArgs := genShardTaskArgs(t1, "") + err := mgr.CompleteTask(ctx, taskArgs) + require.NoError(t, err) + + // no task in queue + err = mgr.CompleteTask(ctx, taskArgs) + require.Error(t, err) + + // update success + mgr.clusterMgrCli.(*MockClusterMgrAPI).EXPECT().UpdateMigrateTask(any, any).Return(nil) + t2 := mockGenShardMigrateTask(101, proto.TaskTypeShardDiskRepair, "z0", 5, proto.ShardTaskStatePrepared, MockMigrateShardInfoMap) + mgr.workQueue.AddPreparedTask(idc, t2.TaskID, t2) + args := genShardTaskArgs(t2, "") + err = mgr.CompleteTask(ctx, args) + require.NoError(t, err) + } +} + +func TestRenewalShardMigrateTask(t *testing.T) { + ctx := context.Background() + idc := "z0" + { + // task switch is close + mgr := newShardMigrateMgr(t) + mgr.taskSwitch.(*mocks.MockSwitcher).EXPECT().Enabled().Return(false) + err := mgr.RenewalTask(ctx, idc, "") + require.True(t, errors.Is(err, proto.ErrTaskPaused)) + } + { + // no task + mgr := newShardMigrateMgr(t) + mgr.taskSwitch.(*mocks.MockSwitcher).EXPECT().Enabled().Return(true) + err := mgr.RenewalTask(ctx, idc, "") + require.Error(t, err) + } + { + mgr := newShardMigrateMgr(t) + t1 := mockGenShardMigrateTask(100, proto.TaskTypeShardDiskRepair, "z0", + 4, proto.ShardTaskStatePrepared, MockMigrateShardInfoMap) + mgr.taskSwitch.(*mocks.MockSwitcher).EXPECT().Enabled().Return(true) + mgr.workQueue.AddPreparedTask(idc, t1.TaskID, t1) + err := mgr.RenewalTask(ctx, idc, t1.TaskID) + require.NoError(t, err) + } +} diff --git a/blobstore/scheduler/startup.go b/blobstore/scheduler/startup.go index 912035f61..2cf20e077 100644 --- a/blobstore/scheduler/startup.go +++ b/blobstore/scheduler/startup.go @@ -156,7 +156,7 @@ func NewService(conf *Config) (svr *Service, err error) { return nil, err } - // all migrate manager + // //===========blobnode module migrate manager=============== taskLogger, err := recordlog.NewEncoder(&conf.TaskLog) if err != nil { return nil, err @@ -191,11 +191,20 @@ func NewService(conf *Config) (svr *Service, err error) { } inspectMgr := NewVolumeInspectMgr(clusterMgrCli, mqProxy, inspectorTaskSwitch, &conf.VolumeInspect) + //===========shard module migrate manager=============== + // new shard disk repair manager + shardDiskRepairTaskSwitch, err := switchMgr.AddSwitch(proto.TaskTypeShardDiskRepair.String()) + if err != nil { + return nil, err + } + shardDiskRepairMgr := NewShardDiskRepairMgr(&conf.ShardDiskRepair, clusterMgrCli, shardDiskRepairTaskSwitch) + svr.balanceMgr = balanceMgr svr.diskDropMgr = diskDropMgr svr.manualMigMgr = manualMigMgr svr.diskRepairMgr = diskRepairMgr svr.inspectMgr = inspectMgr + svr.shardDiskRepairMgr = shardDiskRepairMgr err = svr.waitAndLoad() if err != nil { @@ -228,6 +237,9 @@ func (svr *Service) load() (err error) { if err = svr.manualMigMgr.Load(); err != nil { return } + if err = svr.shardDiskRepairMgr.Load(); err != nil { + return + } return } @@ -245,13 +257,14 @@ func (svr *Service) register(cfg ServiceRegisterConfig) error { return svr.clusterMgrCli.Register(context.Background(), info) } -// Run run task +// Run task func (svr *Service) Run() { svr.diskRepairMgr.Run() svr.balanceMgr.Run() svr.diskDropMgr.Run() svr.manualMigMgr.Run() svr.inspectMgr.Run() + svr.shardDiskRepairMgr.Run() } // RunTask run shard repair and blob delete tasks diff --git a/blobstore/scheduler/startup_test.go b/blobstore/scheduler/startup_test.go index eb8867ab5..01cb3f2d3 100644 --- a/blobstore/scheduler/startup_test.go +++ b/blobstore/scheduler/startup_test.go @@ -104,17 +104,21 @@ func newMockServiceWithOpts(ctr *gomock.Controller, isLeader bool) *Service { clusterTopology := NewMockClusterTopology(ctr) volumeUpdater := NewMockVolumeUpdater(ctr) + shardDiskRepair := NewMockShardDisMigrator(ctr) + balanceMgr.EXPECT().Close().AnyTimes().Return() diskRepairMgr.EXPECT().Close().AnyTimes().Return() diskDropMgr.EXPECT().Close().AnyTimes().Return() manualMgr.EXPECT().Close().AnyTimes().Return() inspecterMgr.EXPECT().Close().AnyTimes().Return() + shardDiskRepair.EXPECT().Close().AnyTimes().Return() balanceMgr.EXPECT().Run().AnyTimes().Return() diskDropMgr.EXPECT().Run().AnyTimes().Return() diskRepairMgr.EXPECT().Run().AnyTimes().Return() inspecterMgr.EXPECT().Run().AnyTimes().Return() manualMgr.EXPECT().Run().AnyTimes().Return() + shardDiskRepair.EXPECT().Run().AnyTimes().Return() clusterTopology.EXPECT().LoadVolumes().AnyTimes().Return(nil) shardRepairMgr.EXPECT().Run().AnyTimes().Return() @@ -126,6 +130,7 @@ func newMockServiceWithOpts(ctr *gomock.Controller, isLeader bool) *Service { diskRepairMgr.EXPECT().Load().AnyTimes().Return(nil) diskDropMgr.EXPECT().Load().AnyTimes().Return(nil) manualMgr.EXPECT().Load().AnyTimes().Return(nil) + shardDiskRepair.EXPECT().Load().AnyTimes().Return(nil) blobDeleteMgr.EXPECT().GetErrorStats().AnyTimes().Return([]string{}, uint64(0)) blobDeleteMgr.EXPECT().GetTaskStats().AnyTimes().Return([counter.SLOT]int{}, [counter.SLOT]int{}) @@ -144,6 +149,9 @@ func newMockServiceWithOpts(ctr *gomock.Controller, isLeader bool) *Service { manualMgr.EXPECT().Stats().AnyTimes().Return(api.MigrateTasksStat{}) inspecterMgr.EXPECT().GetTaskStats().AnyTimes().Return([counter.SLOT]int{}, [counter.SLOT]int{}) inspecterMgr.EXPECT().Enabled().AnyTimes().Return(true) + shardDiskRepair.EXPECT().Stats().AnyTimes().Return(api.ShardTaskStat{}) + shardDiskRepair.EXPECT().Progress(any).AnyTimes().Return([]proto.DiskID{proto.DiskID(1)}, 0, 0) + shardDiskRepair.EXPECT().Enabled().AnyTimes().Return(true) volumeUpdater.EXPECT().UpdateFollowerVolumeCache(any, any, any).AnyTimes().Return(nil) volumeUpdater.EXPECT().UpdateLeaderVolumeCache(any, any).AnyTimes().Return(nil) @@ -152,23 +160,25 @@ func newMockServiceWithOpts(ctr *gomock.Controller, isLeader bool) *Service { diskRepairMgr.EXPECT().AcquireTask(any, any).AnyTimes().Return(&proto.Task{}, errMock) diskDropMgr.EXPECT().AcquireTask(any, any).AnyTimes().Return(&proto.Task{}, errMock) balanceMgr.EXPECT().AcquireTask(any, any).AnyTimes().Return(&proto.Task{}, errMock) + shardDiskRepair.EXPECT().AcquireTask(any, any).AnyTimes().Return(&proto.Task{}, errMock) clusterTopology.EXPECT().UpdateVolume(any).AnyTimes().Return(&client.VolumeInfoSimple{}, nil) clusterMgrCli.EXPECT().GetConfig(any, any).AnyTimes().Return("", errMock) service := &Service{ - ClusterID: 1, - leader: isLeader, - balanceMgr: balanceMgr, - diskDropMgr: diskDropMgr, - manualMigMgr: manualMgr, - diskRepairMgr: diskRepairMgr, - inspectMgr: inspecterMgr, - shardRepairMgr: shardRepairMgr, - blobDeleteMgr: blobDeleteMgr, - clusterTopology: clusterTopology, - volumeUpdater: volumeUpdater, - clusterMgrCli: clusterMgrCli, + ClusterID: 1, + leader: isLeader, + balanceMgr: balanceMgr, + diskDropMgr: diskDropMgr, + manualMigMgr: manualMgr, + diskRepairMgr: diskRepairMgr, + inspectMgr: inspecterMgr, + shardRepairMgr: shardRepairMgr, + blobDeleteMgr: blobDeleteMgr, + clusterTopology: clusterTopology, + volumeUpdater: volumeUpdater, + clusterMgrCli: clusterMgrCli, + shardDiskRepairMgr: shardDiskRepair, } return service }