fix(bssdk): old leader disk is broken, change leader disk

with: #1000014262, #1000010625

Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
mawei029 2025-03-13 15:16:57 +08:00 committed by slasher
parent ff9a307699
commit c577b51612
5 changed files with 106 additions and 82 deletions

View File

@ -244,7 +244,7 @@ func (s *shardControllerImpl) UpdateRoute(ctx context.Context) error {
// UpdateShard update leader disk id and units info
func (s *shardControllerImpl) UpdateShard(ctx context.Context, sd shardnode.ShardStats) error {
span := trace.SpanFromContextSafe(ctx)
span.Debugf("will update shard=%+v", sd)
span.Debugf("will update shard, leaderDiskID=%d, LeaderSuid=%d, version=%d, suid=%d", sd.LeaderDiskID, sd.LeaderSuid, sd.RouteVersion, sd.Suid)
_, err, _ := s.groupRun.Do("shardID-"+sd.Suid.ShardID().ToString(), func() (interface{}, error) {
if isInvalidShardStat(sd) {
@ -257,16 +257,20 @@ func (s *shardControllerImpl) UpdateShard(ctx context.Context, sd shardnode.Shar
// skip old route version
oldShard, exist := s.getShardNoLock(sd.Suid.ShardID())
if !exist {
span.Warnf("dont need update shard, exist:%t, current shard:%v, replace shard:%v", exist, oldShard, sd)
span.Warnf("dont need update shard, exist:%t, current shard:%+v, replace shard:%+v", exist, oldShard, sd)
return nil, errcode.ErrAccessNotFoundShard
}
// only update leader disk id/suid
// only update leader diskID/suid ; may be sd.LeaderDiskID is not in units
// don't need to judge or change RouteVersion, when switch the primary shardNode. only update version in cm GetCatalogChanges
oldShard.leaderDiskID = sd.LeaderDiskID
oldShard.leaderSuid = sd.LeaderSuid
return nil, nil
if sd.LeaderSuid.Epoch() > oldShard.units[sd.LeaderSuid.Index()].Suid.Epoch() {
oldShard.leaderDiskID = sd.LeaderDiskID
oldShard.leaderSuid = sd.LeaderSuid
return nil, nil
} else {
span.Warnf("skip update shard, leader suid epoch is less than old. old:%d, new:%d", oldShard.leaderSuid.Epoch(), sd.LeaderSuid.Epoch())
return nil, errCatalogNoLeader
}
})
return err
}
@ -451,6 +455,10 @@ func (s *shardControllerImpl) handleShardUpdate(ctx context.Context, item cluste
span.Warnf("catalog skip invalid item update, item:%+v", val)
return errCatalogInvalid
}
if val.Unit.Suid.Epoch() <= info.units[val.Unit.Suid.Index()].Suid.Epoch() {
span.Warnf("catalog skip invalid item update, old:%+v, new:%+v", info.units[val.Unit.Suid.Index()], val)
return errCatalogInvalid
}
// fix leader disk is 0: use new disk unit as leaderDisk, and we will fetch the correct leader later from sn
if errors.Is(err, errCatalogNoLeader) {
@ -459,8 +467,8 @@ func (s *shardControllerImpl) handleShardUpdate(ctx context.Context, item cluste
}
// we will update shard, after fix catalog val
s.setShardByID(info, &val)
span.Debugf("handle one catalog item update:%+v", val)
s.setShardByID(ctx, info, &val)
span.Debugf("handle one catalog item update:%+v, shard:%+v", val, *info)
return err
}
@ -490,7 +498,7 @@ func (s *shardControllerImpl) getShardByID(shardID proto.ShardID) (*shard, bool)
return info, ok
}
func (s *shardControllerImpl) setShardByID(info *shard, val *clustermgr.CatalogChangeShardUpdate) {
func (s *shardControllerImpl) setShardByID(ctx context.Context, info *shard, val *clustermgr.CatalogChangeShardUpdate) {
info.version = val.RouteVersion
// info.rangeExt = val.Unit.Range // todo: will update range next version
@ -510,7 +518,10 @@ func (s *shardControllerImpl) setShardByID(info *shard, val *clustermgr.CatalogC
}
}
// not find leader disk in units
// leader disk not in units, fix it
span := trace.SpanFromContextSafe(ctx)
span.Infof("leader disk not in units. old leader:(%d, %d), new leader:%d, units:%+v",
info.leaderDiskID, info.leaderSuid, val.Unit.LeaderDiskID, info.units)
info.leaderDiskID = info.units[0].DiskID
info.leaderSuid = info.units[0].Suid
}
@ -563,6 +574,9 @@ func (i *shard) GetRange() sharding.Range {
}
func (i *shard) GetMember(ctx context.Context, mode acapi.GetShardMode, exclude proto.DiskID) (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 {
return i.getMemberExcluded(ctx, exclude)
@ -580,25 +594,6 @@ func (i *shard) getMemberExcluded(ctx context.Context, diskID proto.DiskID) (Sha
}
func (i *shard) getMemberLeader(ctx context.Context) (ShardOpInfo, error) {
span := trace.SpanFromContextSafe(ctx)
// leader disk status: normal->EIO->broken->repairing->repaired. it greater than broken will mark punished
// we will mark bad disk punished, at punishAndUpdate(stream_blob.go)
disk, err := i.punishCtrl.GetShardnodeHost(ctx, i.leaderDiskID)
if err != nil {
return ShardOpInfo{}, err
}
// if bad disk(punished), we select another disk as leader.
// and then call sn return NoLeader, and we fetch and update new leader shard
// 1. leader: repaired(eio, or status >= broken), sn will remove disk and return DiskNotFound
// 2. leader: broken(eio, broken, repairing), sn return DiskBroken
if disk.Punished {
span.Warnf("leader is punished: %+v, diskID:%d, suid:%d", disk, i.leaderDiskID, i.leaderSuid)
return i.getMemberRandom(ctx, i.leaderDiskID)
}
// if leader disk is normal, not punished
return ShardOpInfo{
DiskID: i.leaderDiskID,
Suid: i.leaderSuid,
@ -677,17 +672,7 @@ func convertShardUnitInfo(units []clustermgr.ShardUnitInfo) []clustermgr.ShardUn
}
func isInvalidShardStat(sd shardnode.ShardStats) bool {
if sd.Suid == 0 || sd.RouteVersion == 0 || sd.LeaderDiskID == 0 {
return true
}
for _, unit := range sd.Units {
if unit.DiskID == 0 || unit.Suid == 0 {
return true
}
}
if sd.Range.Type == sharding.RangeType_RangeTypeUNKNOWN || sd.Range.IsEmpty() {
if sd.Suid == 0 || sd.RouteVersion == 0 || sd.LeaderDiskID == 0 || sd.LeaderSuid == 0 {
return true
}

View File

@ -158,9 +158,9 @@ func TestShardController(t *testing.T) {
DiskID: 4,
})
err = svr.UpdateShard(ctx, shardnode.ShardStats{
Suid: proto.EncodeSuid(newShard.shardID, 1, 0),
Suid: proto.EncodeSuid(newShard.shardID, 1, 1),
LeaderDiskID: newShard.leaderDiskID,
LeaderSuid: proto.EncodeSuid(newShard.shardID, 1, 0),
LeaderSuid: proto.EncodeSuid(newShard.shardID, 1, 1),
RouteVersion: newShard.version,
Range: newShard.rangeExt,
Units: newShard.units,
@ -174,15 +174,11 @@ 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])
cmCli.EXPECT().ShardNodeDiskInfo(gAny, proto.DiskID(2)).Return(&clustermgr.ShardNodeDiskInfo{
DiskInfo: clustermgr.DiskInfo{Host: "testHost1", Idc: "test-idc"},
ShardNodeDiskHeartbeatInfo: clustermgr.ShardNodeDiskHeartbeatInfo{DiskID: 2},
}, nil)
opInfo, err := si.GetMember(ctx, acapi.GetShardModeLeader, 0)
require.NoError(t, err)
require.Equal(t, ShardOpInfo{
DiskID: 2,
Suid: proto.EncodeSuid(1, 1, 0),
Suid: proto.EncodeSuid(1, 1, 1),
RouteVersion: 1,
}, opInfo)
@ -193,12 +189,12 @@ func TestShardController(t *testing.T) {
newShard.units[2] = newShard.units[3]
newShard.units = newShard.units[:3]
newShard.leaderDiskID = 4
newShard.leaderSuid = proto.EncodeSuid(newShard.shardID, 2, 1)
newShard.leaderSuid = proto.EncodeSuid(newShard.shardID, 2, 2)
newShard.version = si.version
err = svr.UpdateShard(ctx, shardnode.ShardStats{
Suid: newShard.units[2].Suid,
LeaderDiskID: newShard.leaderDiskID,
LeaderSuid: newShard.units[2].Suid,
LeaderSuid: newShard.leaderSuid,
RouteVersion: newShard.version,
Range: newShard.rangeExt,
Units: newShard.units,
@ -326,9 +322,9 @@ func TestShardUpdate(t *testing.T) {
val4 := clustermgr.CatalogChangeShardUpdate{
ShardID: 1,
RouteVersion: 3,
RouteVersion: 4,
Unit: clustermgr.ShardUnitInfo{
Suid: proto.EncodeSuid(1, 0, 1),
Suid: proto.EncodeSuid(1, 0, 2),
DiskID: 3,
LeaderDiskID: 3,
RouteVersion: 3,
@ -367,14 +363,14 @@ func TestShardUpdate(t *testing.T) {
err = svr.UpdateRoute(ctx)
require.Equal(t, errCatalogNoLeader, err)
require.Equal(t, proto.RouteVersion(3), svr.version)
require.Equal(t, proto.RouteVersion(4), svr.version)
require.Equal(t, 2, len(svr.shards))
}
{
// add shard id=9
const shardID = 9
const version = 4
const version = 5
val := clustermgr.CatalogChangeShardAdd{
ShardID: shardID,
RouteVersion: version,
@ -507,6 +503,7 @@ func TestShardUpdate(t *testing.T) {
}
{
// update version, leader disk, suid
rv := svr.version + 1
shardID := proto.ShardID(9)
oldShard, exist := svr.getShardNoLock(shardID)
@ -536,7 +533,8 @@ func TestShardUpdate(t *testing.T) {
Value: data,
},
}
svr.handleShardUpdate(ctx, item)
err = svr.handleShardUpdate(ctx, item)
require.NoError(t, err)
sd, exist := svr.getShardNoLock(shardID)
require.True(t, exist)
@ -561,6 +559,56 @@ func TestShardUpdate(t *testing.T) {
}
{
// update version, leader disk not in units(old leader)
rv := svr.version + 1
shardID := proto.ShardID(9)
oldShard, exist := svr.getShardNoLock(shardID)
require.True(t, exist)
val := clustermgr.CatalogChangeShardUpdate{
ShardID: shardID,
RouteVersion: rv,
Unit: clustermgr.ShardUnitInfo{
Suid: proto.EncodeSuid(shardID, 2, 2),
DiskID: 5,
LeaderDiskID: 3,
Range: *ranges[8],
RouteVersion: rv,
Host: "testHost3",
Learner: false,
},
}
data, err := val.Marshal()
require.NoError(t, err)
item := clustermgr.CatalogChangeItem{
RouteVersion: svr.version + 1,
Type: proto.CatalogChangeItemUpdateShard,
Item: &types.Any{
TypeUrl: "",
Value: data,
},
}
err = svr.handleShardUpdate(ctx, item)
require.NoError(t, err)
sd, exist := svr.getShardNoLock(shardID)
require.True(t, exist)
opHeader, err := sd.GetMember(ctx, acapi.GetShardModeLeader, 0)
require.NoError(t, err)
expect := ShardOpInfo{
DiskID: 1,
Suid: proto.EncodeSuid(shardID, 0, 0),
RouteVersion: rv,
}
require.Equal(t, expect, opHeader)
require.Equal(t, oldShard.leaderDiskID, expect.DiskID)
}
{
// leader disk is 0
rv := svr.version + 1
shardID := proto.ShardID(9)
oldShard, exist := svr.getShardNoLock(shardID)
@ -588,7 +636,8 @@ func TestShardUpdate(t *testing.T) {
Value: data,
},
}
svr.handleShardUpdate(ctx, item)
err = svr.handleShardUpdate(ctx, item)
require.ErrorIs(t, errCatalogNoLeader, err)
sd, exist := svr.getShardNoLock(shardID)
require.True(t, exist)
@ -724,7 +773,7 @@ func TestShardGetShard(t *testing.T) {
require.NoError(t, err)
require.Equal(t, proto.ShardID(1), sd.GetShardID())
newDisk, err := sd.GetMember(ctx, 0, 2)
newDisk, err := sd.GetMember(ctx, acapi.GetShardModeRandom, 2)
require.NoError(t, err)
require.NotEqual(t, proto.DiskID(2), newDisk.DiskID)
require.Contains(t, []proto.DiskID{1, 3}, newDisk.DiskID)

View File

@ -479,8 +479,14 @@ func (h *Handler) punishAndUpdate(ctx context.Context, args *punishArgs) (bool,
// This error is coming from the shardnode interface, and we want to make sure that the error can be parsed into an error code
code := rpc.DetectStatusCode(args.err)
// leader disk status: normal->EIO->broken->repairing->repaired. it greater than broken will mark punished
// if bad disk(punished), we select another disk as leader; old leader is not in disk units, after update route(replace suid index unit)
// and then call sn return NoLeader, and we fetch and update new leader shard
// 1. old leader: repaired(eio, or status >= broken), sn will remove disk and return DiskNotFound, update route
// 2. old leader: broken(eio, broken, repairing), sn return DiskBroken
// cm catalog units is always correct, but its leaderDiskID may be wrong
switch code {
case errcode.CodeDiskBroken: // read disk, but disk is reparing
case errcode.CodeDiskBroken: // read shard at bad disk, but shard/disk is reparing
// if follow node broken disk, it will not election, just try again, change other shard;
// if leader node broken disk, it cant get shard stats, wait new leader
h.punishShardnodeDisk(ctx, args.clusterID, args.DiskID, args.host, "Broken")
@ -493,7 +499,7 @@ func (h *Handler) punishAndUpdate(ctx context.Context, args *punishArgs) (bool,
return false, args.err
// update route and punish
case errcode.CodeShardNodeDiskNotFound: // read old disk, but old broken disk is repaired
case errcode.CodeShardNodeDiskNotFound: // read shard at bad disk, but old broken disk is repaired, all shard repaired
h.punishShardnodeDisk(ctx, args.clusterID, args.DiskID, args.host, "NotFound")
if err1 := h.updateShardRoute(ctx, args.clusterID); err1 != nil {
span.Warnf("fail to update shard route, cluster:%d, err:%+v", args.clusterID, err1)
@ -501,7 +507,7 @@ func (h *Handler) punishAndUpdate(ctx context.Context, args *punishArgs) (bool,
return false, args.err
// update route
case errcode.CodeShardDoesNotExist, // shard is removed, disk is repairing ; suid not match disk id
case errcode.CodeShardDoesNotExist, // intermediate state disk, not a final state; shard is removed, disk is repairing/repaired ; suid not match disk id
errcode.CodeShardRouteVersionNeedUpdate: // header op version less than shardnode version
if err1 := h.updateShardRoute(ctx, args.clusterID); err1 != nil {
span.Warnf("fail to update shard route, cluster:%d, err:%+v", args.clusterID, err1)
@ -608,28 +614,14 @@ func (h *Handler) getLeaderShardInfo(ctx context.Context, clusterID proto.Cluste
return shardnode.ShardStats{}, err
}
// skip bad host. LeaderDiskID means in the election. bad disk is last leader, not start election yet
// skip bad host. LeaderDiskID is 0 means in the election. bad disk is last leader, not start election yet
if leader.LeaderDiskID == 0 || leader.LeaderDiskID == badDisk {
span.Warnf("shard node is in the election, host:%s, disk:%d, suid:%d, badDisk:%d", host, diskID, suid, badDisk)
time.Sleep(time.Millisecond * time.Duration(h.ShardnodeRetryIntervalMS))
continue
}
// 2. get leader ShardNode host
leaderHost, err := h.getShardHost(ctx, clusterID, leader.LeaderDiskID)
if err != nil {
return shardnode.ShardStats{}, err
}
// 3. get leader shard stat, with leader host
ret, err := h.shardnodeClient.GetShardStats(ctx, leaderHost, shardnode.GetShardArgs{
DiskID: leader.LeaderDiskID,
Suid: leader.LeaderSuid,
})
if err != nil {
return shardnode.ShardStats{}, err
}
return ret, nil
return leader, nil
}
return shardnode.ShardStats{}, errcode.ErrShardNoLeader

View File

@ -429,14 +429,13 @@ func TestStreamBlobOther(t *testing.T) {
svrCtrl := NewMockServiceController(ctr)
svrCtrl.EXPECT().PunishShardnode(gAny, gAny, gAny).Times(2)
svrCtrl.EXPECT().GetShardnodeHost(gAny, gAny).Return(&controller.HostIDC{Host: "host"}, nil).Times(1)
shardMgr := NewMockShardController(ctr)
shardMgr.EXPECT().UpdateRoute(gAny).Return(nil).Times(2)
shardMgr.EXPECT().UpdateShard(gAny, gAny).Return(nil)
clu := NewMockClusterController(ctr)
clu.EXPECT().GetServiceController(gAny).Return(svrCtrl, nil).Times(3)
clu.EXPECT().GetServiceController(gAny).Return(svrCtrl, nil).Times(2)
clu.EXPECT().GetShardController(gAny).Return(shardMgr, nil).Times(3)
h := &Handler{
@ -462,7 +461,7 @@ func TestStreamBlobOther(t *testing.T) {
require.ErrorIs(t, err1, io.EOF)
shardnodeClient := mocks.NewMockShardnodeAccess(ctr)
shardnodeClient.EXPECT().GetShardStats(gAny, gAny, gAny).Return(shardnode.ShardStats{LeaderDiskID: 11}, nil).Times(2)
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{
@ -488,20 +487,18 @@ func TestStreamBlobOther(t *testing.T) {
RouteVersion: 1,
}
clu.EXPECT().GetShardController(gAny).Return(shardMgr, nil)
clu.EXPECT().GetServiceController(gAny).Return(svrCtrl, nil).Times(3)
clu.EXPECT().GetServiceController(gAny).Return(svrCtrl, nil).Times(2)
shardInfo := NewMockShard(ctr)
shardInfo.EXPECT().GetMember(gAny, gAny, proto.DiskID(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)
svrCtrl.EXPECT().GetShardnodeHost(gAny, proto.DiskID(102)).Return(&controller.HostIDC{Host: "host102"}, nil)
svrCtrl.EXPECT().PunishShardnode(gAny, gAny, gAny)
shardnodeClient.EXPECT().GetShardStats(gAny, gAny, gAny).Return(shardnode.ShardStats{LeaderDiskID: 1}, nil)
shardnodeClient.EXPECT().GetShardStats(gAny, "host101", shardnode.GetShardArgs{
DiskID: proto.DiskID(101),
Suid: info.Suid,
}).Return(shardnode.ShardStats{LeaderDiskID: 102}, nil)
shardnodeClient.EXPECT().GetShardStats(gAny, "host102", gAny).Return(shardnode.ShardStats{LeaderDiskID: 102}, nil)
interrupt, err1 = h.punishAndUpdate(ctx, &punishArgs{
ShardOpHeader: shardnode.ShardOpHeader{
DiskID: 1,

View File

@ -440,6 +440,7 @@ func (s *sdkHandler) getBlobData(ctx context.Context, args *acapi.GetArgs) (io.R
}
defer s.limiter.Release(name)
// TODO next version, supports GetBlob data that has not yet been sealed
// means blob not seal, not support get data, return error
if args.Location.Size_ == 0 {
return noopBody{}, errcode.ErrReaderError