refact(blobnode): we need re-add all disk, when first add node

Signed-off-by: mawei029 <mawei2@oppo.com>

with: #22265425 of #22160501
This commit is contained in:
mawei029 2024-06-05 16:17:25 +08:00 committed by 梁曟風
parent 4c8f8892be
commit 399613b0f9
4 changed files with 26 additions and 21 deletions

View File

@ -26,6 +26,7 @@ import (
"github.com/cubefs/cubefs/blobstore/blobnode/core"
"github.com/cubefs/cubefs/blobstore/blobnode/db"
"github.com/cubefs/cubefs/blobstore/cmd"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/rpc"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/util/defaulter"
@ -110,6 +111,7 @@ func configInit(config *Config) {
}
defaulter.LessOrEqual(&config.InspectConf.IntervalSec, DefaultChunkInspectIntervalSec)
defaulter.LessOrEqual(&config.InspectConf.RateLimit, DefaultInspectRate)
defaulter.LessOrEqual(&config.HostInfo.DiskType, proto.DiskTypeHDD)
}
func (s *Service) changeLimit(ctx context.Context, c Config) {

View File

@ -96,6 +96,7 @@ type HostInfo struct {
Host string `json:"host"`
DiskType proto.DiskType `json:"disk_type,omitempty"` // On a node, there is only one type of disk, and no other types
NodeID proto.NodeID `json:"-"` // A node is a process
ReAddDisk bool `json:"re_add_disk"` // need to re-register all disks under the node. temp switch
}
type Config struct {
@ -159,7 +160,7 @@ func InitConfig(conf *Config) error {
}
func CheckNodeConf(conf *HostInfo) error {
if conf.DiskType == 0 {
if conf.DiskType <= 0 {
return errors.New("disk type is not specified")
}

View File

@ -140,10 +140,10 @@ func (s *Service) handleDiskIOError(ctx context.Context, diskID proto.DiskID, di
err := s.ClusterMgrClient.SetDisk(ctx, diskID, proto.DiskStatusBroken)
// error is nil or already broken status
if err == nil || rpc.DetectStatusCode(err) == bloberr.CodeChangeDiskStatusNotAllow {
span.Infof("set disk(%d) broken success, err:%v", diskID, err)
span.Infof("set disk(%d) broken success, err:%+v", diskID, err)
break
}
span.Errorf("set disk(%d) broken failed: %v", diskID, err)
span.Errorf("set disk(%d) broken failed: %+v", diskID, err)
time.Sleep(3 * time.Second)
}
@ -162,7 +162,7 @@ func (s *Service) handleDiskIOError(ctx context.Context, diskID proto.DiskID, di
return nil, nil
})
span.Debugf("diskID:%d diskErr: %v, shared:%v", diskID, diskErr, shared)
span.Debugf("diskID:%d diskErr: %+v, shared:%v", diskID, diskErr, shared)
}
func (s *Service) waitRepairAndClose(ctx context.Context, disk core.DiskAPI) {
@ -182,7 +182,7 @@ func (s *Service) waitRepairAndClose(ctx context.Context, disk core.DiskAPI) {
info, err := s.ClusterMgrClient.DiskInfo(ctx, diskID)
if err != nil {
span.Errorf("Failed get clustermgr diskinfo %d, err:%v", diskID, err)
span.Errorf("Failed get clustermgr diskinfo %d, err:%+v", diskID, err)
continue
}
@ -277,7 +277,7 @@ func NewService(conf Config) (svr *Service, err error) {
}
err = clusterMgrCli.RegisterService(ctx, node, TickInterval, HeartbeatTicks, ExpiresTicks)
if err != nil {
span.Fatalf("blobnode register to clusterMgr error:%v", err)
span.Fatalf("blobnode register to clusterMgr error:%+v", err)
}
if err = registerNode(ctx, clusterMgrCli, &conf); err != nil {
@ -287,7 +287,7 @@ func NewService(conf Config) (svr *Service, err error) {
registeredDisks, err := clusterMgrCli.ListHostDisk(ctx, conf.Host)
if err != nil {
span.Errorf("Failed ListDisk from clusterMgr. err:%v", err)
span.Errorf("Failed ListDisk from clusterMgr. err:%+v", err)
return nil, err
}
span.Infof("registered disks: %v", registeredDisks)
@ -341,7 +341,7 @@ func NewService(conf Config) (svr *Service, err error) {
lost := atomic.AddInt32(&lostCnt, 1)
svr.reportLostDisk(&diskConf.HostInfo, diskConf.Path) // startup check lost disk
// skip
span.Errorf("Path is not mount point:%s, err:%v. skip init", diskConf.Path, err)
span.Errorf("Path is not mount point:%s, err:%+v. skip init", diskConf.Path, err)
if lost >= LostDiskCount {
log.Fatalf("lost disk count:%d over threshold:%d", lost, LostDiskCount)
}
@ -351,7 +351,7 @@ func NewService(conf Config) (svr *Service, err error) {
format, err := readFormatInfo(ctx, diskConf.Path)
if err != nil {
// todo: report to ums
span.Errorf("Failed read diskMeta:%s, err:%v. skip init", diskConf.Path, err)
span.Errorf("Failed read diskMeta:%s, err:%+v. skip init", diskConf.Path, err)
err = nil // skip
return
}
@ -365,22 +365,23 @@ func NewService(conf Config) (svr *Service, err error) {
nonNormal := foundInCluster && diskInfo.Status != proto.DiskStatusNormal
if nonNormal {
// todo: report to ums
span.Warnf("disk(%v):path(%v) is not normal, skip init", format.DiskID, diskConf.Path)
span.Warnf("disk(%d):path(%s) is not normal, skip init", format.DiskID, diskConf.Path)
return
}
ds, err := disk.NewDiskStorage(svr.ctx, diskConf)
if err != nil {
span.Errorf("Failed Open DiskStorage. conf:%v, err:%v", diskConf, err)
span.Errorf("Failed Open DiskStorage. conf:%v, err:%+v", diskConf, err)
return
}
if !foundInCluster {
span.Warnf("diskInfo:%v not found in clusterMgr, will register to cluster", diskInfo)
if !foundInCluster || conf.HostInfo.ReAddDisk { // need to re-register all disks
span.Warnf("diskInfo:%v not found in cm, will register to cm, nodeID:%d", diskInfo, conf.NodeID)
diskInfo := ds.DiskInfo() // get nodeID to add disk
err := clusterMgrCli.AddDisk(ctx, &diskInfo)
if err != nil {
span.Errorf("Failed register disk: %v, err:%v", diskInfo, err)
// if it need re-register disk, it is necessary to ignore duplicate registrations
if err != nil && (conf.HostInfo.ReAddDisk && rpc.DetectStatusCode(err) != http.StatusCreated) {
span.Errorf("Failed register disk: %v, err:%+v", diskInfo, err)
return
}
}
@ -390,7 +391,7 @@ func NewService(conf Config) (svr *Service, err error) {
svr.lock.Unlock()
svr.reportOnlineDisk(&diskConf.HostInfo, diskConf.Path) // restart, normal disk
span.Infof("Init disk storage, cluster:%v, diskID:%v", conf.ClusterID, format.DiskID)
span.Infof("Init disk storage, cluster:%d, formatID:%d, diskID:%d", conf.ClusterID, format.DiskID, ds.ID())
}(diskConf)
}
wg.Wait()
@ -457,6 +458,6 @@ func registerNode(ctx context.Context, clusterMgrCli *cmapi.Client, conf *Config
return err
}
conf.NodeID = nodeID
conf.NodeID = nodeID // we update nodeID, which can be used in the subsequent process. e.g. to add disk
return nil
}

View File

@ -936,12 +936,9 @@ func TestService_RegisterNode(t *testing.T) {
Conf: &conf,
}
err := registerNode(ctx, svr.ClusterMgrClient, svr.Conf)
require.NotNil(t, err)
// first register
svr.Conf.DiskType = proto.DiskTypeHDD
err = registerNode(ctx, svr.ClusterMgrClient, svr.Conf)
err := registerNode(ctx, svr.ClusterMgrClient, svr.Conf)
require.NoError(t, err)
require.Equal(t, proto.NodeID(1), svr.Conf.HostInfo.NodeID)
@ -960,4 +957,8 @@ func TestService_RegisterNode(t *testing.T) {
require.NoError(t, err)
require.Equal(t, proto.NodeID(2), svr2.Conf.HostInfo.NodeID)
require.NotEqual(t, svr.Conf.NodeID, svr2.Conf.NodeID)
svr.Conf.DiskType = 0
err = registerNode(ctx, svr.ClusterMgrClient, svr.Conf)
require.NotNil(t, err)
}