fix(blobnode): blobnode start up change ip, add node with local node id

with: #1000464074

Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
mawei029 2026-01-08 18:33:31 +08:00 committed by 梁曟風
parent 8d37f14b22
commit ac7ac4b7bb
2 changed files with 120 additions and 4 deletions

View File

@ -322,6 +322,13 @@ func (s *Service) registerNode(ctx context.Context, conf *Config) error {
return err
}
// diskNodeID: 1. is zero: a.new node; b. old version upgrade; c. all borken disks;
// 2. not zero: v1.5.2 node restart(change IP or not change)
localNodeID, err := getNodeID(ctx, conf)
if err != nil {
return err
}
nodeToCm := cmapi.BlobNodeInfo{
NodeInfo: cmapi.NodeInfo{
ClusterID: conf.ClusterID,
@ -330,6 +337,7 @@ func (s *Service) registerNode(ctx context.Context, conf *Config) error {
Rack: conf.Rack,
Host: conf.Host,
Role: proto.NodeRoleBlobNode,
NodeID: localNodeID,
},
}
@ -338,6 +346,11 @@ func (s *Service) registerNode(ctx context.Context, conf *Config) error {
return err
}
// assert nodeID == localNodeID
if localNodeID != 0 && nodeID != localNodeID {
return errors.Newf("node id not match, cm:%d, local:%d]", nodeID, localNodeID)
}
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
@ -396,6 +409,37 @@ func readFormatInfo(ctx context.Context, diskRootPath string, nodeID proto.NodeI
return formatInfo, err
}
func getNodeID(ctx context.Context, conf *Config) (proto.NodeID, error) {
span := trace.SpanFromContextSafe(ctx)
nodeID := proto.NodeID(0)
// If a broken disk is detected during boot:
// 1. If the IP address hasn't been changed, we can query historical bad disk records in cm;
// 2. If the IP address has been changed, we can not query it in cm.
for _, diskConf := range conf.Disks {
// bad disk
_, err := os.ReadDir(diskConf.Path)
if err != nil {
span.Warnf("read disk[%s] root path error:%+v", diskConf.Path, err)
continue
}
formatInfo, err := core.ReadFormatInfo(ctx, diskConf.Path)
if err != nil {
// new disk; format.info v1 upgrade to v2; bad disk; or other error
span.Warnf("Failed read disk[%s] format info, err:%+v", diskConf.Path, err)
continue
}
if formatInfo.NodeID != 0 && nodeID != 0 && nodeID != formatInfo.NodeID {
return 0, errors.Newf("disk[%s] node id not match", diskConf.Path)
}
nodeID = formatInfo.NodeID
}
return nodeID, nil
}
func isAllInConfig(ctx context.Context, registeredDisks []*cmapi.BlobNodeDiskInfo, conf *Config) bool {
span := trace.SpanFromContextSafe(ctx)
configDiskMap := make(map[string]struct{})
@ -441,6 +485,10 @@ func startBlobnodeService(ctx context.Context, svr *Service, conf Config) (err e
span := trace.SpanFromContextSafe(ctx)
span.Debug("start blobnode service...")
if err = svr.registerNode(ctx, &conf); err != nil {
span.Fatalf("fail to register node to clusterMgr, err:%+v", err)
}
node := cmapi.ServiceNode{
ClusterID: uint64(conf.ClusterID),
Name: proto.ServiceNameBlobNode,
@ -452,10 +500,6 @@ func startBlobnodeService(ctx context.Context, svr *Service, conf Config) (err e
span.Fatalf("blobnode register to clusterMgr error:%+v", err)
}
if err = svr.registerNode(ctx, &conf); err != nil {
span.Fatalf("fail to register node to clusterMgr, err:%+v", err)
}
registeredDisks, err := svr.ClusterMgrClient.ListHostDisk(ctx, conf.Host)
if err != nil {
span.Errorf("Failed ListDisk from clusterMgr. err:%+v", err)

View File

@ -1576,3 +1576,75 @@ func TestService_DataInspect(t *testing.T) {
log.Infof("inspect stat: %+v\n", data)
}
}
func TestService_Blobnode_registerNode(t *testing.T) {
ctx := context.Background()
workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"AddNode")
require.NoError(t, err)
defer os.RemoveAll(workDir)
path1 := filepath.Join(workDir, "path1")
err = os.MkdirAll(path1, 0o755)
require.NoError(t, err)
path2 := filepath.Join(workDir, "path2")
err = os.MkdirAll(path2, 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, MaxChunks: 700}, MetaConfig: db.MetaConfig{}},
},
}
A := gomock.Any()
ctr := gomock.NewController(t)
cmCli := mocks.NewMockClientAPI(ctr)
cmCli.EXPECT().GetConfig(A, A).Return("[]", nil).AnyTimes()
cmCli.EXPECT().AddNode(A, A).Return(proto.NodeID(1), nil).Times(1)
format := &core.FormatInfo{}
format.FormatInfoProtectedField = core.FormatInfoProtectedField{
DiskID: proto.DiskID(1),
Version: 1,
Format: core.FormatMetaTypeV1,
Ctime: time.Now().UnixNano(),
}
format.NodeID = 1
format.NodeCtime = format.Ctime
err = format.CalCheckSum()
require.NoError(t, err)
err = core.SaveDiskFormatInfo(ctx, path1, format)
require.NoError(t, err)
format2 := *format
err = format2.CalCheckSum()
require.NoError(t, err)
err = core.SaveDiskFormatInfo(ctx, path2, &format2)
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 = svr.registerNode(ctx, &conf)
require.NoError(t, err)
require.Equal(t, proto.NodeID(1), conf.NodeID)
// error, node id not match
conf.NodeID = 0
format2.NodeID = 999
err = format2.CalCheckSum()
require.NoError(t, err)
err = core.SaveDiskFormatInfo(ctx, path2, &format2)
require.NoError(t, err)
err = svr.registerNode(ctx, &conf)
require.NotNil(t, err)
}