feat(scheduler): add shard disk repair manager

1. finish shard task process and shard disk repair

with: #22452167

Signed-off-by: JasonHu520 <huzongchao@oppo.com>
This commit is contained in:
JasonHu520 2024-07-31 17:12:54 +08:00 committed by slasher
parent 19e3ec09e4
commit 67277bdf0e
27 changed files with 2141 additions and 156 deletions

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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