feat(shardnode): init execute repair message process

with: #1000420040

Signed-off-by: xiejian <xiejian3@oppo.com>
This commit is contained in:
xiejian 2025-11-11 17:33:08 +08:00 committed by 梁曟風
parent 073b7aee84
commit 4d5b05ec0d
24 changed files with 1916 additions and 270 deletions

View File

@ -34,7 +34,26 @@ func (c *Client) DeleteBlobRaw(ctx context.Context, host string, args DeleteBlob
return c.doRequest(ctx, host, "/blob/delete/raw", &args, nil)
}
func (c *Client) DeleteBlobStats(ctx context.Context, host string, args DeleteBlobStatsArgs) (ret DeleteBlobStatsRet, err error) {
func (c *Client) DeleteBlobStats(ctx context.Context, host string, args ShardnodeTaskStatsArgs) (ret ShardnodeTaskStatsRet, err error) {
err = c.doRequest(ctx, host, "/blob/delete/stats", &args, &ret)
return
}
func (r *RepairSliceArgs) GetShardKeys(tagNum int) []string {
if r == nil {
return nil
}
keys := make([]string, util.Max(2, tagNum))
keys[0] = util.Any2String(r.Vid)
keys[1] = util.Any2String(r.Bid)
return keys
}
func (c *Client) RepairSlice(ctx context.Context, host string, args RepairSliceArgs) error {
return c.doRequest(ctx, host, "/slice/repair", &args, nil)
}
func (c *Client) RepairSliceStats(ctx context.Context, host string, args ShardnodeTaskStatsArgs) (ret ShardnodeTaskStatsRet, err error) {
err = c.doRequest(ctx, host, "/slice/repair/stats", &args, &ret)
return
}

File diff suppressed because it is too large Load Diff

View File

@ -262,9 +262,19 @@ message DeleteBlobRawArgs {
message DeleteBlobRawRet {}
message DeleteBlobStatsArgs {}
message RepairSliceArgs {
ShardOpHeader header = 1 [(gogoproto.nullable) = false];
uint64 bid = 2 [(gogoproto.customname) = "Bid", (gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.BlobID"];
uint32 vid = 3 [(gogoproto.customname) = "Vid", (gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.Vid"];
repeated uint32 bad_idx = 4 [(gogoproto.customname) = "BadIdx"];
string reason = 5;
}
message DeleteBlobStatsRet {
message RepairSliceRet {}
message ShardnodeTaskStatsArgs {}
message ShardnodeTaskStatsRet {
bool enable = 1;
string success_per_min = 2;
string failed_per_min = 3;

View File

@ -30,8 +30,8 @@ var types = map[string]func() rpc2.Codec{
"shardnode.GetBlobArgs": func() rpc2.Codec { return new(shardnode.GetBlobArgs) },
"shardnode.GetBlobRet": func() rpc2.Codec { return new(shardnode.GetBlobRet) },
"shardnode.DeleteBlobStatsArgs": func() rpc2.Codec { return new(shardnode.DeleteBlobStatsArgs) },
"shardnode.DeleteBlobStatsRet": func() rpc2.Codec { return new(shardnode.DeleteBlobStatsRet) },
"shardnode.ShardnodeTaskStatsArgs": func() rpc2.Codec { return new(shardnode.ShardnodeTaskStatsArgs) },
"shardnode.ShardnodeTaskStatsRet": func() rpc2.Codec { return new(shardnode.ShardnodeTaskStatsRet) },
"nil": func() rpc2.Codec { return nil },
"rpc2.NoParameter": func() rpc2.Codec { return rpc2.NoParameter },

View File

@ -19,6 +19,10 @@ import (
"sort"
"strings"
"sync"
"github.com/prometheus/client_golang/prometheus"
"github.com/cubefs/cubefs/blobstore/common/proto"
)
// ErrorStats error stats
@ -98,3 +102,75 @@ func errStrFormat(err error) string {
strSlice := strings.Split(err.Error(), ":")
return strings.TrimSpace(strSlice[len(strSlice)-1])
}
const (
Namespace = "shardnode"
ShardRepair = "shard_repair" // same as scheduler
ChunkMissMigrateAbnormal = "chunk_miss_migrate"
)
type AbnormalReporter struct {
lock sync.RWMutex
abnormalReporter *prometheus.GaugeVec
reportedVuids map[proto.Vuid]struct{}
}
// NewAbnormalReporter returns abnormal reporter
func NewAbnormalReporter(clusterID proto.ClusterID, taskType string, abnormalKind string) *AbnormalReporter {
abnormalReporter := prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: Namespace,
Subsystem: "task",
Name: "abnormal_task",
Help: "abnormal task",
ConstLabels: map[string]string{
"cluster_id": fmt.Sprintf("%d", clusterID),
"task_type": taskType,
"kind": abnormalKind,
},
}, []string{"diskID", "vuid"})
if err := prometheus.Register(abnormalReporter); err != nil {
if are, ok := err.(prometheus.AlreadyRegisteredError); ok {
abnormalReporter = are.ExistingCollector.(*prometheus.GaugeVec)
return &AbnormalReporter{
abnormalReporter: abnormalReporter,
reportedVuids: make(map[proto.Vuid]struct{}),
}
}
panic(err)
}
return &AbnormalReporter{
abnormalReporter: abnormalReporter,
reportedVuids: make(map[proto.Vuid]struct{}),
}
}
// ReportAbnormal report abnormal task
func (abr *AbnormalReporter) ReportAbnormal(diskID proto.DiskID, vuid proto.Vuid) {
abr.abnormalReporter.WithLabelValues(
fmt.Sprintf("%d", diskID),
fmt.Sprintf("%d", vuid),
).Set(1)
}
// CancelAbnormal cancel abnormal report
func (abr *AbnormalReporter) CancelAbnormal(diskID proto.DiskID, vuid proto.Vuid) {
abr.abnormalReporter.WithLabelValues(
fmt.Sprintf("%d", diskID),
fmt.Sprintf("%d", vuid),
).Set(0)
}
// IsVuidReported check if vuid abnormal is reported
func (abr *AbnormalReporter) IsVuidReported(vuid proto.Vuid) bool {
abr.lock.RLock()
defer abr.lock.RUnlock()
_, ok := abr.reportedVuids[vuid]
return ok
}
// SetVuidReported set vuid reported
func (abr *AbnormalReporter) SetVuidReported(vuid proto.Vuid) {
abr.lock.Lock()
abr.reportedVuids[vuid] = struct{}{}
abr.lock.Unlock()
}

View File

@ -20,6 +20,8 @@ import (
"testing"
"github.com/stretchr/testify/require"
"github.com/cubefs/cubefs/blobstore/common/proto"
)
func TestErrorStats(t *testing.T) {
@ -61,3 +63,15 @@ func TestErrStrFormat(t *testing.T) {
require.Equal(t, "fake error", errStrFormat(err2))
require.Equal(t, "", errStrFormat(err3))
}
func TestAbnormalReport(t *testing.T) {
rp := NewAbnormalReporter(proto.ClusterID(1), "", "")
diskID := proto.DiskID(1)
vuid := proto.Vuid(100)
rp.ReportAbnormal(diskID, vuid)
rp.CancelAbnormal(diskID, vuid)
rp.SetVuidReported(vuid)
require.True(t, rp.IsVuidReported(vuid))
require.False(t, rp.IsVuidReported(vuid+1))
}

View File

@ -34,6 +34,8 @@ type (
GetConfig(ctx context.Context, key string) (string, error)
ShardReport(ctx context.Context, reports []clustermgr.ShardUnitInfo) ([]clustermgr.ShardTask, error)
GetRouteUpdate(ctx context.Context, routeVersion proto.RouteVersion) (proto.RouteVersion, []clustermgr.CatalogChangeItem, error)
GetService(ctx context.Context, name string) ([]string, error)
GetBlobnodeDiskInfo(ctx context.Context, diskID proto.DiskID) (*clustermgr.BlobNodeDiskInfo, error)
NodeTransport
SpaceTransport
AllocVolTransport
@ -78,6 +80,7 @@ type (
BlobTransport interface {
DeleteSliceUnit(ctx context.Context, info proto.VunitLocation, bid proto.BlobID) (err error)
MarkDeleteSliceUnit(ctx context.Context, info proto.VunitLocation, bid proto.BlobID) (err error)
RepairSlice(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) (err error)
}
VolumeTransport interface {
@ -291,6 +294,18 @@ func (t *transport) GetConfig(ctx context.Context, key string) (string, error) {
return t.cmClient.GetConfig(ctx, key)
}
func (t *transport) GetService(ctx context.Context, name string) ([]string, error) {
nodes, err := t.cmClient.GetService(ctx, clustermgr.GetServiceArgs{Name: name})
if err != nil {
return nil, err
}
hosts := make([]string, 0, len(nodes.Nodes))
for _, node := range nodes.Nodes {
hosts = append(hosts, node.Host)
}
return hosts, nil
}
func (t *transport) AllocBid(ctx context.Context, count uint64) (proto.BlobID, error) {
ret, err := t.cmClient.AllocBid(ctx, &clustermgr.BidScopeArgs{Count: count})
if err != nil {
@ -360,6 +375,17 @@ func (t *transport) MarkDeleteSliceUnit(ctx context.Context, info proto.VunitLoc
})
}
func (t *transport) RepairSlice(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) (err error) {
task := proto.ShardRepairTask{
Bid: repairMsg.Bid,
CodeMode: volInfo.CodeMode,
Sources: volInfo.VunitLocations,
BadIdxs: sliceUint32ToInt(repairMsg.BadIdx),
Reason: repairMsg.Reason,
}
return t.bnClient.RepairShard(ctx, host, &task)
}
func (t *transport) ListVolume(ctx context.Context, marker proto.Vid, count int) (volInfo []*snproto.VolumeInfoSimple, retVid proto.Vid, err error) {
vols, err := t.cmClient.ListVolume(ctx, &clustermgr.ListVolumeArgs{Marker: marker, Count: count})
if err != nil {
@ -383,3 +409,15 @@ func (t *transport) GetVolumeInfo(ctx context.Context, vid proto.Vid) (ret *snpr
ret.Set(info)
return ret, nil
}
func (t *transport) GetBlobnodeDiskInfo(ctx context.Context, diskID proto.DiskID) (*clustermgr.BlobNodeDiskInfo, error) {
return t.cmClient.DiskInfo(ctx, diskID)
}
func sliceUint32ToInt(s []uint32) []uint8 {
var ret []uint8
for _, e := range s {
ret = append(ret, uint8(e))
}
return ret
}

View File

@ -114,6 +114,14 @@ func (s *service) deleteBlobRaw(ctx context.Context, req *shardnode.DeleteBlobRa
return s.blobDelMgr.Delete(ctx, req)
}
func (s *service) deleteBlobStats() *shardnode.DeleteBlobStatsRet {
func (s *service) deleteBlobStats() *shardnode.ShardnodeTaskStatsRet {
return s.blobDelMgr.Stats()
}
func (s *service) repairSlice(ctx context.Context, req *shardnode.RepairSliceArgs) error {
return s.sliceRepairMgr.Repair(ctx, req)
}
func (s *service) repairSliceStats() *shardnode.ShardnodeTaskStatsRet {
return s.sliceRepairMgr.Stats()
}

View File

@ -15,9 +15,6 @@
package shardnode
import (
"encoding/json"
"net/http"
"github.com/cubefs/cubefs/blobstore/api/shardnode"
"github.com/cubefs/cubefs/blobstore/common/rpc"
"github.com/cubefs/cubefs/blobstore/common/trace"
@ -43,25 +40,17 @@ func (s *HttpService) HttpShardStats(c *rpc.Context) {
c.RespondError(err)
return
}
data, err := json.Marshal(ret)
if err != nil {
c.RespondError(err)
return
}
c.RespondWith(http.StatusOK, rpc.MIMEJSON, data)
c.RespondJSON(ret)
}
func (s *HttpService) HttpDeleteBlobStats(c *rpc.Context) {
ret := s.deleteBlobStats()
data, err := json.Marshal(ret)
if err != nil {
c.RespondError(err)
return
}
c.RespondJSON(ret)
}
c.RespondWith(http.StatusOK, rpc.MIMEJSON, data)
func (s *HttpService) HttpRepairSliceStats(c *rpc.Context) {
ret := s.repairSliceStats()
c.RespondJSON(ret)
}
func newHttpHandler(service *HttpService) *rpc.Router {
@ -69,6 +58,7 @@ func newHttpHandler(service *HttpService) *rpc.Router {
rpc.GET("/shard/stats", service.HttpShardStats, rpc.OptArgsQuery())
rpc.GET("/blob/delete/stats", service.HttpDeleteBlobStats)
rpc.GET("/slice/repair/stats", service.HttpRepairSliceStats)
return rpc.DefaultRouter
}

View File

@ -84,12 +84,20 @@ func TestHttpService_HTTPGet(t *testing.T) {
require.Equal(t, suid, ret1.Suid)
// Test delete blob stats
ret2 := &shardnode.DeleteBlobStatsRet{}
ret2 := &shardnode.ShardnodeTaskStatsRet{}
resp2, err := client.Get(ctx, "http://127.0.0.1:11000/blob/delete/stats")
require.Nil(t, err)
json.NewDecoder(resp2.Body).Decode(ret2)
resp2.Body.Close()
require.NotNil(t, ret2)
// Test repair shard stats
ret3 := &shardnode.ShardnodeTaskStatsRet{}
resp3, err := client.Get(ctx, "http://127.0.0.1:11000/slice/repair/stats")
require.Nil(t, err)
json.NewDecoder(resp3.Body).Decode(ret3)
resp3.Body.Close()
require.NotNil(t, ret3)
}
// Test setUpHttp function

View File

@ -1,43 +0,0 @@
// Copyright 2025 The CubeFS Authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
// implied. See the License for the specific language governing
// permissions and limitations under the License.
package message
import (
"context"
"fmt"
snapi "github.com/cubefs/cubefs/blobstore/api/shardnode"
"github.com/cubefs/cubefs/blobstore/common/proto"
)
func (m *BlobDeleteMgr) SlicesToDeleteMsgItems(ctx context.Context, slices []proto.Slice, shardKeys []string) ([]snapi.Item, error) {
return m.slicesToDeleteMsgItems(ctx, slices, shardKeys)
}
func (m *BlobDeleteMgr) Delete(ctx context.Context, req *snapi.DeleteBlobRawArgs) error {
return m.insertDeleteMsg(ctx, req)
}
func (m *BlobDeleteMgr) Stats() *snapi.DeleteBlobStatsRet {
deleteSuccessCounter, deleteFailedCounter := m.getTaskStats()
delErrStats, delTotalErrCnt := m.getErrorStats()
return &snapi.DeleteBlobStatsRet{
Enable: m.enabled(),
SuccessPerMin: fmt.Sprint(deleteSuccessCounter),
FailedPerMin: fmt.Sprint(deleteFailedCounter),
TotalErrCnt: delTotalErrCnt,
ErrStats: delErrStats,
}
}

View File

@ -16,6 +16,7 @@ package message
import (
"context"
"fmt"
"time"
snapi "github.com/cubefs/cubefs/blobstore/api/shardnode"
@ -58,6 +59,26 @@ func (m *BlobDeleteMgr) Start() {
m.messageMgr.run()
}
func (m *BlobDeleteMgr) SlicesToDeleteMsgItems(ctx context.Context, slices []proto.Slice, shardKeys []string) ([]snapi.Item, error) {
return m.slicesToDeleteMsgItems(ctx, slices, shardKeys)
}
func (m *BlobDeleteMgr) Delete(ctx context.Context, req *snapi.DeleteBlobRawArgs) error {
return m.insertDeleteMsg(ctx, req)
}
func (m *BlobDeleteMgr) Stats() *snapi.ShardnodeTaskStatsRet {
deleteSuccessCounter, deleteFailedCounter := m.getTaskStats()
delErrStats, delTotalErrCnt := m.getErrorStats()
return &snapi.ShardnodeTaskStatsRet{
Enable: m.enabled(),
SuccessPerMin: fmt.Sprint(deleteSuccessCounter),
FailedPerMin: fmt.Sprint(deleteFailedCounter),
TotalErrCnt: delTotalErrCnt,
ErrStats: delErrStats,
}
}
func (m *BlobDeleteMgr) ItemToMessageExt(item interface{}) (snproto.MessageExt, error) {
msgItem, ok := item.(messageItem)
if !ok {
@ -259,10 +280,10 @@ func (m *BlobDeleteMgr) deleteSliceUnit(ctx context.Context, info proto.VunitLoc
if markerDel {
stage = DeleteStageMarkDelete
err = m.cfg.Transport.MarkDeleteSliceUnit(ctx, info, sliceId)
err = m.cfg.BlobTransport.MarkDeleteSliceUnit(ctx, info, sliceId)
} else {
stage = DeleteStageDelete
err = m.cfg.Transport.DeleteSliceUnit(ctx, info, sliceId)
err = m.cfg.BlobTransport.DeleteSliceUnit(ctx, info, sliceId)
}
if err != nil {

View File

@ -51,7 +51,7 @@ func TestNewTestBlobDeleteMgr(t *testing.T) {
cfg := &BlobDelMgrConfig{
TaskSwitchMgr: taskSwitchMgr,
ShardGetter: sg,
Transport: mocks.NewMockBlobTransport(ctr(t)),
BlobTransport: mocks.NewMockBlobTransport(ctr(t)),
VolCache: mocks.NewMockIVolumeCache(ctr(t)),
MessageCfg: MessageCfg{
FailedMsgChannelSize: 5,
@ -337,8 +337,8 @@ func TestBlobDeleteMgr_DeleteSlice_UpdateVolume(t *testing.T) {
tp.EXPECT().MarkDeleteSliceUnit(any, any, any).Return(nil).Times(3)
mgr := newMinimalBlobDeleteMgr(&MessageMgrConfig{
Transport: tp,
VolCache: vc,
BlobTransport: tp,
VolCache: vc,
})
msgExt := &delMsgExt{msg: snproto.DeleteMsg{
@ -368,8 +368,8 @@ func TestBlobDeleteMgr_DeleteSlice_UpdateVolume2(t *testing.T) {
}).Times(6)
mgr := newMinimalBlobDeleteMgr(&MessageMgrConfig{
Transport: tp,
VolCache: vc,
BlobTransport: tp,
VolCache: vc,
})
msgExt := &delMsgExt{msg: snproto.DeleteMsg{
@ -389,7 +389,7 @@ func TestBlobDeleteMgr_DeleteShard_BackToInitStage(t *testing.T) {
volInfo := newMockSimpleVolumeInfo(vid)
mgr := newMinimalBlobDeleteMgr(&MessageMgrConfig{
Transport: tp,
BlobTransport: tp,
})
idx := 0
@ -415,7 +415,7 @@ func TestBlobDeleteMgr_DeleteShard_AssumeSuccess(t *testing.T) {
volInfo := newMockSimpleVolumeInfo(vid)
mgr := newMinimalBlobDeleteMgr(&MessageMgrConfig{
Transport: tp,
BlobTransport: tp,
})
idx := 0

View File

@ -84,7 +84,7 @@ type MessageMgrConfig struct {
executor MessageExecutor
TaskSwitchMgr *taskswitch.SwitchMgr
ShardGetter ShardGetter
Transport base.BlobTransport
BlobTransport base.BlobTransport
VolCache base.IVolumeCache
MessageCfg
}

View File

@ -58,7 +58,7 @@ func TestBlobMessageMgr_New(t *testing.T) {
cfg := &BlobDelMgrConfig{
TaskSwitchMgr: taskSwitchMgr,
ShardGetter: sg,
Transport: mocks.NewMockBlobTransport(ctr(t)),
BlobTransport: mocks.NewMockBlobTransport(ctr(t)),
VolCache: mocks.NewMockIVolumeCache(ctr(t)),
MessageCfg: MessageCfg{
FailedMsgChannelSize: 5,
@ -531,10 +531,10 @@ func newTestPriorityConfig() TierConfig {
func newTestBlobMessageMgr(t *testing.T, sg ShardGetter, tp base.BlobTransport, vc base.IVolumeCache, executor MessageExecutor, messageType snproto.MessageType) *messageMgr {
cfg := &MessageMgrConfig{
executor: executor,
ShardGetter: sg,
Transport: tp,
VolCache: vc,
executor: executor,
ShardGetter: sg,
BlobTransport: tp,
VolCache: vc,
MessageCfg: MessageCfg{
RetryTimes: 3,
FailedMsgChannelSize: 5,

View File

@ -16,32 +16,65 @@ package message
import (
"context"
"fmt"
"time"
"github.com/cubefs/cubefs/blobstore/api/scheduler"
snapi "github.com/cubefs/cubefs/blobstore/api/shardnode"
apierr "github.com/cubefs/cubefs/blobstore/common/errors"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/rpc2"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/shardnode/base"
snproto "github.com/cubefs/cubefs/blobstore/shardnode/proto"
"github.com/cubefs/cubefs/blobstore/shardnode/storage"
"github.com/cubefs/cubefs/blobstore/util/errors"
"github.com/cubefs/cubefs/blobstore/util/retry"
"github.com/cubefs/cubefs/blobstore/util/selector"
)
// ErrBlobnodeServiceUnavailable worker service unavailable
var ErrBlobnodeServiceUnavailable = errors.New("blobnode service unavailable")
// OrphanSlice orphan slice identification.
type OrphanSlice struct {
ClusterID proto.ClusterID `json:"cluster_id"`
Vid proto.Vid `json:"vid"`
Bid proto.BlobID `json:"bid"`
}
type SliceRepairMgrConfig struct {
*MessageMgrConfig
Transport base.Transport
BlobNodeSelector selector.Selector
SCClient scheduler.ITaskInfoNotifier
}
// SliceRepairMgr handles repair messages
type SliceRepairMgr struct {
*messageMgr
transport base.Transport
scClient scheduler.ITaskInfoNotifier
blobNodeSelector selector.Selector
chunkMissMigrateReporter *base.AbnormalReporter
}
func NewShardRepairMgr(cfg *MessageMgrConfig) (*SliceRepairMgr, error) {
func NewSliceRepairMgr(cfg *SliceRepairMgrConfig) (*SliceRepairMgr, error) {
if cfg.messageType == 0 {
cfg.messageType = snproto.MessageTypeRepair
}
repairMgr := &SliceRepairMgr{}
cfg.executor = repairMgr
msgMgr, err := newMessageMgr(cfg)
msgMgr, err := newMessageMgr(cfg.MessageMgrConfig)
if err != nil {
return nil, err
}
repairMgr.messageMgr = msgMgr
repairMgr.Start()
return repairMgr, nil
}
@ -49,6 +82,22 @@ func (m *SliceRepairMgr) Start() {
m.messageMgr.run()
}
func (r *SliceRepairMgr) Repair(ctx context.Context, req *snapi.RepairSliceArgs) error {
return r.insertRepairMsg(ctx, req)
}
func (r *SliceRepairMgr) Stats() *snapi.ShardnodeTaskStatsRet {
repairSuccessCounter, repairFailedCounter := r.getTaskStats()
repairErrStats, repairTotalErrCnt := r.getErrorStats()
return &snapi.ShardnodeTaskStatsRet{
Enable: r.enabled(),
SuccessPerMin: fmt.Sprint(repairSuccessCounter),
FailedPerMin: fmt.Sprint(repairFailedCounter),
TotalErrCnt: repairTotalErrCnt,
ErrStats: repairErrStats,
}
}
func (m *SliceRepairMgr) ItemToMessageExt(item interface{}) (snproto.MessageExt, error) {
msgItem, ok := item.(messageItem)
if !ok {
@ -66,10 +115,203 @@ func (m *SliceRepairMgr) ItemToMessageExt(item interface{}) (snproto.MessageExt,
}
// ExecuteWithCheckVolConsistency implements MessageExecutor interface for SliceRepairMgr
// TODO: implement repair logic
func (m *SliceRepairMgr) ExecuteWithCheckVolConsistency(ctx context.Context, vid proto.Vid, ret interface{}) error {
// TODO: implement repair logic
return errors.New("repair logic not implemented yet")
return m.cfg.VolCache.DoubleCheckedRun(ctx, vid, func(info *snproto.VolumeInfoSimple) (newVol *snproto.VolumeInfoSimple, _err error) {
executeRet, ok := ret.(*executeRet)
if !ok {
return nil, errors.New("invalid execute ret")
}
me, ok := executeRet.msgExt.(*repairMsgExt)
if !ok {
return nil, errors.New("not a repair message")
}
return m.tryRepair(executeRet.ctx, info, &me.msg)
})
}
func (m *SliceRepairMgr) tryRepair(ctx context.Context, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) (*snproto.VolumeInfoSimple, error) {
span := trace.SpanFromContextSafe(ctx)
newVol, err := m.repairSlice(ctx, volInfo, repairMsg)
if err == nil {
return newVol, nil
}
if err == ErrBlobnodeServiceUnavailable {
return volInfo, err
}
if isErrDiskNotFound(err) {
m.processDiskNotFoundErr(ctx, volInfo, repairMsg)
}
newVol, err1 := m.cfg.VolCache.UpdateVolume(volInfo.Vid)
if err1 != nil || newVol.EqualWith(volInfo) {
// if update volInfo failed or volInfo not updated, don't need retry
span.Warnf("new volInfo is same or clusterTopology.UpdateVolume failed: vid[%d], vol cache update err[%+v], repair err[%+v]",
volInfo.Vid, err1, err)
return volInfo, err
}
return m.repairSlice(ctx, newVol, repairMsg)
}
func (m *SliceRepairMgr) repairSlice(ctx context.Context, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) (*snproto.VolumeInfoSimple, error) {
span := trace.SpanFromContextSafe(ctx)
span.Infof("repair slice: msg[%+v], vol info[%+v]", repairMsg, volInfo)
hosts := m.blobNodeSelector.GetRandomN(1)
if len(hosts) == 0 {
return volInfo, ErrBlobnodeServiceUnavailable
}
workerHost := hosts[0]
err := m.cfg.BlobTransport.RepairSlice(ctx, workerHost, volInfo, repairMsg)
if err == nil {
return volInfo, nil
}
if isOrphanSlice(err) {
m.saveOrphanSlice(ctx, repairMsg)
}
return volInfo, err
}
func (m *SliceRepairMgr) saveOrphanSlice(ctx context.Context, repairMsg *snproto.SliceRepairMsg) {
span := trace.SpanFromContextSafe(ctx)
slice := OrphanSlice{
ClusterID: m.cfg.ClusterID,
Vid: repairMsg.Vid,
Bid: repairMsg.Bid,
}
span.Infof("save orphan slice: [%+v]", slice)
insistOn(ctx, "save orphan slice", func() error {
return m.executeLogger.Encode(slice)
})
}
func (m *SliceRepairMgr) processDiskNotFoundErr(ctx context.Context, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) {
span := trace.SpanFromContextSafe(ctx)
for _, idx := range repairMsg.BadIdx {
vunitInfo := volInfo.VunitLocations[idx]
if m.chunkMissMigrateReporter.IsVuidReported(vunitInfo.Vuid) {
span.Warnf("chunk is miss migrate and already reported, vunitInfo: %+v", vunitInfo)
continue
}
// check disk status
disk, err := m.transport.GetBlobnodeDiskInfo(ctx, vunitInfo.DiskID)
if err != nil {
span.Errorf("get diskinfo failed, vunitInfo: %+v, err: %s", volInfo, err.Error())
continue
}
// maybe disk is broken but not repaired, and restarted, retry repair next time
if disk.Status <= proto.DiskStatusRepairing {
continue
}
// disk is repaired or dropped, means volInfo is too old, get new volInfo from cm
vol, err := m.transport.GetVolumeInfo(ctx, vunitInfo.Vuid.Vid())
if err != nil {
span.Errorf("get volumeinfo failed, vunitInfo: %+v, err: %s", volInfo, err.Error())
continue
}
if !vol.EqualWith(volInfo) {
continue
}
ret, err := m.scClient.CheckTaskExist(ctx, &scheduler.CheckTaskExistArgs{
TaskType: proto.TaskTypeManualMigrate,
DiskID: vunitInfo.DiskID,
Vuid: vunitInfo.Vuid,
})
if err != nil {
span.Errorf("check task exist failed, vunitInfo: %+v, err: %s", volInfo, err.Error())
continue
}
if ret.Exist {
m.chunkMissMigrateReporter.SetVuidReported(vunitInfo.Vuid)
continue
}
m.chunkMissMigrateReporter.ReportAbnormal(vunitInfo.DiskID, vunitInfo.Vuid)
m.chunkMissMigrateReporter.SetVuidReported(vunitInfo.Vuid)
}
}
func (m *SliceRepairMgr) insertRepairMsg(ctx context.Context, req *snapi.RepairSliceArgs) error {
shard, err := m.getShard(req.Header.DiskID, req.Header.Suid)
if err != nil {
return err
}
span := trace.SpanFromContextSafe(ctx)
ts := m.tsGen.GenerateTs()
shardKeys := req.GetShardKeys(shard.ShardingSubRangeCount())
tier := snproto.TierSingleIdx
if len(req.BadIdx) > 1 {
tier = snproto.TierMultiIdx
}
mk := newMsgKey()
defer mk.release()
mk.msgType = snproto.MessageTypeRepair
mk.tier = tier
mk.ts = ts
mk.vid = req.Vid
mk.bid = req.Bid
mk.shardKeys = shardKeys
key := mk.encode()
msg := snproto.SliceRepairMsg{
Bid: req.Bid,
Vid: req.Vid,
BadIdx: req.BadIdx,
Reason: req.Reason,
ReqId: span.TraceID(),
}
raw, err := msg.Marshal()
if err != nil {
return err
}
itm := snapi.Item{
ID: string(key),
Fields: []snapi.Field{
{
ID: snproto.SliceRepairMsgFieldID,
Value: raw,
},
},
}
oph := storage.OpHeader{
RouteVersion: req.Header.RouteVersion,
ShardKeys: shardKeys,
}
return shard.InsertItem(ctx, oph, []byte(itm.ID), itm)
}
func isOrphanSlice(err error) bool {
return rpc2.DetectStatusCode(err) == apierr.CodeOrphanShard
}
func isErrDiskNotFound(err error) bool {
return rpc2.DetectStatusCode(err) == apierr.CodeDiskNotFound
}
func insistOn(ctx context.Context, errMsg string, on func() error) {
span := trace.SpanFromContextSafe(ctx)
attempt := 0
retry.Insist(time.Second, on, func(err error) {
attempt++
span.Errorf("insist attempt-%d: %s %s", attempt, errMsg, err.Error())
})
}
func itemToRepairMsg(itm snapi.Item) (msg snproto.SliceRepairMsg, err error) {

View File

@ -0,0 +1,511 @@
// Copyright 2025 The CubeFS Authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
// implied. See the License for the specific language governing
// permissions and limitations under the License.
package message
import (
"context"
"os"
"testing"
"github.com/stretchr/testify/require"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/api/scheduler"
snapi "github.com/cubefs/cubefs/blobstore/api/shardnode"
apierr "github.com/cubefs/cubefs/blobstore/common/errors"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/recordlog"
"github.com/cubefs/cubefs/blobstore/common/rpc2"
"github.com/cubefs/cubefs/blobstore/common/taskswitch"
"github.com/cubefs/cubefs/blobstore/shardnode/base"
snproto "github.com/cubefs/cubefs/blobstore/shardnode/proto"
"github.com/cubefs/cubefs/blobstore/shardnode/storage"
"github.com/cubefs/cubefs/blobstore/testing/mocks"
mock "github.com/cubefs/cubefs/blobstore/testing/mockshardnode"
"github.com/cubefs/cubefs/blobstore/util/errors"
"github.com/cubefs/cubefs/blobstore/util/selector"
)
func TestNewSliceRepairMgr(t *testing.T) {
cmClient := mocks.NewMockClientAPI(ctr(t))
cmClient.EXPECT().GetConfig(any, any).Return("", nil).AnyTimes()
taskSwitchMgr := taskswitch.NewSwitchMgr(cmClient)
sh := mock.NewMockSpaceShardHandler(ctr(t))
sg := mock.NewMockMessageMgrShardGetter(ctr(t))
sg.EXPECT().GetAllShards().Return([]storage.ShardHandler{sh}).AnyTimes()
testDir, err := os.MkdirTemp(os.TempDir(), "repair_log")
require.NoError(t, err)
defer os.RemoveAll(testDir)
cfg := &SliceRepairMgrConfig{
MessageMgrConfig: &MessageMgrConfig{
TaskSwitchMgr: taskSwitchMgr,
ShardGetter: sg,
BlobTransport: mocks.NewMockBlobTransport(ctr(t)),
VolCache: mocks.NewMockIVolumeCache(ctr(t)),
MessageCfg: MessageCfg{
FailedMsgChannelSize: 5,
ProduceTaskPoolSize: 4,
RateLimit: 100.0,
RateLimitBurst: 10,
MaxExecuteSliceNum: 1000,
TierConfig: newTestPriorityConfig(),
MessageLog: recordlog.Config{Dir: testDir},
},
},
}
mgr, err := NewSliceRepairMgr(cfg)
require.Nil(t, err)
// test stats
ret := mgr.Stats()
require.NotNil(t, ret)
mgr.Close()
}
func TestSliceRepairMgr_ExecuteWithCheckVolConsistency(t *testing.T) {
// mock blob transport
blobTp := mocks.NewMockBlobTransport(ctr(t))
blobTp.EXPECT().RepairSlice(any, any, any, any).DoAndReturn(func(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) error {
return nil
}).AnyTimes()
// mock blobnode selector
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{"blobnode-1"}).AnyTimes()
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
// mock volume cache transport
volTp := mocks.NewMockVolumeTransport(ctr(t))
volTp.EXPECT().GetVolumeInfo(any, any).Return(volInfo, nil).AnyTimes()
vc := base.NewVolumeCache(volTp, 1)
mgr := newTestSliceRepairMgr(t, nil, blobTp, vc, nil, nil, blobNodeSelector)
bid := proto.BlobID(1)
msg := &repairMsgExt{
msg: snproto.SliceRepairMsg{
Bid: bid,
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
},
}
repairRet := &executeRet{
msgExt: msg,
ctx: context.Background(),
}
// repair
err := mgr.ExecuteWithCheckVolConsistency(context.Background(), vid, repairRet)
require.Nil(t, err)
}
func TestSliceRepairMgr_ExecuteWithCheckVolConsistency_BlobnodeUnavailable(t *testing.T) {
// mock blobnode selector - return empty list
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{}).AnyTimes()
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
// mock volume cache transport
volTp := mocks.NewMockVolumeTransport(ctr(t))
volTp.EXPECT().GetVolumeInfo(any, any).Return(volInfo, nil).AnyTimes()
vc := base.NewVolumeCache(volTp, 1)
mgr := newTestSliceRepairMgr(t, nil, nil, vc, nil, nil, blobNodeSelector)
bid := proto.BlobID(1)
msg := &repairMsgExt{
msg: snproto.SliceRepairMsg{
Bid: bid,
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
},
}
repairRet := &executeRet{
msgExt: msg,
ctx: context.Background(),
}
// repair should fail with ErrBlobnodeServiceUnavailable
err := mgr.ExecuteWithCheckVolConsistency(context.Background(), vid, repairRet)
require.NotNil(t, err)
require.Equal(t, ErrBlobnodeServiceUnavailable, err)
}
func TestSliceRepairMgr_ExecuteWithCheckVolConsistency_OrphanShard(t *testing.T) {
// mock blob transport - return orphan slice error
blobTp := mocks.NewMockBlobTransport(ctr(t))
blobTp.EXPECT().RepairSlice(any, any, any, any).DoAndReturn(func(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) error {
return rpc2.NewError(apierr.CodeOrphanShard, "orphan slice", "")
}).AnyTimes()
// mock blobnode selector
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{"blobnode-1"}).AnyTimes()
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
// mock volume cache transport
volTp := mocks.NewMockVolumeTransport(ctr(t))
volTp.EXPECT().GetVolumeInfo(any, any).Return(volInfo, nil).AnyTimes()
// mock execute logger
executeLogger := mocks.NewMockRecordLogEncoder(ctr(t))
executeLogger.EXPECT().Encode(any).Return(nil).AnyTimes()
vc := base.NewVolumeCache(volTp, 1)
mgr := newTestSliceRepairMgr(t, nil, blobTp, vc, nil, nil, blobNodeSelector)
mgr.executeLogger = executeLogger
bid := proto.BlobID(1)
msg := &repairMsgExt{
msg: snproto.SliceRepairMsg{
Bid: bid,
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
},
}
repairRet := &executeRet{
msgExt: msg,
ctx: context.Background(),
}
mgr.cfg.ClusterID = 1
err := mgr.ExecuteWithCheckVolConsistency(context.Background(), vid, repairRet)
require.NotNil(t, err)
require.Equal(t, apierr.CodeOrphanShard, rpc2.DetectStatusCode(err))
}
func TestSliceRepairMgr_ExecuteWithCheckVolConsistency_DiskNotFound(t *testing.T) {
// mock blob transport - return disk not found error
blobTp := mocks.NewMockBlobTransport(ctr(t))
blobTp.EXPECT().RepairSlice(any, any, any, any).DoAndReturn(func(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) error {
return rpc2.NewError(apierr.CodeDiskNotFound, "disk not found", "")
}).AnyTimes()
// mock blobnode selector
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{"blobnode-1"}).AnyTimes()
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
// mock volume cache transport
volTp := mocks.NewMockVolumeTransport(ctr(t))
volTp.EXPECT().GetVolumeInfo(any, any).Return(volInfo, nil).AnyTimes()
// mock transport for GetBlobnodeDiskInfo and GetVolumeInfo
transport := mocks.NewMockTransport(ctr(t))
transport.EXPECT().GetBlobnodeDiskInfo(any, any).Return(&clustermgr.BlobNodeDiskInfo{
DiskInfo: clustermgr.DiskInfo{
Status: proto.DiskStatusNormal,
},
DiskHeartBeatInfo: clustermgr.DiskHeartBeatInfo{
DiskID: volInfo.VunitLocations[0].DiskID,
},
}, nil).AnyTimes()
transport.EXPECT().GetVolumeInfo(any, any).Return(volInfo, nil).AnyTimes()
// mock scheduler client
scClient := mocks.NewMockIScheduler(ctr(t))
scClient.EXPECT().CheckTaskExist(any, any).Return(&scheduler.CheckTaskExistResp{Exist: false}, nil).AnyTimes()
vc := base.NewVolumeCache(volTp, 1)
mgr := newTestSliceRepairMgr(t, nil, blobTp, vc, transport, scClient, blobNodeSelector)
bid := proto.BlobID(1)
msg := &repairMsgExt{
msg: snproto.SliceRepairMsg{
Bid: bid,
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
},
}
repairRet := &executeRet{
msgExt: msg,
ctx: context.Background(),
}
// repair should fail with disk not found error
err := mgr.ExecuteWithCheckVolConsistency(context.Background(), vid, repairRet)
require.NotNil(t, err)
require.Equal(t, apierr.CodeDiskNotFound, rpc2.DetectStatusCode(err))
}
func TestSliceRepairMgr_ExecuteWithCheckVolConsistency_UpdateVolume(t *testing.T) {
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
newVolInfo := newMockSimpleVolumeInfo(vid)
newVolInfo.VunitLocations[0].DiskID = volInfo.VunitLocations[0].DiskID + 100
// mock volume cache transport
firstTime := true
volTp := mocks.NewMockVolumeTransport(ctr(t))
volTp.EXPECT().GetVolumeInfo(any, any).DoAndReturn(func(ctx context.Context, vid proto.Vid) (*snproto.VolumeInfoSimple, error) {
if firstTime {
firstTime = false
return volInfo, nil
}
return newVolInfo, nil
}).Times(2)
vc := base.NewVolumeCache(volTp, 1)
// mock blob transport - first time fail, second time success
firstTimeRepair := true
blobTp := mocks.NewMockBlobTransport(ctr(t))
blobTp.EXPECT().RepairSlice(any, any, any, any).DoAndReturn(func(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) error {
if firstTimeRepair {
firstTimeRepair = false
return rpc2.NewError(int32(apierr.CodeDiskBroken), "disk broken", "")
}
return nil
}).Times(2)
// mock blobnode selector
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{"blobnode-1"}).AnyTimes()
mgr := newMinimalSliceRepairMgr(&SliceRepairMgrConfig{
MessageMgrConfig: &MessageMgrConfig{
BlobTransport: blobTp,
VolCache: vc,
},
BlobNodeSelector: blobNodeSelector,
})
bid := proto.BlobID(1)
msg := &repairMsgExt{
msg: snproto.SliceRepairMsg{
Bid: bid,
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
},
}
repairRet := &executeRet{
msgExt: msg,
ctx: context.Background(),
}
// repair should succeed after volume update
err := mgr.ExecuteWithCheckVolConsistency(context.Background(), vid, repairRet)
require.Nil(t, err)
}
func TestSliceRepairMgr_ExecuteWithCheckVolConsistency_UpdateVolume2(t *testing.T) {
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
newVolInfo := newMockSimpleVolumeInfo(vid)
newVolInfo.VunitLocations[0].DiskID = volInfo.VunitLocations[0].DiskID + 100
// mock volume cache transport
firstTime := true
volTp := mocks.NewMockVolumeTransport(ctr(t))
volTp.EXPECT().GetVolumeInfo(any, any).DoAndReturn(func(ctx context.Context, vid proto.Vid) (*snproto.VolumeInfoSimple, error) {
if firstTime {
firstTime = false
return volInfo, nil
}
return newVolInfo, nil
}).Times(2)
vc := base.NewVolumeCache(volTp, 1)
// mock blob transport - always fail
blobTp := mocks.NewMockBlobTransport(ctr(t))
blobTp.EXPECT().RepairSlice(any, any, any, any).DoAndReturn(func(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) error {
return rpc2.NewError(int32(apierr.CodeDiskBroken), "disk broken", "")
}).AnyTimes()
// mock blobnode selector
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{"blobnode-1"}).AnyTimes()
mgr := newMinimalSliceRepairMgr(&SliceRepairMgrConfig{
MessageMgrConfig: &MessageMgrConfig{
BlobTransport: blobTp,
VolCache: vc,
},
BlobNodeSelector: blobNodeSelector,
})
bid := proto.BlobID(1)
msg := &repairMsgExt{
msg: snproto.SliceRepairMsg{
Bid: bid,
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
},
}
repairRet := &executeRet{
msgExt: msg,
ctx: context.Background(),
}
// repair should fail even after volume update (volume not changed)
err := mgr.ExecuteWithCheckVolConsistency(context.Background(), vid, repairRet)
require.NotNil(t, err)
}
func TestSliceRepairMgr_RepairShard(t *testing.T) {
// mock blob transport
blobTp := mocks.NewMockBlobTransport(ctr(t))
blobTp.EXPECT().RepairSlice(any, any, any, any).DoAndReturn(func(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) error {
return nil
}).AnyTimes()
// mock blobnode selector
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{"blobnode-1"}).AnyTimes()
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
mgr := newMinimalSliceRepairMgr(&SliceRepairMgrConfig{
MessageMgrConfig: &MessageMgrConfig{
BlobTransport: blobTp,
},
BlobNodeSelector: blobNodeSelector,
})
repairMsg := &snproto.SliceRepairMsg{
Bid: proto.BlobID(1),
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
}
// repair should succeed
newVol, err := mgr.repairSlice(context.Background(), volInfo, repairMsg)
require.Nil(t, err)
require.Equal(t, volInfo, newVol)
}
func TestSliceRepairMgr_RepairShard_Error(t *testing.T) {
// mock blob transport - return error
blobTp := mocks.NewMockBlobTransport(ctr(t))
blobTp.EXPECT().RepairSlice(any, any, any, any).DoAndReturn(func(ctx context.Context, host string, volInfo *snproto.VolumeInfoSimple, repairMsg *snproto.SliceRepairMsg) error {
return errors.New("repair failed")
}).AnyTimes()
// mock blobnode selector
blobNodeSelector := mocks.NewMockSelector(ctr(t))
blobNodeSelector.EXPECT().GetRandomN(1).Return([]string{"blobnode-1"}).AnyTimes()
// mock volume info
vid := proto.Vid(1)
volInfo := newMockSimpleVolumeInfo(vid)
mgr := newMinimalSliceRepairMgr(&SliceRepairMgrConfig{
MessageMgrConfig: &MessageMgrConfig{
BlobTransport: blobTp,
},
BlobNodeSelector: blobNodeSelector,
})
repairMsg := &snproto.SliceRepairMsg{
Bid: proto.BlobID(1),
Vid: vid,
BadIdx: []uint32{0},
Reason: "test",
}
// repair should fail
newVol, err := mgr.repairSlice(context.Background(), volInfo, repairMsg)
require.NotNil(t, err)
require.Equal(t, volInfo, newVol)
}
func TestSliceRepairMgr_InsertRepairMsg(t *testing.T) {
// mock shard
shard := mock.NewMockSpaceShardHandler(ctr(t))
shard.EXPECT().InsertItem(any, any, any, any).Return(nil).AnyTimes()
shard.EXPECT().ShardingSubRangeCount().Return(2).AnyTimes()
// mock shard getter
sg := mock.NewMockMessageMgrShardGetter(ctr(t))
sg.EXPECT().GetShard(any, any).Return(shard, nil).AnyTimes()
mgr := newTestSliceRepairMgr(t, sg, nil, nil, nil, nil, nil)
args := &snapi.RepairSliceArgs{
Header: snapi.ShardOpHeader{
DiskID: 1,
Suid: 1,
},
Bid: proto.BlobID(1),
Vid: proto.Vid(1),
BadIdx: []uint32{0, 1},
Reason: "test",
}
err := mgr.Repair(context.Background(), args)
require.Nil(t, err)
}
// newMinimalSliceRepairMgr creates a minimal SliceRepairMgr for testing
func newMinimalSliceRepairMgr(cfg *SliceRepairMgrConfig) *SliceRepairMgr {
baseMgr := &messageMgr{
cfg: cfg.MessageMgrConfig,
}
mgr := &SliceRepairMgr{
messageMgr: baseMgr,
transport: cfg.Transport,
blobNodeSelector: cfg.BlobNodeSelector,
scClient: cfg.SCClient,
}
baseMgr.cfg.executor = mgr
return mgr
}
func newTestSliceRepairMgr(t *testing.T, sg ShardGetter, tp base.BlobTransport, vc base.IVolumeCache, transport base.Transport, scClient scheduler.IScheduler, selector selector.Selector) *SliceRepairMgr {
baseMgr := newTestBlobMessageMgr(t, sg, tp, vc, nil, snproto.MessageTypeRepair)
mgr := &SliceRepairMgr{
messageMgr: baseMgr,
transport: transport,
scClient: scClient,
blobNodeSelector: selector,
chunkMissMigrateReporter: base.NewAbnormalReporter(0, base.ShardRepair, base.ChunkMissMigrateAbnormal),
}
// Set executor to ShardRepairMgr itself
baseMgr.cfg.executor = mgr
return mgr
}

View File

@ -289,24 +289,23 @@ func TestShardListReader_IsProtected(t *testing.T) {
require.True(t, reader.isProtected(time.Hour))
}
func TestDelMsgExt_IsProtected(t *testing.T) {
func TestMessageExt_IsProtected(t *testing.T) {
// test msg unprotected
oldTime := time.Now().Add(-2 * time.Hour).Unix()
unprotectedMsg := &delMsgExt{
msg := &delMsgExt{
msg: snproto.DeleteMsg{
Time: oldTime,
},
}
require.False(t, unprotectedMsg.IsProtected(time.Hour))
require.False(t, msg.IsProtected(time.Hour))
// test msg protected
now := time.Now().Unix()
protectedMsg := &delMsgExt{
msg: snproto.DeleteMsg{
Time: now,
msg2 := &repairMsgExt{
msg: snproto.SliceRepairMsg{
Time: time.Now().Unix(),
},
}
require.True(t, protectedMsg.IsProtected(time.Hour))
require.True(t, msg2.IsProtected(time.Hour))
}
func TestItemToDelMsg(t *testing.T) {
@ -343,8 +342,56 @@ func TestItemToDelMsg(t *testing.T) {
require.Error(t, err)
}
func TestItemToRepairMsg(t *testing.T) {
// test valid item
msg := snproto.SliceRepairMsg{
Bid: proto.BlobID(100),
Vid: proto.Vid(1),
BadIdx: []uint32{0, 1},
Reason: "test",
ReqId: "test_req_id",
}
item := snapi.Item{
ID: "test_key",
Fields: []snapi.Field{
{
ID: snproto.SliceRepairMsgFieldID,
Value: marshalRepairMsg(msg),
},
},
}
result, err := itemToRepairMsg(item)
require.NoError(t, err)
require.Equal(t, msg.Bid, result.Bid)
require.Equal(t, msg.Vid, result.Vid)
require.Equal(t, msg.BadIdx, result.BadIdx)
require.Equal(t, msg.Reason, result.Reason)
require.Equal(t, msg.ReqId, result.ReqId)
// test invalid item
invalidItem := snapi.Item{
ID: "invalid_key",
Fields: []snapi.Field{
{
ID: snproto.SliceRepairMsgFieldID,
Value: []byte("invalid_data"),
},
},
}
_, err = itemToRepairMsg(invalidItem)
require.Error(t, err)
}
// marshal DeleteMsg to bytes
func marshalDeleteMsg(msg snproto.DeleteMsg) []byte {
data, _ := msg.Marshal()
return data
}
// marshal RepairMsg to bytes
func marshalRepairMsg(msg snproto.SliceRepairMsg) []byte {
data, _ := msg.Marshal()
return data
}

View File

@ -418,6 +418,22 @@ func (s *RpcService) DBStats(w rpc2.ResponseWriter, req *rpc2.Request) error {
return w.WriteOK(&ret)
}
func (s *RpcService) RepairSlice(w rpc2.ResponseWriter, req *rpc2.Request) error {
ctx := req.Context()
span := req.Span()
args := &shardnode.RepairSliceArgs{}
if err := req.ParseParameter(args); err != nil {
return err
}
span.Infof("receive RepairSlice request, args:%+v", args)
return s.repairSlice(ctx, args)
}
func (s *RpcService) RepairSliceStats(w rpc2.ResponseWriter, req *rpc2.Request) error {
ret := s.repairSliceStats()
return w.WriteOK(ret)
}
func initConfig(args []string) (*cmd.Config, error) {
config.Init("f", "", "shardnode.conf")
if err := config.Load(&conf); err != nil {
@ -460,6 +476,9 @@ func newHandler(s *RpcService) *rpc2.Router {
handler.Register("/blob/delete/raw", s.DeleteBlobRaw)
handler.Register("/blob/delete/stats", s.DeleteBlobStats)
handler.Register("/slice/repair", s.RepairSlice)
handler.Register("/slice/repair/stats", s.RepairSliceStats)
return handler
}

View File

@ -145,6 +145,10 @@ func newMockService(t *testing.T, cfg mockServiceCfg) (*service, func(), error)
require.NoError(t, err)
defer os.RemoveAll(delLogDir)
repairLogDir, err := os.MkdirTemp(os.TempDir(), "repair_log")
require.NoError(t, err)
defer os.RemoveAll(repairLogDir)
dm, _ := message.NewBlobDeleteMgr(&message.BlobDelMgrConfig{
TaskSwitchMgr: taskSwitchMgr,
ShardGetter: sg2,
@ -156,6 +160,15 @@ func newMockService(t *testing.T, cfg mockServiceCfg) (*service, func(), error)
})
s.blobDelMgr = dm
shardRepairMgr, _ := message.NewSliceRepairMgr(&message.SliceRepairMgrConfig{
MessageMgrConfig: &message.MessageMgrConfig{
TaskSwitchMgr: taskSwitchMgr,
ShardGetter: sg2,
MessageCfg: message.MessageCfg{MessageLog: recordlog.Config{Dir: repairLogDir}},
},
})
s.sliceRepairMgr = shardRepairMgr
// set disk
s.disks = make(map[proto.DiskID]*storage.Disk)
for id, d := range cfg.disks {
@ -308,7 +321,25 @@ func TestRpcService_Blob(t *testing.T) {
require.Nil(t, err)
// delete blob stats
stats, err := cli.DeleteBlobStats(context.Background(), tcpAddrBlob, shardnode.DeleteBlobStatsArgs{})
stats, err := cli.DeleteBlobStats(context.Background(), tcpAddrBlob, shardnode.ShardnodeTaskStatsArgs{})
require.Nil(t, err)
require.NotNil(t, stats)
// repair shard
err = cli.RepairSlice(context.Background(), tcpAddrBlob, shardnode.RepairSliceArgs{
Header: shardnode.ShardOpHeader{
DiskID: 1,
Suid: 1,
},
Bid: proto.BlobID(1),
Vid: proto.Vid(1),
BadIdx: []uint32{0, 1},
Reason: "test",
})
require.Nil(t, err)
// repair shard stats
stats, err = cli.RepairSliceStats(context.Background(), tcpAddrBlob, shardnode.ShardnodeTaskStatsArgs{})
require.Nil(t, err)
require.NotNil(t, stats)
}

View File

@ -402,6 +402,13 @@ func initServiceConfig(cfg *Config) {
defaulter.LessOrEqual(&cfg.DeleteBlobCfg.MaxExecuteSliceNum, uint64(64))
defaulter.LessOrEqual(&cfg.DeleteBlobCfg.TierConfig.SafeMessageTimeout.Duration, 12*time.Hour)
defaulter.LessOrEqual(&cfg.DeleteBlobCfg.TierConfig.PunishTimeout.Duration, 1*time.Minute)
defaulter.LessOrEqual(&cfg.SliceRepairCfg.RateLimit, float64(1024))
defaulter.LessOrEqual(&cfg.SliceRepairCfg.RateLimitBurst, 64)
defaulter.LessOrEqual(&cfg.SliceRepairCfg.FailedMsgChannelSize, 10<<10)
defaulter.LessOrEqual(&cfg.SliceRepairCfg.ProduceTaskPoolSize, 16)
defaulter.LessOrEqual(&cfg.SliceRepairCfg.MaxExecuteSliceNum, uint64(64))
defaulter.LessOrEqual(&cfg.SliceRepairCfg.TierConfig.SafeMessageTimeout.Duration, 12*time.Hour)
defaulter.LessOrEqual(&cfg.SliceRepairCfg.TierConfig.PunishTimeout.Duration, 1*time.Minute)
}
func isDiskInfoMatch(a, b clustermgr.ShardNodeDiskInfo) bool {

View File

@ -24,6 +24,7 @@ import (
bnapi "github.com/cubefs/cubefs/blobstore/api/blobnode"
cmapi "github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/api/scheduler"
shardnodeapi "github.com/cubefs/cubefs/blobstore/api/shardnode"
"github.com/cubefs/cubefs/blobstore/cmd"
apierr "github.com/cubefs/cubefs/blobstore/common/errors"
@ -40,6 +41,7 @@ import (
"github.com/cubefs/cubefs/blobstore/shardnode/storage"
"github.com/cubefs/cubefs/blobstore/shardnode/storage/store"
"github.com/cubefs/cubefs/blobstore/util/closer"
"github.com/cubefs/cubefs/blobstore/util/selector"
"github.com/cubefs/cubefs/blobstore/util/taskpool"
)
@ -79,7 +81,9 @@ type Config struct {
WaitReOpenDiskIntervalS int64 `json:"wait_re_open_disk_interval_s"`
ShardCheckAndClearIntervalH int64 `json:"shard_check_and_clear_interval_h"`
DeleteBlobCfg message.MessageCfg `json:"blob_delete_cfg"`
DeleteBlobCfg message.MessageCfg `json:"blob_delete_cfg"`
ScClientConfig scheduler.Config `json:"sc_client_config"`
SliceRepairCfg message.MessageCfg `json:"slice_repair_cfg"`
}
// newService returns the singleton service instance
@ -155,19 +159,39 @@ func createService(cfg *Config) *service {
})
svr.catalog = c
cfg.DeleteBlobCfg.ClusterID = cfg.NodeConfig.ClusterID
taskSwitchMgr := taskswitch.NewSwitchMgr(cmClient)
dm, err := message.NewBlobDeleteMgr(&message.BlobDelMgrConfig{
cfg.DeleteBlobCfg.ClusterID = cfg.NodeConfig.ClusterID
blboDeleteMgr, err := message.NewBlobDeleteMgr(&message.BlobDelMgrConfig{
TaskSwitchMgr: taskSwitchMgr,
ShardGetter: svr,
Transport: transport,
BlobTransport: transport,
VolCache: base.NewVolumeCache(transport, 10*time.Second),
MessageCfg: cfg.DeleteBlobCfg,
})
if err != nil {
span.Fatalf("new blob delete mgr failed, err: %s", err.Error())
}
svr.blobDelMgr = dm
svr.blobDelMgr = blboDeleteMgr
cfg.SliceRepairCfg.ClusterID = cfg.NodeConfig.ClusterID
sliceRepairMgr, err := message.NewSliceRepairMgr(&message.SliceRepairMgrConfig{
MessageMgrConfig: &message.MessageMgrConfig{
TaskSwitchMgr: taskSwitchMgr,
ShardGetter: svr,
BlobTransport: transport,
VolCache: base.NewVolumeCache(transport, 10*time.Second),
MessageCfg: cfg.SliceRepairCfg,
},
Transport: transport,
BlobNodeSelector: selector.MakeSelector(60*1000, func() (hosts []string, err error) {
return transport.GetService(context.Background(), proto.ServiceNameWorker)
}),
SCClient: scheduler.New(&cfg.ScClientConfig, cmClient, cfg.NodeConfig.ClusterID),
})
if err != nil {
span.Fatalf("new slice repair mgr failed, err: %s", err.Error())
}
svr.sliceRepairMgr = sliceRepairMgr
go svr.loop(ctx)
span.Infof("service started success")
@ -182,7 +206,8 @@ type service struct {
taskPool taskpool.TaskPool
groupRun singleflight.Group
blobDelMgr *message.BlobDeleteMgr
blobDelMgr *message.BlobDeleteMgr
sliceRepairMgr *message.SliceRepairMgr
cfg Config
lock sync.RWMutex
@ -219,5 +244,6 @@ func (s *service) getAllDisks() []*storage.Disk {
func (s *service) close() {
s.closer.Close()
s.blobDelMgr.Close()
s.sliceRepairMgr.Close()
s.cfg.RaftConfig.Transport.Close()
}

View File

@ -82,6 +82,11 @@ func TestSvr_Loop(t *testing.T) {
defer os.RemoveAll(delLogDir)
cfg.DeleteBlobCfg.MessageLog.Dir = delLogDir
repairLogDir, err := os.MkdirTemp(os.TempDir(), "repair_log")
require.NoError(t, err)
defer os.RemoveAll(repairLogDir)
cfg.SliceRepairCfg.MessageLog.Dir = repairLogDir
s := newService(cfg)
time.Sleep(3 * time.Second)
s.closer.Close()
@ -99,6 +104,11 @@ func TestSvr_HandleEIO(t *testing.T) {
defer os.RemoveAll(delLogDir)
cfg.DeleteBlobCfg.MessageLog.Dir = delLogDir
repairLogDir, err := os.MkdirTemp(os.TempDir(), "repair_log")
require.NoError(t, err)
defer os.RemoveAll(repairLogDir)
cfg.SliceRepairCfg.MessageLog.Dir = repairLogDir
s := newService(cfg)
disk := &storage.Disk{}
disk.SetDiskInfo(cmapi.ShardNodeDiskInfo{
@ -128,6 +138,11 @@ func TestSingletonPattern(t *testing.T) {
defer os.RemoveAll(delLogDir)
cfg.DeleteBlobCfg.MessageLog.Dir = delLogDir
repairLogDir, err := os.MkdirTemp(os.TempDir(), "repair_log")
require.NoError(t, err)
defer os.RemoveAll(repairLogDir)
cfg.SliceRepairCfg.MessageLog.Dir = repairLogDir
// First call to newService
service1 := newService(cfg)
require.NotNil(t, service1)

View File

@ -113,6 +113,21 @@ func (mr *MockTransportMockRecorder) GetAllSpaces(ctx interface{}) *gomock.Call
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetAllSpaces", reflect.TypeOf((*MockTransport)(nil).GetAllSpaces), ctx)
}
// GetBlobnodeDiskInfo mocks base method.
func (m *MockTransport) GetBlobnodeDiskInfo(ctx context.Context, diskID proto.DiskID) (*clustermgr.BlobNodeDiskInfo, error) {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "GetBlobnodeDiskInfo", ctx, diskID)
ret0, _ := ret[0].(*clustermgr.BlobNodeDiskInfo)
ret1, _ := ret[1].(error)
return ret0, ret1
}
// GetBlobnodeDiskInfo indicates an expected call of GetBlobnodeDiskInfo.
func (mr *MockTransportMockRecorder) GetBlobnodeDiskInfo(ctx, diskID interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetBlobnodeDiskInfo", reflect.TypeOf((*MockTransport)(nil).GetBlobnodeDiskInfo), ctx, diskID)
}
// GetConfig mocks base method.
func (m *MockTransport) GetConfig(ctx context.Context, key string) (string, error) {
m.ctrl.T.Helper()
@ -188,6 +203,21 @@ func (mr *MockTransportMockRecorder) GetRouteUpdate(ctx, routeVersion interface{
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetRouteUpdate", reflect.TypeOf((*MockTransport)(nil).GetRouteUpdate), ctx, routeVersion)
}
// GetService mocks base method.
func (m *MockTransport) GetService(ctx context.Context, name string) ([]string, error) {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "GetService", ctx, name)
ret0, _ := ret[0].([]string)
ret1, _ := ret[1].(error)
return ret0, ret1
}
// GetService indicates an expected call of GetService.
func (mr *MockTransportMockRecorder) GetService(ctx, name interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetService", reflect.TypeOf((*MockTransport)(nil).GetService), ctx, name)
}
// GetSpace mocks base method.
func (m *MockTransport) GetSpace(ctx context.Context, sid proto.SpaceID) (*clustermgr.Space, error) {
m.ctrl.T.Helper()
@ -334,6 +364,20 @@ func (mr *MockTransportMockRecorder) RegisterDisk(ctx, disk interface{}) *gomock
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RegisterDisk", reflect.TypeOf((*MockTransport)(nil).RegisterDisk), ctx, disk)
}
// RepairSlice mocks base method.
func (m *MockTransport) RepairSlice(ctx context.Context, host string, volInfo *proto0.VolumeInfoSimple, repairMsg *proto0.SliceRepairMsg) error {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "RepairSlice", ctx, host, volInfo, repairMsg)
ret0, _ := ret[0].(error)
return ret0
}
// RepairSlice indicates an expected call of RepairSlice.
func (mr *MockTransportMockRecorder) RepairSlice(ctx, host, volInfo, repairMsg interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RepairSlice", reflect.TypeOf((*MockTransport)(nil).RepairSlice), ctx, host, volInfo, repairMsg)
}
// ResolveNodeAddr mocks base method.
func (m *MockTransport) ResolveNodeAddr(ctx context.Context, diskID proto.DiskID) (string, error) {
m.ctrl.T.Helper()
@ -873,6 +917,20 @@ func (mr *MockBlobTransportMockRecorder) MarkDeleteSliceUnit(ctx, info, bid inte
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "MarkDeleteSliceUnit", reflect.TypeOf((*MockBlobTransport)(nil).MarkDeleteSliceUnit), ctx, info, bid)
}
// RepairSlice mocks base method.
func (m *MockBlobTransport) RepairSlice(ctx context.Context, host string, volInfo *proto0.VolumeInfoSimple, repairMsg *proto0.SliceRepairMsg) error {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "RepairSlice", ctx, host, volInfo, repairMsg)
ret0, _ := ret[0].(error)
return ret0
}
// RepairSlice indicates an expected call of RepairSlice.
func (mr *MockBlobTransportMockRecorder) RepairSlice(ctx, host, volInfo, repairMsg interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RepairSlice", reflect.TypeOf((*MockBlobTransport)(nil).RepairSlice), ctx, host, volInfo, repairMsg)
}
// MockVolumeTransport is a mock of VolumeTransport interface.
type MockVolumeTransport struct {
ctrl *gomock.Controller