fix(bssdk): update shard node leader, need change next shard node host

with: #1000347594

Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
mawei029 2025-09-17 16:58:36 +08:00 committed by slasher
parent 43780a871a
commit 0e239bbda4
2 changed files with 78 additions and 17 deletions

View File

@ -481,7 +481,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.waitShardnodeNextLeader(ctx, args.clusterID, args.Suid, args.DiskID)
err1 := h.updateLeaderFromNewHost(ctx, args.clusterID, args.Suid, args.DiskID)
if err1 != nil {
span.Warnf("fail to change other shard node, cluster:%d, err:%+v", args.clusterID, err1)
}
@ -506,7 +506,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.updateShard(ctx, args); err1 != nil {
if err1 := h.updateLeaderFromNewHost(ctx, args.clusterID, args.Suid, args.DiskID); err1 != nil {
span.Warnf("fail to update shard, cluster:%d, err:%+v", args.clusterID, err1)
}
return false, args.err
@ -515,12 +515,12 @@ func (h *Handler) punishAndUpdate(ctx context.Context, args *punishArgs) (bool,
}
// err:dial tcp 127.0.0.1:9100: connect: connection refused code:500
if errorConnectionRefused(args.err) {
span.Warnf("shardnode connection refused, args:%+v, err:%+v", *args, args.err)
if errorConnectionRefused(args.err) || errorTimeout(args.err) {
span.Warnf("shardnode connection refused/timeout, args:%+v, err:%+v", *args, args.err)
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.waitShardnodeNextLeader(ctx, args.clusterID, args.Suid, args.DiskID)
err1 := h.updateLeaderFromNewHost(ctx, args.clusterID, args.Suid, args.DiskID)
if err1 != nil {
span.Warnf("fail to change other shard node, cluster:%d, err:%+v", args.clusterID, err1)
}
@ -542,7 +542,8 @@ func (h *Handler) updateShardRoute(ctx context.Context, clusterID proto.ClusterI
return shardMgr.UpdateRoute(ctx)
}
func (h *Handler) updateShard(ctx context.Context, args *punishArgs) error {
// updateLeaderFromCurrentHost from old current shard host/disk, get leader and update shard
func (h *Handler) updateLeaderFromCurrentHost(ctx context.Context, args *punishArgs) error {
shardMgr, err := h.clusterController.GetShardController(args.clusterID)
if err != nil {
return err
@ -556,7 +557,8 @@ func (h *Handler) updateShard(ctx context.Context, args *punishArgs) error {
return shardMgr.UpdateShard(ctx, shardStat)
}
func (h *Handler) waitShardnodeNextLeader(ctx context.Context, clusterID proto.ClusterID, suid proto.Suid, badDisk proto.DiskID) error {
// 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)
if err != nil {
return err

View File

@ -454,16 +454,6 @@ func TestStreamBlobOther(t *testing.T) {
require.Equal(t, false, interrupt)
require.ErrorIs(t, err1, io.EOF)
shardnodeClient := mocks.NewMockShardnodeAccess(ctr)
shardnodeClient.EXPECT().GetShardStats(gAny, gAny, gAny).Return(shardnode.ShardStats{LeaderDiskID: 11}, nil).Times(1)
h.shardnodeClient = shardnodeClient
h.ShardnodeRetryTimes = defaultShardnodeRetryTimes
interrupt, err1 = h.punishAndUpdate(ctx, &punishArgs{
err: errcode.ErrShardNodeNotLeader,
})
require.Equal(t, false, interrupt)
require.ErrorIs(t, err1, errcode.ErrShardNodeNotLeader)
// broken disk
interrupt, err1 = h.punishAndUpdate(ctx, &punishArgs{
ShardOpHeader: shardnode.ShardOpHeader{},
@ -474,7 +464,17 @@ func TestStreamBlobOther(t *testing.T) {
require.Equal(t, false, interrupt)
require.ErrorIs(t, err1, errcode.ErrDiskBroken)
shardnodeClient := mocks.NewMockShardnodeAccess(ctr)
shardnodeClient.EXPECT().GetShardStats(gAny, gAny, gAny).Return(shardnode.ShardStats{LeaderDiskID: 11}, nil).Times(1)
h.shardnodeClient = shardnodeClient
h.ShardnodeRetryTimes = defaultShardnodeRetryTimes
err1 = h.updateLeaderFromCurrentHost(ctx, &punishArgs{
err: errcode.ErrShardNodeNotLeader,
})
require.NoError(t, err1)
// wait connect refused
h.ShardnodeRetryTimes = defaultShardnodeRetryTimes
info := controller.ShardOpInfo{
DiskID: 101,
Suid: proto.EncodeSuid(1, 0, 1),
@ -508,3 +508,62 @@ func TestStreamBlobOther(t *testing.T) {
// require.Equal(t, true, interrupt)
// require.ErrorIs(t, err1, errcode.ErrCallShardNodeFail)
}
func TestStreamBlob_NotLeader_RetrySuccess(t *testing.T) {
ctx := context.Background()
ctr := gomock.NewController(t)
gAny := gomock.Any()
// old leader, Not Leader
oldInfo := controller.ShardOpInfo{
DiskID: 101,
Suid: proto.EncodeSuid(1, 0, 1),
RouteVersion: 1,
}
// new leader
newInfo := controller.ShardOpInfo{
DiskID: 102,
Suid: proto.EncodeSuid(1, 0, 2),
RouteVersion: 1,
}
// 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().GetShardID().Return(proto.ShardID(1)).AnyTimes()
// shard controller mock
shardMgr := NewMockShardController(ctr)
shardMgr.EXPECT().GetShard(gAny, gAny).Return(shard, nil).AnyTimes()
shardMgr.EXPECT().GetShardByID(gAny, proto.ShardID(1)).Return(shard, nil).AnyTimes()
shardMgr.EXPECT().GetSpaceID().Return(proto.SpaceID(1)).AnyTimes()
shardMgr.EXPECT().UpdateRoute(gAny).Return(nil).AnyTimes()
shardMgr.EXPECT().GetShardSubRangeCount(gAny).Return(2).AnyTimes()
shardMgr.EXPECT().UpdateShard(gAny, gAny).Return(nil).Times(1) // waitShardnodeNextLeader finally
// service controller mock
svrCtrl := NewMockServiceController(ctr)
svrCtrl.EXPECT().GetShardnodeHost(gAny, proto.DiskID(oldInfo.DiskID)).Return(&controller.HostIDC{Host: "host-old"}, nil).AnyTimes()
svrCtrl.EXPECT().GetShardnodeHost(gAny, proto.DiskID(newInfo.DiskID)).Return(&controller.HostIDC{Host: "host-new"}, nil).AnyTimes()
// cluster controller mock
clu := NewMockClusterController(ctr)
clu.EXPECT().GetShardController(gAny).Return(shardMgr, nil).AnyTimes()
clu.EXPECT().GetServiceController(gAny).Return(svrCtrl, nil).AnyTimes()
// shardnode client mockfirst NotLeaderand then success
shardCli := mocks.NewMockShardnodeAccess(ctr)
shardCli.EXPECT().DeleteBlob(gAny, gAny, gAny).Return(errcode.ErrShardNodeNotLeader)
shardCli.EXPECT().GetShardStats(gAny, gAny, gAny).Return(shardnode.ShardStats{LeaderDiskID: newInfo.DiskID}, nil)
shardCli.EXPECT().DeleteBlob(gAny, gAny, gAny).Return(nil)
h := &Handler{
clusterController: clu,
shardnodeClient: shardCli,
}
h.ShardnodeRetryTimes = defaultShardnodeRetryTimes
args := acapi.DelBlobArgs{ClusterID: 1, BlobName: "blob-notleader"}
require.NoError(t, h.DeleteBlob(ctx, &args))
}