refactor(blobnode): added multiple startup of disk ut individual tests

with: #1000212202

Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
mawei029 2025-07-03 16:41:52 +08:00 committed by slasher
parent f12c6f5acf
commit c113ecfa69
5 changed files with 581 additions and 120 deletions

View File

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

View File

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

View File

@ -44,7 +44,7 @@ type Service struct {
WorkerService *WorkerService
// client handler
ClusterMgrClient *cmapi.Client
ClusterMgrClient cmapi.APIBlobnode
Conf *Config
inspectMgr *DataInspectMgr

View File

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

View File

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