From c113ecfa69019f69121bafdf864335c55070b006 Mon Sep 17 00:00:00 2001 From: mawei029 Date: Thu, 3 Jul 2025 16:41:52 +0800 Subject: [PATCH] refactor(blobnode): added multiple startup of disk ut individual tests with: #1000212202 Signed-off-by: mawei029 --- blobstore/api/clustermgr/proto.go | 20 ++ blobstore/blobnode/startup.go | 232 ++++++++--------- blobstore/blobnode/svr.go | 2 +- blobstore/blobnode/svr_test.go | 302 +++++++++++++++++++++- blobstore/testing/mocks/api_clustermgr.go | 145 +++++++++++ 5 files changed, 581 insertions(+), 120 deletions(-) diff --git a/blobstore/api/clustermgr/proto.go b/blobstore/api/clustermgr/proto.go index c4d5c9461..35d33d243 100644 --- a/blobstore/api/clustermgr/proto.go +++ b/blobstore/api/clustermgr/proto.go @@ -53,6 +53,7 @@ func GetConsulClusterPath(region string) string { type ClientAPI interface { APIAccess APIProxy + APIBlobnode } // APIAccess sub of cluster manager api for access @@ -81,4 +82,23 @@ type APIProxy interface { // APIService sub of cluster manager api for service type APIService interface { GetService(ctx context.Context, args GetServiceArgs) (ServiceInfo, error) + RegisterService(ctx context.Context, node ServiceNode, tickInterval, heartbeatTicks, expiresTicks uint32) (err error) +} + +type APIBlobnode interface { + APIService + GetConfig(ctx context.Context, key string) (value string, err error) + SetConfig(ctx context.Context, key, value string) error + AddNode(ctx context.Context, info *BlobNodeInfo) (proto.NodeID, error) + ListHostDisk(ctx context.Context, host string) (ret []*BlobNodeDiskInfo, err error) + ListDisk(ctx context.Context, options *ListOptionArgs) (ret ListDiskRet, err error) + AddDisk(ctx context.Context, info *BlobNodeDiskInfo) (err error) + DiskInfo(ctx context.Context, id proto.DiskID) (ret *BlobNodeDiskInfo, err error) + SetDisk(ctx context.Context, id proto.DiskID, status proto.DiskStatus) (err error) + AllocDiskID(ctx context.Context) (proto.DiskID, error) + SetCompactChunk(ctx context.Context, args *SetCompactChunkArgs) (err error) + ListVolumeUnit(ctx context.Context, args *ListVolumeUnitArgs) ([]*VolumeUnitInfo, error) + GetVolumeInfo(ctx context.Context, args *GetVolumeArgs) (ret *VolumeInfo, err error) + ReportChunk(ctx context.Context, args *ReportChunkArgs) (err error) + HeartbeatDisk(ctx context.Context, infos []*DiskHeartBeatInfo) (ret []*DiskHeartbeatRet, err error) } diff --git a/blobstore/blobnode/startup.go b/blobstore/blobnode/startup.go index e6ac40056..d1b4d84f0 100644 --- a/blobstore/blobnode/startup.go +++ b/blobstore/blobnode/startup.go @@ -49,44 +49,32 @@ const ( LostDiskCount = 3 ) -func readFormatInfo(ctx context.Context, diskRootPath string) ( - formatInfo *core.FormatInfo, err error, -) { - span := trace.SpanFromContextSafe(ctx) - _, err = os.ReadDir(diskRootPath) - if err != nil { - span.Errorf("read disk root path error:%s", diskRootPath) - return nil, err +func NewService(conf Config) (svr *Service, err error) { + _, ctx := trace.StartSpanFromContext(context.Background(), "NewBlobNodeService") + + configInit(&conf) + + clusterMgrCli := cmapi.New(conf.Clustermgr) + + svr = &Service{ + ClusterMgrClient: clusterMgrCli, + Disks: make(map[proto.DiskID]core.DiskAPI), + Conf: &conf, + closeCh: make(chan struct{}), } - formatInfo, err = core.ReadFormatInfo(ctx, diskRootPath) - if err != nil { - if os.IsNotExist(err) { - span.Warnf("format file not exist. must be first register") - return new(core.FormatInfo), nil - } - return nil, err + svr.ctx, svr.cancel = context.WithCancel(context.Background()) + + // start worker service + if conf.StartMode == proto.ServiceNameWorker || conf.StartMode == defaultServiceBothBlobNodeWorker { + startWorkerService(ctx, svr, conf) } - return formatInfo, err -} + // start blobndoe service + if conf.StartMode == proto.ServiceNameBlobNode || conf.StartMode == defaultServiceBothBlobNodeWorker { + err = startBlobnodeService(ctx, svr, conf) + } -func isAllInConfig(ctx context.Context, registeredDisks []*cmapi.BlobNodeDiskInfo, conf *Config) bool { - span := trace.SpanFromContextSafe(ctx) - configDiskMap := make(map[string]struct{}) - for i := range conf.Disks { - configDiskMap[conf.Disks[i].Path] = struct{}{} - } - // check all registered normal disks are in config - for _, registeredDisk := range registeredDisks { - if registeredDisk.Status != proto.DiskStatusNormal { - continue - } - if _, ok := configDiskMap[registeredDisk.Path]; !ok { - span.Errorf("disk registered to clustermgr, but is not in config: %v", registeredDisk.Path) - return false - } - } - return true + return } // call by heartbeat single, or datafile read/write concurrence @@ -267,15 +255,6 @@ func (s *Service) handleDiskDrop(ctx context.Context, ds core.DiskAPI) { }() } -func setDefaultIOStat(dryRun bool) error { - ios, err := flow.NewIOFlowStat("default", dryRun) - if err != nil { - return errors.New("init stat failed") - } - flow.SetupDefaultIOStat(ios) - return nil -} - func (s *Service) fixDiskConf(config *core.Config) { config.AllocDiskID = s.ClusterMgrClient.AllocDiskID config.NotifyCompacting = s.ClusterMgrClient.SetCompactChunk @@ -315,47 +294,95 @@ func (s *Service) handleStartDiskError(ctx context.Context, allUniqDiskPathMap m if diskInfo.Status != proto.DiskStatusNormal { span.Warnf("disk[id:%d,status:%d,path:%s] is not normal, err:%+v. skip init", diskInfo.DiskID, diskInfo.Status, diskPath, err) return - } else { - // set broken disk, may be already mark broken - span.Errorf("open normal disk[%d:%s] failed, err:%+v. skip init", diskInfo.DiskID, diskPath, err) - _err := s.ClusterMgrClient.SetDisk(ctx, diskInfo.DiskID, proto.DiskStatusBroken) - if _err != nil && rpc.DetectStatusCode(_err) != bloberr.CodeChangeDiskStatusNotAllow { - span.Fatalf("set disk[%d:%s] broken to cm failed: %s", diskInfo.DiskID, diskPath, _err) - } - return + } + + // set broken disk, may be already mark broken + span.Errorf("open normal disk[%d:%s] failed, err:%+v. skip init", diskInfo.DiskID, diskPath, err) + _err := s.ClusterMgrClient.SetDisk(ctx, diskInfo.DiskID, proto.DiskStatusBroken) + if _err != nil && rpc.DetectStatusCode(_err) != bloberr.CodeChangeDiskStatusNotAllow { + span.Fatalf("set disk[%d:%s] broken to cm failed: %s", diskInfo.DiskID, diskPath, _err) + } + return + } +} + +func (s *Service) registerNode(ctx context.Context, conf *Config) error { + span := trace.SpanFromContextSafe(ctx) + if err := core.CheckNodeConf(&conf.HostInfo); err != nil { + return err + } + + nodeToCm := cmapi.BlobNodeInfo{ + NodeInfo: cmapi.NodeInfo{ + ClusterID: conf.ClusterID, + DiskType: conf.DiskType, + Idc: conf.IDC, + Rack: conf.Rack, + Host: conf.Host, + Role: proto.NodeRoleBlobNode, + }, + } + + nodeID, err := s.ClusterMgrClient.AddNode(ctx, &nodeToCm) + if err != nil && rpc.DetectStatusCode(err) != http.StatusCreated { + return err + } + + conf.NodeID = nodeID // we update nodeID, which can be used in the subsequent process. e.g. to add disk + span.Infof("add node success, nodeID=%d", nodeID) + return nil +} + +func setDefaultIOStat(dryRun bool) error { + ios, err := flow.NewIOFlowStat("default", dryRun) + if err != nil { + return errors.New("init stat failed") + } + flow.SetupDefaultIOStat(ios) + return nil +} + +func readFormatInfo(ctx context.Context, diskRootPath string) ( + formatInfo *core.FormatInfo, err error, +) { + span := trace.SpanFromContextSafe(ctx) + _, err = os.ReadDir(diskRootPath) + if err != nil { + span.Errorf("read disk root path error:%s", diskRootPath) + return nil, err + } + formatInfo, err = core.ReadFormatInfo(ctx, diskRootPath) + if err != nil { + if os.IsNotExist(err) { + span.Warnf("format file not exist. must be first register") + return new(core.FormatInfo), nil + } + return nil, err + } + + return formatInfo, err +} + +func isAllInConfig(ctx context.Context, registeredDisks []*cmapi.BlobNodeDiskInfo, conf *Config) bool { + span := trace.SpanFromContextSafe(ctx) + configDiskMap := make(map[string]struct{}) + for i := range conf.Disks { + configDiskMap[conf.Disks[i].Path] = struct{}{} + } + // check all registered normal disks are in config + for _, registeredDisk := range registeredDisks { + if registeredDisk.Status != proto.DiskStatusNormal { + continue + } + if _, ok := configDiskMap[registeredDisk.Path]; !ok { + span.Errorf("disk registered to clustermgr, but is not in config: %v", registeredDisk.Path) + return false } } + return true } -func NewService(conf Config) (svr *Service, err error) { - _, ctx := trace.StartSpanFromContext(context.Background(), "NewBlobNodeService") - - configInit(&conf) - - clusterMgrCli := cmapi.New(conf.Clustermgr) - - svr = &Service{ - ClusterMgrClient: clusterMgrCli, - Disks: make(map[proto.DiskID]core.DiskAPI), - Conf: &conf, - closeCh: make(chan struct{}), - } - svr.ctx, svr.cancel = context.WithCancel(context.Background()) - - // start worker service - if conf.StartMode == proto.ServiceNameWorker || conf.StartMode == defaultServiceBothBlobNodeWorker { - startWorkerService(ctx, svr, conf, clusterMgrCli) - } - - // start blobndoe service - if conf.StartMode == proto.ServiceNameBlobNode || conf.StartMode == defaultServiceBothBlobNodeWorker { - err = startBlobnodeService(ctx, svr, conf, clusterMgrCli) - } - - return -} - -func startWorkerService(ctx context.Context, svr *Service, conf Config, clusterMgrCli *cmapi.Client) { +func startWorkerService(ctx context.Context, svr *Service, conf Config) { span := trace.SpanFromContextSafe(ctx) span.Debug("start worker service...") @@ -366,18 +393,18 @@ func startWorkerService(ctx context.Context, svr *Service, conf Config, clusterM Idc: conf.IDC, } - err := clusterMgrCli.RegisterService(ctx, node, TickInterval, HeartbeatTicks, ExpiresTicks) + err := svr.ClusterMgrClient.RegisterService(ctx, node, TickInterval, HeartbeatTicks, ExpiresTicks) if err != nil { span.Fatalf("worker register to clusterMgr error:%+v", err) } - svr.WorkerService, err = NewWorkerService(&conf.WorkerConfig, clusterMgrCli, conf.ClusterID, conf.IDC) + svr.WorkerService, err = NewWorkerService(&conf.WorkerConfig, svr.ClusterMgrClient, conf.ClusterID, conf.IDC) if err != nil { span.Fatalf("Failed to new worker service, err: %v", err) } } -func startBlobnodeService(ctx context.Context, svr *Service, conf Config, clusterMgrCli *cmapi.Client) (err error) { +func startBlobnodeService(ctx context.Context, svr *Service, conf Config) (err error) { span := trace.SpanFromContextSafe(ctx) span.Debug("start blobnode service...") @@ -387,7 +414,7 @@ func startBlobnodeService(ctx context.Context, svr *Service, conf Config, cluste Host: conf.Host, Idc: conf.IDC, } - if err = cmapi.LoadExtendCodemode(ctx, clusterMgrCli); err != nil { + if err = cmapi.LoadExtendCodemode(ctx, svr.ClusterMgrClient); err != nil { span.Fatalf("load extend codemode from clusterMgr error:%+v", err) } for _, ecmode := range codemode.GetECCodeModes() { @@ -396,16 +423,16 @@ func startBlobnodeService(ctx context.Context, svr *Service, conf Config, cluste } } - err = clusterMgrCli.RegisterService(ctx, node, TickInterval, HeartbeatTicks, ExpiresTicks) + err = svr.ClusterMgrClient.RegisterService(ctx, node, TickInterval, HeartbeatTicks, ExpiresTicks) if err != nil { span.Fatalf("blobnode register to clusterMgr error:%+v", err) } - if err = registerNode(ctx, clusterMgrCli, &conf); err != nil { + if err = svr.registerNode(ctx, &conf); err != nil { span.Fatalf("fail to register node to clusterMgr, err:%+v", err) } - registeredDisks, err := clusterMgrCli.ListHostDisk(ctx, conf.Host) + registeredDisks, err := svr.ClusterMgrClient.ListHostDisk(ctx, conf.Host) if err != nil { span.Errorf("Failed ListDisk from clusterMgr. err:%+v", err) return err @@ -427,7 +454,7 @@ func startBlobnodeService(ctx context.Context, svr *Service, conf Config, cluste svr.InspectLimiterPerKey = keycount.New(1) svr.BrokenLimitPerDisk = keycount.New(1) - switchMgr := taskswitch.NewSwitchMgr(clusterMgrCli) + switchMgr := taskswitch.NewSwitchMgr(svr.ClusterMgrClient) svr.inspectMgr, err = NewDataInspectMgr(svr, conf.InspectConf, switchMgr) if err != nil { return err @@ -508,7 +535,7 @@ func startBlobnodeService(ctx context.Context, svr *Service, conf Config, cluste if format.DiskID == 0 || !foundIDInCluster { span.Warnf("diskInfo:%v not found in cm, will register to cm, nodeID:%d", diskInfo, conf.NodeID) dsInfo := ds.DiskInfo() // get nodeID to add disk - err = clusterMgrCli.AddDisk(ctx, &dsInfo) + err = svr.ClusterMgrClient.AddDisk(ctx, &dsInfo) if err != nil { span.Fatalf("Failed register disk: %v, err:%+v", dsInfo, err) return @@ -554,30 +581,3 @@ func startBlobnodeService(ctx context.Context, svr *Service, conf Config, cluste return } - -func registerNode(ctx context.Context, clusterMgrCli *cmapi.Client, conf *Config) error { - span := trace.SpanFromContextSafe(ctx) - if err := core.CheckNodeConf(&conf.HostInfo); err != nil { - return err - } - - nodeToCm := cmapi.BlobNodeInfo{ - NodeInfo: cmapi.NodeInfo{ - ClusterID: conf.ClusterID, - DiskType: conf.DiskType, - Idc: conf.IDC, - Rack: conf.Rack, - Host: conf.Host, - Role: proto.NodeRoleBlobNode, - }, - } - - nodeID, err := clusterMgrCli.AddNode(ctx, &nodeToCm) - if err != nil && rpc.DetectStatusCode(err) != http.StatusCreated { - return err - } - - conf.NodeID = nodeID // we update nodeID, which can be used in the subsequent process. e.g. to add disk - span.Infof("add node success, nodeID=%d", nodeID) - return nil -} diff --git a/blobstore/blobnode/svr.go b/blobstore/blobnode/svr.go index 6ccb45281..f891d297e 100644 --- a/blobstore/blobnode/svr.go +++ b/blobstore/blobnode/svr.go @@ -44,7 +44,7 @@ type Service struct { WorkerService *WorkerService // client handler - ClusterMgrClient *cmapi.Client + ClusterMgrClient cmapi.APIBlobnode Conf *Config inspectMgr *DataInspectMgr diff --git a/blobstore/blobnode/svr_test.go b/blobstore/blobnode/svr_test.go index 508bcd828..24c103e23 100644 --- a/blobstore/blobnode/svr_test.go +++ b/blobstore/blobnode/svr_test.go @@ -17,19 +17,24 @@ package blobnode import ( "context" "encoding/json" + "fmt" "io" "math" "net/http" "net/http/httptest" "os" "path/filepath" + "reflect" "strings" "sync" "sync/atomic" + "syscall" "testing" "time" + "github.com/agiledragon/gomonkey/v2" "github.com/golang/mock/gomock" + "github.com/opentracing/opentracing-go" "github.com/stretchr/testify/require" bnapi "github.com/cubefs/cubefs/blobstore/api/blobnode" @@ -44,6 +49,7 @@ import ( "github.com/cubefs/cubefs/blobstore/common/recordlog" "github.com/cubefs/cubefs/blobstore/common/rpc" "github.com/cubefs/cubefs/blobstore/common/trace" + "github.com/cubefs/cubefs/blobstore/testing/mocks" "github.com/cubefs/cubefs/blobstore/util/errors" "github.com/cubefs/cubefs/blobstore/util/log" ) @@ -1140,7 +1146,7 @@ func TestService_RegisterNode(t *testing.T) { // first register svr.Conf.DiskType = proto.DiskTypeHDD - err := registerNode(ctx, svr.ClusterMgrClient, svr.Conf) + err := svr.registerNode(ctx, svr.Conf) require.NoError(t, err) require.Equal(t, proto.NodeID(1), svr.Conf.HostInfo.NodeID) @@ -1155,13 +1161,13 @@ func TestService_RegisterNode(t *testing.T) { Conf: &conf2, } svr2.Conf.DiskType = proto.DiskTypeSSD - err = registerNode(ctx, svr2.ClusterMgrClient, svr2.Conf) + err = svr2.registerNode(ctx, svr2.Conf) 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) + err = svr.registerNode(ctx, svr.Conf) require.NotNil(t, err) } @@ -1242,3 +1248,293 @@ func TestService_OnlyBlobnode(t *testing.T) { _, err = NewService(conf) require.NoError(t, err) } + +func TestService_OnlyBlobnode_OpenFailedEIO(t *testing.T) { + ctx := context.Background() + workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"OnlyBlobnode") + require.NoError(t, err) + defer os.RemoveAll(workDir) + + path1 := filepath.Join(workDir, "path1") + path2 := filepath.Join(workDir, "path2") + path3 := filepath.Join(workDir, "path3") + path4 := filepath.Join(workDir, "path4") + for _, path := range []string{workDir, path1, path2, path3, path4} { + err = os.MkdirAll(path, 0o755) + require.NoError(t, err) + } + + conf := Config{ + HostInfo: core.HostInfo{ + IDC: "testIdc", + Rack: "testRack", + DiskType: proto.DiskTypeHDD, + }, + Disks: []core.Config{ + {BaseConfig: core.BaseConfig{Path: path1, AutoFormat: true, MaxChunks: 700}, MetaConfig: db.MetaConfig{}}, + {BaseConfig: core.BaseConfig{Path: path2, AutoFormat: true, MaxChunks: 700}, MetaConfig: db.MetaConfig{}}, + }, + DiskConfig: core.RuntimeConfig{DiskReservedSpaceB: 1, CompactReservedSpaceB: 1}, + HeartbeatIntervalSec: 600, + InspectConf: DataInspectConf{Record: recordlog.Config{Dir: filepath.Join(workDir, "inspect")}}, + } + + // open readFormat eio, report broken disk + diskInfo1 := &cmapi.BlobNodeDiskInfo{ + DiskInfo: cmapi.DiskInfo{ + Path: path1, + Status: proto.DiskStatusNormal, + }, + DiskHeartBeatInfo: cmapi.DiskHeartBeatInfo{DiskID: proto.DiskID(1)}, + } + // open readFormat eio, status repaired, skip + diskInfo2 := &cmapi.BlobNodeDiskInfo{ + DiskInfo: cmapi.DiskInfo{ + Path: path2, + Status: proto.DiskStatusRepaired, + }, + DiskHeartBeatInfo: cmapi.DiskHeartBeatInfo{DiskID: proto.DiskID(2)}, + } + + A := gomock.Any() + ctr := gomock.NewController(t) + cmCli := mocks.NewMockClientAPI(ctr) + cmCli.EXPECT().GetConfig(A, A).Return("[]", nil).AnyTimes() + cmCli.EXPECT().RegisterService(A, A, A, A, A).Return(nil).Times(2) + cmCli.EXPECT().AddNode(A, A).Return(proto.NodeID(1), nil).Times(2) + cmCli.EXPECT().ListHostDisk(A, A).Return([]*cmapi.BlobNodeDiskInfo{diskInfo1, diskInfo2}, nil) + cmCli.EXPECT().SetDisk(A, A, A).Return(nil) + + patches := gomonkey.ApplyFunc(readFormatInfo, func(ctx context.Context, path string) (*core.FormatInfo, error) { + if path == path1 || path == path2 { + return nil, syscall.EIO + } + return &core.FormatInfo{CheckSum: 1}, nil + }) + defer patches.Reset() + patches2 := gomonkey.ApplyFunc(disk.NewDiskStorage, func(ctx context.Context, diskConf core.Config) (*disk.DiskStorageWrapper, error) { + if diskConf.Path == path3 || diskConf.Path == path4 { + return nil, syscall.EIO + } + disk2 := &disk.DiskStorageWrapper{DiskStorage: &disk.DiskStorage{DiskID: proto.DiskID(2)}} + return disk2, nil + }) + defer patches2.Reset() + + configInit(&conf) + svr := &Service{ + ClusterMgrClient: cmCli, + Disks: make(map[proto.DiskID]core.DiskAPI), + Conf: &conf, + closeCh: make(chan struct{}), + } + svr.ctx, svr.cancel = context.WithCancel(ctx) + + err = startBlobnodeService(ctx, svr, conf) + require.Nil(t, err) + require.Equal(t, 0, len(svr.Disks)) + + conf.Disks = []core.Config{ + {BaseConfig: core.BaseConfig{Path: path3, AutoFormat: true}, MetaConfig: db.MetaConfig{}}, + } + svr.Conf = &conf + + // newDiskStorage status repaired, skip + diskInfo3 := &cmapi.BlobNodeDiskInfo{ + DiskInfo: cmapi.DiskInfo{ + Path: path3, + Status: proto.DiskStatusRepaired, + }, + DiskHeartBeatInfo: cmapi.DiskHeartBeatInfo{DiskID: proto.DiskID(3)}, + } + cmCli.EXPECT().ListHostDisk(A, A).Return([]*cmapi.BlobNodeDiskInfo{diskInfo3}, nil) + err = startBlobnodeService(ctx, svr, conf) + require.NoError(t, err) + require.Equal(t, 0, len(svr.Disks)) +} + +func TestService_OnlyBlobnode_OpenDiskNormal(t *testing.T) { + ctx := context.Background() + workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"OnlyBlobnode") + require.NoError(t, err) + defer os.RemoveAll(workDir) + + path1 := filepath.Join(workDir, "path1") + err = os.MkdirAll(path1, 0o755) + require.NoError(t, err) + + conf := Config{ + Disks: []core.Config{{BaseConfig: core.BaseConfig{Path: path1, AutoFormat: true, MaxChunks: 700}, MetaConfig: db.MetaConfig{}}}, + InspectConf: DataInspectConf{Record: recordlog.Config{Dir: filepath.Join(workDir, "inspect")}}, + } + + // open disk success + diskInfo1 := &cmapi.BlobNodeDiskInfo{ + DiskInfo: cmapi.DiskInfo{ + Path: path1, + Status: proto.DiskStatusRepaired, + }, + DiskHeartBeatInfo: cmapi.DiskHeartBeatInfo{DiskID: proto.DiskID(1)}, + } + + A := gomock.Any() + ctr := gomock.NewController(t) + cmCli := mocks.NewMockClientAPI(ctr) + cmCli.EXPECT().GetConfig(A, A).Return("[]", nil).AnyTimes() + cmCli.EXPECT().RegisterService(A, A, A, A, A).Return(nil).Times(1) + cmCli.EXPECT().AddNode(A, A).Return(proto.NodeID(1), nil).Times(1) + cmCli.EXPECT().ListHostDisk(A, A).Return([]*cmapi.BlobNodeDiskInfo{diskInfo1}, nil) + cmCli.EXPECT().AllocDiskID(A).Return(proto.DiskID(101), nil) + cmCli.EXPECT().AddDisk(A, A).Return(nil) + + configInit(&conf) + svr := &Service{ + ClusterMgrClient: cmCli, + Disks: make(map[proto.DiskID]core.DiskAPI), + Conf: &conf, + closeCh: make(chan struct{}), + } + svr.ctx, svr.cancel = context.WithCancel(ctx) + + err = startBlobnodeService(ctx, svr, conf) + require.NoError(t, err) + require.Equal(t, 1, len(svr.Disks)) +} + +func TestService_OnlyBlobnode_Fatal(t *testing.T) { + ctx := context.Background() + workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"OnlyBlobnode") + require.NoError(t, err) + defer os.RemoveAll(workDir) + + path1 := filepath.Join(workDir, "path1") + path2 := filepath.Join(workDir, "path2") + for _, path := range []string{path1, path2} { + err = os.MkdirAll(path, 0o755) + require.NoError(t, err) + } + + conf := Config{ + Disks: []core.Config{ + {BaseConfig: core.BaseConfig{Path: path1, AutoFormat: true, MaxChunks: 700}, MetaConfig: db.MetaConfig{}}, + // {BaseConfig: core.BaseConfig{Path: path2, AutoFormat: true}, MetaConfig: db.MetaConfig{}}, + // {BaseConfig: core.BaseConfig{Path: "wrongPath", AutoFormat: true}, MetaConfig: db.MetaConfig{}}, + }, + InspectConf: DataInspectConf{Record: recordlog.Config{Dir: filepath.Join(workDir, "inspect")}}, + } + + // new disk, read meta fake error + diskInfo1 := &cmapi.BlobNodeDiskInfo{ + DiskInfo: cmapi.DiskInfo{ + Path: path1, + Status: proto.DiskStatusRepaired, + }, + DiskHeartBeatInfo: cmapi.DiskHeartBeatInfo{DiskID: proto.DiskID(1)}, + } + // old disk is repairing + diskInfo2 := &cmapi.BlobNodeDiskInfo{ + DiskInfo: cmapi.DiskInfo{ + Path: path2, + Status: proto.DiskStatusRepairing, + }, + DiskHeartBeatInfo: cmapi.DiskHeartBeatInfo{DiskID: proto.DiskID(2)}, + } + + A := gomock.Any() + ctr := gomock.NewController(t) + cmCli := mocks.NewMockClientAPI(ctr) + cmCli.EXPECT().GetConfig(A, A).Return("[]", nil).AnyTimes() + cmCli.EXPECT().RegisterService(A, A, A, A, A).Return(nil).Times(1) + cmCli.EXPECT().AddNode(A, A).Return(proto.NodeID(1), nil).Times(1) + cmCli.EXPECT().ListHostDisk(A, A).Return([]*cmapi.BlobNodeDiskInfo{diskInfo1, diskInfo2}, nil) + // cmCli.EXPECT().AllocDiskID(A).Return(proto.DiskID(102), nil) + + patches := gomonkey.ApplyFunc(readFormatInfo, func(ctx context.Context, path string) (*core.FormatInfo, error) { + if path == path1 { + return nil, errMock + } + return &core.FormatInfo{}, nil + }) + defer patches.Reset() + + mockSpan := opentracing.GlobalTracer().StartSpan("") + patches2 := gomonkey.ApplyMethod(reflect.TypeOf(mockSpan), "Fatalf", func(xx interface{}, format string, v ...interface{}) { + fmt.Println("startBlobnodeService fatal") + }) + defer patches2.Reset() + + configInit(&conf) + svr := &Service{ + ClusterMgrClient: cmCli, + Disks: make(map[proto.DiskID]core.DiskAPI), + Conf: &conf, + closeCh: make(chan struct{}), + } + svr.ctx, svr.cancel = context.WithCancel(ctx) + + // require.Panics(t, func() { startBlobnodeService(ctx, svr, conf) }) + err = startBlobnodeService(ctx, svr, conf) + require.NoError(t, err) + require.Equal(t, 0, len(svr.Disks)) +} + +func TestService_OnlyBlobnode_OpenOldDisk(t *testing.T) { + ctx := context.Background() + workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"OnlyBlobnode") + require.NoError(t, err) + defer os.RemoveAll(workDir) + + path1 := filepath.Join(workDir, "path1") + err = os.MkdirAll(path1, 0o755) + require.NoError(t, err) + + conf := Config{ + Disks: []core.Config{ + {BaseConfig: core.BaseConfig{Path: path1, AutoFormat: true, MaxChunks: 700}, MetaConfig: db.MetaConfig{}}, + }, + InspectConf: DataInspectConf{Record: recordlog.Config{Dir: filepath.Join(workDir, "inspect")}}, + } + + // old disk, repairing, skip + diskInfo1 := &cmapi.BlobNodeDiskInfo{ + DiskInfo: cmapi.DiskInfo{ + Path: path1, + Status: proto.DiskStatusRepairing, + }, + DiskHeartBeatInfo: cmapi.DiskHeartBeatInfo{DiskID: proto.DiskID(1)}, + } + + A := gomock.Any() + ctr := gomock.NewController(t) + cmCli := mocks.NewMockClientAPI(ctr) + cmCli.EXPECT().GetConfig(A, A).Return("[]", nil).AnyTimes() + cmCli.EXPECT().RegisterService(A, A, A, A, A).Return(nil).Times(1) + cmCli.EXPECT().AddNode(A, A).Return(proto.NodeID(1), nil).Times(1) + cmCli.EXPECT().ListHostDisk(A, A).Return([]*cmapi.BlobNodeDiskInfo{diskInfo1}, nil) + + format := &core.FormatInfo{ + FormatInfoProtectedField: core.FormatInfoProtectedField{ + DiskID: proto.DiskID(1), + Version: 1, + Format: core.FormatMetaTypeV1, + }, + } + checkSum, err := format.CalCheckSum() + require.NoError(t, err) + format.CheckSum = checkSum + err = core.SaveDiskFormatInfo(ctx, path1, format) + require.NoError(t, err) + + configInit(&conf) + svr := &Service{ + ClusterMgrClient: cmCli, + Disks: make(map[proto.DiskID]core.DiskAPI), + Conf: &conf, + closeCh: make(chan struct{}), + } + svr.ctx, svr.cancel = context.WithCancel(ctx) + + err = startBlobnodeService(ctx, svr, conf) + require.NoError(t, err) + require.Equal(t, 0, len(svr.Disks)) +} diff --git a/blobstore/testing/mocks/api_clustermgr.go b/blobstore/testing/mocks/api_clustermgr.go index 57b585019..107a39360 100644 --- a/blobstore/testing/mocks/api_clustermgr.go +++ b/blobstore/testing/mocks/api_clustermgr.go @@ -36,6 +36,35 @@ func (m *MockClientAPI) EXPECT() *MockClientAPIMockRecorder { return m.recorder } +// AddDisk mocks base method. +func (m *MockClientAPI) AddDisk(arg0 context.Context, arg1 *clustermgr.BlobNodeDiskInfo) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "AddDisk", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// AddDisk indicates an expected call of AddDisk. +func (mr *MockClientAPIMockRecorder) AddDisk(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AddDisk", reflect.TypeOf((*MockClientAPI)(nil).AddDisk), arg0, arg1) +} + +// AddNode mocks base method. +func (m *MockClientAPI) AddNode(arg0 context.Context, arg1 *clustermgr.BlobNodeInfo) (proto.NodeID, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "AddNode", arg0, arg1) + ret0, _ := ret[0].(proto.NodeID) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// AddNode indicates an expected call of AddNode. +func (mr *MockClientAPIMockRecorder) AddNode(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AddNode", reflect.TypeOf((*MockClientAPI)(nil).AddNode), arg0, arg1) +} + // AllocBid mocks base method. func (m *MockClientAPI) AllocBid(arg0 context.Context, arg1 *clustermgr.BidScopeArgs) (*clustermgr.BidScopeRet, error) { m.ctrl.T.Helper() @@ -51,6 +80,21 @@ func (mr *MockClientAPIMockRecorder) AllocBid(arg0, arg1 interface{}) *gomock.Ca return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AllocBid", reflect.TypeOf((*MockClientAPI)(nil).AllocBid), arg0, arg1) } +// AllocDiskID mocks base method. +func (m *MockClientAPI) AllocDiskID(arg0 context.Context) (proto.DiskID, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "AllocDiskID", arg0) + ret0, _ := ret[0].(proto.DiskID) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// AllocDiskID indicates an expected call of AllocDiskID. +func (mr *MockClientAPIMockRecorder) AllocDiskID(arg0 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AllocDiskID", reflect.TypeOf((*MockClientAPI)(nil).AllocDiskID), arg0) +} + // AllocVolume mocks base method. func (m *MockClientAPI) AllocVolume(arg0 context.Context, arg1 *clustermgr.AllocVolumeArgs) (clustermgr.AllocatedVolumeInfos, error) { m.ctrl.T.Helper() @@ -170,6 +214,21 @@ func (mr *MockClientAPIMockRecorder) GetVolumeInfo(arg0, arg1 interface{}) *gomo return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetVolumeInfo", reflect.TypeOf((*MockClientAPI)(nil).GetVolumeInfo), arg0, arg1) } +// HeartbeatDisk mocks base method. +func (m *MockClientAPI) HeartbeatDisk(arg0 context.Context, arg1 []*clustermgr.DiskHeartBeatInfo) ([]*clustermgr.DiskHeartbeatRet, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "HeartbeatDisk", arg0, arg1) + ret0, _ := ret[0].([]*clustermgr.DiskHeartbeatRet) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// HeartbeatDisk indicates an expected call of HeartbeatDisk. +func (mr *MockClientAPIMockRecorder) HeartbeatDisk(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "HeartbeatDisk", reflect.TypeOf((*MockClientAPI)(nil).HeartbeatDisk), arg0, arg1) +} + // ListDisk mocks base method. func (m *MockClientAPI) ListDisk(arg0 context.Context, arg1 *clustermgr.ListOptionArgs) (clustermgr.ListDiskRet, error) { m.ctrl.T.Helper() @@ -185,6 +244,21 @@ func (mr *MockClientAPIMockRecorder) ListDisk(arg0, arg1 interface{}) *gomock.Ca return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListDisk", reflect.TypeOf((*MockClientAPI)(nil).ListDisk), arg0, arg1) } +// ListHostDisk mocks base method. +func (m *MockClientAPI) ListHostDisk(arg0 context.Context, arg1 string) ([]*clustermgr.BlobNodeDiskInfo, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListHostDisk", arg0, arg1) + ret0, _ := ret[0].([]*clustermgr.BlobNodeDiskInfo) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListHostDisk indicates an expected call of ListHostDisk. +func (mr *MockClientAPIMockRecorder) ListHostDisk(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListHostDisk", reflect.TypeOf((*MockClientAPI)(nil).ListHostDisk), arg0, arg1) +} + // ListShardNodeDisk mocks base method. func (m *MockClientAPI) ListShardNodeDisk(arg0 context.Context, arg1 *clustermgr.ListOptionArgs) (clustermgr.ListShardNodeDiskRet, error) { m.ctrl.T.Helper() @@ -200,6 +274,21 @@ func (mr *MockClientAPIMockRecorder) ListShardNodeDisk(arg0, arg1 interface{}) * return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListShardNodeDisk", reflect.TypeOf((*MockClientAPI)(nil).ListShardNodeDisk), arg0, arg1) } +// ListVolumeUnit mocks base method. +func (m *MockClientAPI) ListVolumeUnit(arg0 context.Context, arg1 *clustermgr.ListVolumeUnitArgs) ([]*clustermgr.VolumeUnitInfo, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ListVolumeUnit", arg0, arg1) + ret0, _ := ret[0].([]*clustermgr.VolumeUnitInfo) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ListVolumeUnit indicates an expected call of ListVolumeUnit. +func (mr *MockClientAPIMockRecorder) ListVolumeUnit(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ListVolumeUnit", reflect.TypeOf((*MockClientAPI)(nil).ListVolumeUnit), arg0, arg1) +} + // RegisterService mocks base method. func (m *MockClientAPI) RegisterService(arg0 context.Context, arg1 clustermgr.ServiceNode, arg2, arg3, arg4 uint32) error { m.ctrl.T.Helper() @@ -214,6 +303,20 @@ func (mr *MockClientAPIMockRecorder) RegisterService(arg0, arg1, arg2, arg3, arg return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RegisterService", reflect.TypeOf((*MockClientAPI)(nil).RegisterService), arg0, arg1, arg2, arg3, arg4) } +// ReportChunk mocks base method. +func (m *MockClientAPI) ReportChunk(arg0 context.Context, arg1 *clustermgr.ReportChunkArgs) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ReportChunk", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// ReportChunk indicates an expected call of ReportChunk. +func (mr *MockClientAPIMockRecorder) ReportChunk(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ReportChunk", reflect.TypeOf((*MockClientAPI)(nil).ReportChunk), arg0, arg1) +} + // RetainVolume mocks base method. func (m *MockClientAPI) RetainVolume(arg0 context.Context, arg1 *clustermgr.RetainVolumeArgs) (clustermgr.RetainVolumes, error) { m.ctrl.T.Helper() @@ -229,6 +332,48 @@ func (mr *MockClientAPIMockRecorder) RetainVolume(arg0, arg1 interface{}) *gomoc return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RetainVolume", reflect.TypeOf((*MockClientAPI)(nil).RetainVolume), arg0, arg1) } +// SetCompactChunk mocks base method. +func (m *MockClientAPI) SetCompactChunk(arg0 context.Context, arg1 *clustermgr.SetCompactChunkArgs) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "SetCompactChunk", arg0, arg1) + ret0, _ := ret[0].(error) + return ret0 +} + +// SetCompactChunk indicates an expected call of SetCompactChunk. +func (mr *MockClientAPIMockRecorder) SetCompactChunk(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SetCompactChunk", reflect.TypeOf((*MockClientAPI)(nil).SetCompactChunk), arg0, arg1) +} + +// SetConfig mocks base method. +func (m *MockClientAPI) SetConfig(arg0 context.Context, arg1, arg2 string) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "SetConfig", arg0, arg1, arg2) + ret0, _ := ret[0].(error) + return ret0 +} + +// SetConfig indicates an expected call of SetConfig. +func (mr *MockClientAPIMockRecorder) SetConfig(arg0, arg1, arg2 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SetConfig", reflect.TypeOf((*MockClientAPI)(nil).SetConfig), arg0, arg1, arg2) +} + +// SetDisk mocks base method. +func (m *MockClientAPI) SetDisk(arg0 context.Context, arg1 proto.DiskID, arg2 proto.DiskStatus) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "SetDisk", arg0, arg1, arg2) + ret0, _ := ret[0].(error) + return ret0 +} + +// SetDisk indicates an expected call of SetDisk. +func (mr *MockClientAPIMockRecorder) SetDisk(arg0, arg1, arg2 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SetDisk", reflect.TypeOf((*MockClientAPI)(nil).SetDisk), arg0, arg1, arg2) +} + // ShardNodeDiskInfo mocks base method. func (m *MockClientAPI) ShardNodeDiskInfo(arg0 context.Context, arg1 proto.DiskID) (*clustermgr.ShardNodeDiskInfo, error) { m.ctrl.T.Helper()