fix(blobnode): update leader from new host, may need exclucde multiple disk

with: #1000348559

Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
mawei029 2025-09-22 19:32:48 +08:00 committed by slasher
parent 7d153130f0
commit e4d5e0f8c3
5 changed files with 36 additions and 28 deletions

View File

@ -550,7 +550,7 @@ type ShardOpInfo struct {
type Shard interface {
GetShardID() proto.ShardID
GetRange() sharding.Range
GetMember(context.Context, acapi.GetShardMode, proto.DiskID) (ShardOpInfo, error)
GetMember(context.Context, acapi.GetShardMode, map[proto.DiskID]struct{}) (ShardOpInfo, error)
}
// shard implement btree.Item interface, shard route information
@ -587,12 +587,12 @@ func (i *shard) GetRange() sharding.Range {
return i.rangeExt
}
func (i *shard) GetMember(ctx context.Context, mode acapi.GetShardMode, exclude proto.DiskID) (ShardOpInfo, error) {
func (i *shard) GetMember(ctx context.Context, mode acapi.GetShardMode, exclude map[proto.DiskID]struct{}) (ShardOpInfo, error) {
span := trace.SpanFromContextSafe(ctx)
span.Debugf("get shard member, mode:%d, exclude:%d, shard:%+v", mode, exclude, *i)
// 1. get member exclude disk id
if exclude != 0 {
if len(exclude) != 0 {
return i.getMemberExcluded(ctx, exclude)
}
@ -600,11 +600,11 @@ func (i *shard) GetMember(ctx context.Context, mode acapi.GetShardMode, exclude
if mode == acapi.GetShardModeLeader {
return i.getMemberLeader(ctx)
}
return i.getMemberRandom(ctx, 0)
return i.getMemberRandom(ctx, nil)
}
func (i *shard) getMemberExcluded(ctx context.Context, diskID proto.DiskID) (ShardOpInfo, error) {
return i.getMemberRandom(ctx, diskID)
func (i *shard) getMemberExcluded(ctx context.Context, exclude map[proto.DiskID]struct{}) (ShardOpInfo, error) {
return i.getMemberRandom(ctx, exclude)
}
func (i *shard) getMemberLeader(ctx context.Context) (ShardOpInfo, error) {
@ -615,7 +615,7 @@ func (i *shard) getMemberLeader(ctx context.Context) (ShardOpInfo, error) {
}, nil
}
func (i *shard) getMemberRandom(ctx context.Context, exclude proto.DiskID) (ShardOpInfo, error) {
func (i *shard) getMemberRandom(ctx context.Context, exclude map[proto.DiskID]struct{}) (ShardOpInfo, error) {
span := trace.SpanFromContextSafe(ctx)
n := len(i.units)
@ -627,7 +627,7 @@ func (i *shard) getMemberRandom(ctx context.Context, exclude proto.DiskID) (Shar
if err != nil {
return ShardOpInfo{}, err
}
if i.units[idx].DiskID != exclude && !disk.Punished && !i.units[idx].Learner {
if _, exist := exclude[i.units[idx].DiskID]; !exist && !disk.Punished && !i.units[idx].Learner {
return i.getShardOpInfo(idx), nil
}

View File

@ -181,7 +181,7 @@ func TestShardController(t *testing.T) {
require.Equal(t, clustermgr.ShardUnit{Suid: proto.EncodeSuid(1, 1, 0), DiskID: 2}, si.units[1])
// require.Equal(t, clustermgr.ShardUnit{Suid: proto.EncodeSuid(1, 2, 1), DiskID: 4}, si.units[2])
opInfo, err := si.GetMember(ctx, acapi.GetShardModeLeader, 0)
opInfo, err := si.GetMember(ctx, acapi.GetShardModeLeader, nil)
require.NoError(t, err)
require.Equal(t, ShardOpInfo{
DiskID: 2,
@ -636,7 +636,7 @@ func TestShardUpdate(t *testing.T) {
sd, exist := svr.getShardNoLock(shardID)
require.True(t, exist)
opHeader, err := sd.GetMember(ctx, acapi.GetShardModeLeader, 0)
opHeader, err := sd.GetMember(ctx, acapi.GetShardModeLeader, nil)
require.NoError(t, err)
expect := ShardOpInfo{
DiskID: 1,
@ -799,11 +799,11 @@ func TestShardGetShard(t *testing.T) {
sd, err = svr.GetShardByID(ctx, proto.ShardID(1))
require.NoError(t, err)
shardInfo, err := sd.GetMember(ctx, acapi.GetShardModeLeader, 0)
shardInfo, err := sd.GetMember(ctx, acapi.GetShardModeLeader, nil)
require.NoError(t, err)
require.Equal(t, shards[0].shardID, shardInfo.Suid.ShardID())
shardInfo, err = sd.GetMember(ctx, acapi.GetShardModeRandom, 0)
shardInfo, err = sd.GetMember(ctx, acapi.GetShardModeRandom, nil)
require.NoError(t, err)
require.Equal(t, shards[0].shardID, shardInfo.Suid.ShardID())
@ -817,10 +817,11 @@ func TestShardGetShard(t *testing.T) {
require.NoError(t, err)
require.Equal(t, proto.ShardID(1), sd.GetShardID())
newDisk, err := sd.GetMember(ctx, acapi.GetShardModeRandom, 2)
newDisk, err := sd.GetMember(ctx, acapi.GetShardModeRandom, map[proto.DiskID]struct{}{2: {}, 1: {}})
require.NoError(t, err)
require.NotEqual(t, proto.DiskID(2), newDisk.DiskID)
require.Contains(t, []proto.DiskID{1, 3}, newDisk.DiskID)
require.NotEqual(t, proto.DiskID(1), newDisk.DiskID)
require.Contains(t, []proto.DiskID{3}, newDisk.DiskID)
}
// get shard by range

View File

@ -553,7 +553,7 @@ func (m *MockShard) EXPECT() *MockShardMockRecorder {
}
// GetMember mocks base method.
func (m *MockShard) GetMember(arg0 context.Context, arg1 access.GetShardMode, arg2 proto.DiskID) (controller.ShardOpInfo, error) {
func (m *MockShard) GetMember(arg0 context.Context, arg1 access.GetShardMode, arg2 map[proto.DiskID]struct{}) (controller.ShardOpInfo, error) {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "GetMember", arg0, arg1, arg2)
ret0, _ := ret[0].(controller.ShardOpInfo)

View File

@ -423,7 +423,7 @@ func (h *Handler) getOpHeaderByShard(ctx context.Context, shardMgr controller.IS
span := trace.SpanFromContextSafe(ctx)
spaceID := shardMgr.GetSpaceID()
info, err := shard.GetMember(ctx, mode, 0)
info, err := shard.GetMember(ctx, mode, nil)
if err != nil {
return shardnode.ShardOpHeader{}, err
}
@ -460,6 +460,7 @@ type punishArgs struct {
clusterID proto.ClusterID
host string
mode acapi.GetShardMode
exclude map[proto.DiskID]struct{}
err error
}
@ -481,7 +482,7 @@ func (h *Handler) punishAndUpdate(ctx context.Context, args *punishArgs) (bool,
// if leader node broken disk, it cant get shard stats, wait new leader
h.punishShardnodeDisk(ctx, args.clusterID, args.DiskID, args.host, "Broken")
if args.mode == acapi.GetShardModeLeader {
err1 := h.updateLeaderFromNewHost(ctx, args.clusterID, args.Suid, args.DiskID)
err1 := h.updateLeaderFromNewHost(ctx, args)
if err1 != nil {
span.Warnf("fail to change other shard node, cluster:%d, err:%+v", args.clusterID, err1)
}
@ -506,7 +507,7 @@ func (h *Handler) punishAndUpdate(ctx context.Context, args *punishArgs) (bool,
// select master
case errcode.CodeShardNodeNotLeader: // leader disk id error when create/delete/seal
if err1 := h.updateLeaderFromNewHost(ctx, args.clusterID, args.Suid, args.DiskID); err1 != nil {
if err1 := h.updateLeaderFromNewHost(ctx, args); err1 != nil {
span.Warnf("fail to update leader and shard info, cluster:%d, err:%+v", args.clusterID, err1)
}
return false, args.err
@ -520,7 +521,7 @@ func (h *Handler) punishAndUpdate(ctx context.Context, args *punishArgs) (bool,
h.groupRun.Do("shardnode-leader-"+args.DiskID.ToString(), func() (interface{}, error) {
// must wait have master leader, block wait
h.punishShardnodeDisk(ctx, args.clusterID, args.DiskID, args.host, "Refused")
err1 := h.updateLeaderFromNewHost(ctx, args.clusterID, args.Suid, args.DiskID)
err1 := h.updateLeaderFromNewHost(ctx, args)
if err1 != nil {
span.Warnf("fail to change other shard node, cluster:%d, err:%+v", args.clusterID, err1)
}
@ -558,30 +559,36 @@ func (h *Handler) updateLeaderFromCurrentHost(ctx context.Context, args *punishA
}
// updateLeaderFromNewHost from other shard host/disk, get leader and update shard
func (h *Handler) updateLeaderFromNewHost(ctx context.Context, clusterID proto.ClusterID, suid proto.Suid, badDisk proto.DiskID) error {
shardMgr, err := h.clusterController.GetShardController(clusterID)
func (h *Handler) updateLeaderFromNewHost(ctx context.Context, args *punishArgs) error {
shardMgr, err := h.clusterController.GetShardController(args.clusterID)
if err != nil {
return err
}
shard, err := shardMgr.GetShardByID(ctx, suid.ShardID())
shard, err := shardMgr.GetShardByID(ctx, args.Suid.ShardID())
if err != nil {
return err
}
if len(args.exclude) == 0 {
args.exclude = make(map[proto.DiskID]struct{})
args.exclude[args.DiskID] = struct{}{}
}
// we get new disk, exclude bad diskID
newDisk, err := shard.GetMember(ctx, acapi.GetShardModeRandom, badDisk)
newDisk, err := shard.GetMember(ctx, acapi.GetShardModeRandom, args.exclude)
if err != nil {
return err
}
newHost, err := h.getShardHost(ctx, clusterID, newDisk.DiskID)
newHost, err := h.getShardHost(ctx, args.clusterID, newDisk.DiskID)
if err != nil {
return err
}
// span := trace.SpanFromContextSafe(ctx)
// span.Debugf("get newDisk:%+v, old host:%s, old disk:%d", newDisk, args.host, args.DiskID)
shardStat, err := h.getLeaderShardInfo(ctx, clusterID, newHost, newDisk.DiskID, newDisk.Suid, badDisk)
shardStat, err := h.getLeaderShardInfo(ctx, args.clusterID, newHost, newDisk.DiskID, newDisk.Suid, args.DiskID)
if err != nil {
args.exclude[newDisk.DiskID] = struct{}{}
return err
}

View File

@ -482,7 +482,7 @@ func TestStreamBlobOther(t *testing.T) {
clu.EXPECT().GetShardController(gAny).Return(shardMgr, nil)
clu.EXPECT().GetServiceController(gAny).Return(svrCtrl, nil).Times(2)
shardInfo := NewMockShard(ctr)
shardInfo.EXPECT().GetMember(gAny, gAny, proto.DiskID(1)).Return(info, nil)
shardInfo.EXPECT().GetMember(gAny, gAny, map[proto.DiskID]struct{}{1: {}}).Return(info, nil)
shardMgr.EXPECT().GetShardByID(gAny, gAny).Return(shardInfo, nil)
shardMgr.EXPECT().UpdateShard(gAny, gAny).Return(nil)
svrCtrl.EXPECT().GetShardnodeHost(gAny, proto.DiskID(101)).Return(&controller.HostIDC{Host: "host101"}, nil)
@ -528,8 +528,8 @@ func TestStreamBlob_NotLeader_RetrySuccess(t *testing.T) {
// shard mock
shard := NewMockShard(ctr)
shard.EXPECT().GetMember(gAny, gAny, gAny).Return(oldInfo, nil).AnyTimes() // getShardOpHeader
shard.EXPECT().GetMember(gAny, acapi.GetShardModeRandom, proto.DiskID(oldInfo.DiskID)).Return(newInfo, nil).AnyTimes() // waitShardnodeNextLeader
shard.EXPECT().GetMember(gAny, gAny, gAny).Return(oldInfo, nil).AnyTimes() // getShardOpHeader
shard.EXPECT().GetMember(gAny, acapi.GetShardModeRandom, map[proto.DiskID]struct{}{oldInfo.DiskID: {}}).Return(newInfo, nil).AnyTimes() // waitShardnodeNextLeader
shard.EXPECT().GetShardID().Return(proto.ShardID(1)).AnyTimes()
// shard controller mock