diff --git a/blobstore/common/proto/shard.go b/blobstore/common/proto/shard.go index d4d71ed93..acdb4b4b4 100644 --- a/blobstore/common/proto/shard.go +++ b/blobstore/common/proto/shard.go @@ -119,6 +119,7 @@ const ( ShardTaskTypeClearShard = ShardTaskType(iota + 1) ShardTaskTypeCheckpoint ShardTaskTypeSyncRouteVersion + ShardTaskTypeCheckAndClear ) type ShardUnitStatus uint8 diff --git a/blobstore/shardnode/base/mock_transport.go b/blobstore/shardnode/base/mock_transport.go index 6eb137ed8..fce41c28e 100644 --- a/blobstore/shardnode/base/mock_transport.go +++ b/blobstore/shardnode/base/mock_transport.go @@ -333,6 +333,21 @@ func (mr *MockTransportMockRecorder) ShardReport(ctx, reports interface{}) *gomo return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ShardReport", reflect.TypeOf((*MockTransport)(nil).ShardReport), ctx, reports) } +// ShardStats mocks base method. +func (m *MockTransport) ShardStats(ctx context.Context, host string, args shardnode.GetShardArgs) (shardnode.ShardStats, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ShardStats", ctx, host, args) + ret0, _ := ret[0].(shardnode.ShardStats) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ShardStats indicates an expected call of ShardStats. +func (mr *MockTransportMockRecorder) ShardStats(ctx, host, args interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ShardStats", reflect.TypeOf((*MockTransport)(nil).ShardStats), ctx, host, args) +} + // UpdateShard mocks base method. func (m *MockTransport) UpdateShard(ctx context.Context, host string, args shardnode.UpdateShardArgs) error { m.ctrl.T.Helper() @@ -688,6 +703,21 @@ func (mr *MockShardTransportMockRecorder) ResolveRaftAddr(ctx, diskID interface{ return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ResolveRaftAddr", reflect.TypeOf((*MockShardTransport)(nil).ResolveRaftAddr), ctx, diskID) } +// ShardStats mocks base method. +func (m *MockShardTransport) ShardStats(ctx context.Context, host string, args shardnode.GetShardArgs) (shardnode.ShardStats, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "ShardStats", ctx, host, args) + ret0, _ := ret[0].(shardnode.ShardStats) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// ShardStats indicates an expected call of ShardStats. +func (mr *MockShardTransportMockRecorder) ShardStats(ctx, host, args interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ShardStats", reflect.TypeOf((*MockShardTransport)(nil).ShardStats), ctx, host, args) +} + // UpdateShard mocks base method. func (m *MockShardTransport) UpdateShard(ctx context.Context, host string, args shardnode.UpdateShardArgs) error { m.ctrl.T.Helper() diff --git a/blobstore/shardnode/base/transport.go b/blobstore/shardnode/base/transport.go index 03196087e..9c4e25346 100644 --- a/blobstore/shardnode/base/transport.go +++ b/blobstore/shardnode/base/transport.go @@ -67,6 +67,7 @@ type ( ResolveRaftAddr(ctx context.Context, diskID proto.DiskID) (string, error) ResolveNodeAddr(ctx context.Context, diskID proto.DiskID) (string, error) UpdateShard(ctx context.Context, host string, args shardnodeapi.UpdateShardArgs) error + ShardStats(ctx context.Context, host string, args shardnodeapi.GetShardArgs) (shardnodeapi.ShardStats, error) } ) @@ -295,3 +296,7 @@ func (t *transport) ResolveNodeAddr(ctx context.Context, diskID proto.DiskID) (s func (t *transport) UpdateShard(ctx context.Context, host string, args shardnodeapi.UpdateShardArgs) error { return t.snClient.UpdateShard(ctx, host, args) } + +func (t *transport) ShardStats(ctx context.Context, host string, args shardnodeapi.GetShardArgs) (shardnodeapi.ShardStats, error) { + return t.snClient.GetShardStats(ctx, host, args) +} diff --git a/blobstore/shardnode/mock/mock_shard.go b/blobstore/shardnode/mock/mock_shard.go index d107363a5..3fa75bba2 100644 --- a/blobstore/shardnode/mock/mock_shard.go +++ b/blobstore/shardnode/mock/mock_shard.go @@ -279,6 +279,20 @@ func (m *MockSpaceShardHandler) EXPECT() *MockSpaceShardHandlerMockRecorder { return m.recorder } +// CheckAndClearShard mocks base method. +func (m *MockSpaceShardHandler) CheckAndClearShard(ctx context.Context) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "CheckAndClearShard", ctx) + ret0, _ := ret[0].(error) + return ret0 +} + +// CheckAndClearShard indicates an expected call of CheckAndClearShard. +func (mr *MockSpaceShardHandlerMockRecorder) CheckAndClearShard(ctx interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "CheckAndClearShard", reflect.TypeOf((*MockSpaceShardHandler)(nil).CheckAndClearShard), ctx) +} + // Checkpoint mocks base method. func (m *MockSpaceShardHandler) Checkpoint(ctx context.Context) error { m.ctrl.T.Helper() diff --git a/blobstore/shardnode/shard.go b/blobstore/shardnode/shard.go index ce40fdd19..85ea27b86 100644 --- a/blobstore/shardnode/shard.go +++ b/blobstore/shardnode/shard.go @@ -180,6 +180,7 @@ func (s *service) loop(ctx context.Context) { reportTicker := time.NewTicker(time.Duration(s.cfg.ReportIntervalS) * time.Second) routeUpdateTicker := time.NewTicker(time.Duration(s.cfg.RouteUpdateIntervalS) * time.Second) checkpointTicker := time.NewTicker(time.Duration(s.cfg.CheckPointIntervalM) * time.Minute) + trashShardCheckTicker := time.NewTicker(time.Duration(s.cfg.ShardCheckAndClearIntervalH) * time.Hour) defer func() { heartbeatTicker.Stop() @@ -192,6 +193,7 @@ func (s *service) loop(ctx context.Context) { diskReports := make([]clustermgr.ShardNodeDiskHeartbeatInfo, 0) shardReports := make([]clustermgr.ShardUnitInfo, 0, 1<<10) shards := make([]storage.ShardHandler, 0, 1<<10) + tasks := make([]clustermgr.ShardTask, 0, 1<<10) for { select { @@ -263,20 +265,44 @@ func (s *service) loop(ctx context.Context) { case <-checkpointTicker.C: span, ctx = trace.StartSpanFromContext(ctx, "do checkpoint") disks := s.getAllDisks() + tasks = tasks[:0] for _, disk := range disks { shards = shards[:0] disk.RangeShard(func(s storage.ShardHandler) bool { - shards = append(shards, s) + tasks = append(tasks, clustermgr.ShardTask{ + TaskType: proto.ShardTaskTypeCheckpoint, + Suid: s.GetSuid(), + DiskID: disk.DiskID(), + }) return true }) - for _, shard := range shards { - err := shard.Checkpoint(ctx) - if err != nil { - span.Errorf("do checkpoint failed: %s", errors.Detail(err)) + for _, task := range tasks { + if err := s.executeShardTask(ctx, task); err != nil { + span.Errorf("execute shard task[%+v] failed: %s", task, errors.Detail(err)) continue } } } + case <-trashShardCheckTicker.C: + span, ctx = trace.StartSpanFromContext(ctx, "trash shard check") + disks := s.getAllDisks() + tasks = tasks[:0] + for _, disk := range disks { + disk.RangeShard(func(s storage.ShardHandler) bool { + tasks = append(tasks, clustermgr.ShardTask{ + TaskType: proto.ShardTaskTypeCheckAndClear, + Suid: s.GetSuid(), + DiskID: disk.DiskID(), + }) + return true + }) + } + for _, task := range tasks { + if err := s.executeShardTask(ctx, task); err != nil { + span.Errorf("execute shard task[%+v] failed: %s", task, errors.Detail(err)) + continue + } + } case <-s.closer.Done(): return } @@ -325,6 +351,21 @@ func (s *service) executeShardTask(ctx context.Context, task clustermgr.ShardTas curVersion, task.OldRouteVersion, task.RouteVersion) } }) + case proto.ShardTaskTypeCheckAndClear: + s.taskPool.Run(func() { + _span, _ctx := trace.StartSpanFromContextWithTraceID(ctx, "", "shard-check-"+task.Suid.ToString()) + if err := shard.CheckAndClearShard(_ctx); err != nil { + _span.Errorf("check trash shard task[%+v] failed: %s", task, errors.Detail(err)) + } + }) + case proto.ShardTaskTypeCheckpoint: + s.taskPool.Run(func() { + _span, _ctx := trace.StartSpanFromContextWithTraceID(ctx, "", "checkpoint-"+task.Suid.ToString()) + if err := shard.Checkpoint(_ctx); err != nil { + _span.Errorf("shard do checkpoint task[%+v] failed: %s", task, errors.Detail(err)) + } + }) + default: } return nil diff --git a/blobstore/shardnode/startup.go b/blobstore/shardnode/startup.go index 0fb1b0dc2..1f1003864 100644 --- a/blobstore/shardnode/startup.go +++ b/blobstore/shardnode/startup.go @@ -373,4 +373,5 @@ func initServiceConfig(cfg *Config) { defaulter.LessOrEqual(&cfg.CheckPointIntervalM, int64(1)) defaulter.LessOrEqual(&cfg.WaitRepairCloseDiskIntervalS, int64(30)) defaulter.LessOrEqual(&cfg.WaitReOpenDiskIntervalS, int64(30)) + defaulter.LessOrEqual(&cfg.ShardCheckAndClearIntervalH, int64(24)) } diff --git a/blobstore/shardnode/storage/shard.go b/blobstore/shardnode/storage/shard.go index b6e8d5359..002ecf691 100644 --- a/blobstore/shardnode/storage/shard.go +++ b/blobstore/shardnode/storage/shard.go @@ -27,6 +27,7 @@ import ( kvstore "github.com/cubefs/cubefs/blobstore/common/kvstorev2" "github.com/cubefs/cubefs/blobstore/common/proto" "github.com/cubefs/cubefs/blobstore/common/raft" + "github.com/cubefs/cubefs/blobstore/common/rpc" "github.com/cubefs/cubefs/blobstore/common/sharding" "github.com/cubefs/cubefs/blobstore/common/trace" "github.com/cubefs/cubefs/blobstore/shardnode/base" @@ -73,6 +74,7 @@ type ( Stats(ctx context.Context) (shardnode.ShardStats, error) GetSuid() proto.Suid GetUnits() []clustermgr.ShardUnit + CheckAndClearShard(ctx context.Context) error } OpHeader struct { RouteVersion proto.RouteVersion @@ -462,7 +464,8 @@ func (s *shard) Stats(ctx context.Context) (shardnode.ShardStats, error) { defer s.shardState.prepRWCheckDone() s.shardInfoMu.RLock() - units := s.shardInfoMu.Units + units := make([]clustermgr.ShardUnit, len(s.shardInfoMu.Units)) + copy(units, s.shardInfoMu.Units) routeVersion := s.shardInfoMu.RouteVersion appliedIndex := s.shardInfoMu.AppliedIndex rg := s.shardInfoMu.Range @@ -696,6 +699,62 @@ func (s *shard) DeleteShard(ctx context.Context, nodeHost string, clearData bool return kvStore.FlushCF(ctx, dataCF) } +func (s *shard) CheckAndClearShard(ctx context.Context) error { + span := trace.SpanFromContext(ctx) + + stats, err := s.Stats(ctx) + if err != nil { + return err + } + if stats.LeaderDiskID != s.diskID { + return nil + } + + raftStats := make(map[uint64]bool) + for _, p := range stats.RaftStat.Peers { + raftStats[p.NodeID] = p.RecentActive + } + + for _, u := range stats.Units { + if u.Suid == s.GetSuid() { + continue + } + + var active, ok bool + if active, ok = raftStats[uint64(u.DiskID)]; !ok { + return errors.Newf("disk[%d] not found in raft peers", u.DiskID) + } + if active { + continue + } + + host, err := s.cfg.Transport.ResolveNodeAddr(ctx, u.DiskID) + if err != nil { + return errors.Info(err, "resolve node address failed", u.DiskID) + } + _, err = s.cfg.Transport.ShardStats(ctx, host, shardnodeapi.GetShardArgs{ + DiskID: u.DiskID, + Suid: u.Suid, + }) + if err == nil { + continue + } + if rpc.DetectStatusCode(err) != apierr.CodeShardNodeDiskNotFound { + return errors.Info(err, "get shard stats failed", u.DiskID, u.Suid) + } + err = s.UpdateShard(ctx, proto.ShardUpdateTypeRemoveMember, clustermgr.ShardUnit{ + Suid: u.Suid, + DiskID: u.DiskID, + }, host) + if err != nil { + return errors.Info(err, "update shard failed", s.suid, u.DiskID) + } + span.Infof("remove shard[%d] suid[%d] from unit[%+v] done", s.suid.ShardID(), s.suid, u) + return err + } + return nil +} + func (s *shard) Start() { // Do nothing because need to be improved later } diff --git a/blobstore/shardnode/svr.go b/blobstore/shardnode/svr.go index 5abc8de60..f84f1a82c 100644 --- a/blobstore/shardnode/svr.go +++ b/blobstore/shardnode/svr.go @@ -71,6 +71,7 @@ type Config struct { CheckPointIntervalM int64 `json:"check_point_interval_m"` WaitRepairCloseDiskIntervalS int64 `json:"wait_repair_close_disk_interval_s"` WaitReOpenDiskIntervalS int64 `json:"wait_re_open_disk_interval_s"` + ShardCheckAndClearIntervalH int64 `json:"shard_check_and_clear_interval_h"` } func newService(cfg *Config) *service {