From 399613b0f9edfd9dfce600dd75e35b19e00a54d8 Mon Sep 17 00:00:00 2001 From: mawei029 Date: Wed, 5 Jun 2024 16:17:25 +0800 Subject: [PATCH] refact(blobnode): we need re-add all disk, when first add node Signed-off-by: mawei029 with: #22265425 of #22160501 --- blobstore/blobnode/config.go | 2 ++ blobstore/blobnode/core/config.go | 3 ++- blobstore/blobnode/startup.go | 33 ++++++++++++++++--------------- blobstore/blobnode/svr_test.go | 9 +++++---- 4 files changed, 26 insertions(+), 21 deletions(-) diff --git a/blobstore/blobnode/config.go b/blobstore/blobnode/config.go index f815ded22..d5dd5eac7 100644 --- a/blobstore/blobnode/config.go +++ b/blobstore/blobnode/config.go @@ -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) { diff --git a/blobstore/blobnode/core/config.go b/blobstore/blobnode/core/config.go index 2e63bbf86..9bb6cd1f5 100644 --- a/blobstore/blobnode/core/config.go +++ b/blobstore/blobnode/core/config.go @@ -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") } diff --git a/blobstore/blobnode/startup.go b/blobstore/blobnode/startup.go index 95147b10d..3b516e9f4 100644 --- a/blobstore/blobnode/startup.go +++ b/blobstore/blobnode/startup.go @@ -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 } diff --git a/blobstore/blobnode/svr_test.go b/blobstore/blobnode/svr_test.go index a6dfde750..4166fb4d0 100644 --- a/blobstore/blobnode/svr_test.go +++ b/blobstore/blobnode/svr_test.go @@ -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) }