From 0e239bbda4c5d35c53ea16376f938c29b72ebcdb Mon Sep 17 00:00:00 2001 From: mawei029 Date: Wed, 17 Sep 2025 16:58:36 +0800 Subject: [PATCH] fix(bssdk): update shard node leader, need change next shard node host with: #1000347594 Signed-off-by: mawei029 --- blobstore/access/stream/stream_blob.go | 16 +++-- blobstore/access/stream/stream_blob_test.go | 79 ++++++++++++++++++--- 2 files changed, 78 insertions(+), 17 deletions(-) diff --git a/blobstore/access/stream/stream_blob.go b/blobstore/access/stream/stream_blob.go index 0f2b52e21..421556d1e 100644 --- a/blobstore/access/stream/stream_blob.go +++ b/blobstore/access/stream/stream_blob.go @@ -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 diff --git a/blobstore/access/stream/stream_blob_test.go b/blobstore/access/stream/stream_blob_test.go index 4a8eeea6c..df85f90e5 100644 --- a/blobstore/access/stream/stream_blob_test.go +++ b/blobstore/access/stream/stream_blob_test.go @@ -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 mock:first NotLeader,and 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)) +}