feat(clustermgr): create and update volume route

with #1000431984

Signed-off-by: tangdeyi <tangdeyi@oppo.com>
This commit is contained in:
tangdeyi 2025-10-17 16:53:27 +08:00 committed by 梁曟風
parent 2b2efeb073
commit be22aad2b7
36 changed files with 3086 additions and 428 deletions

View File

@ -71,6 +71,7 @@ type VolumeInfoBase struct {
Used uint64 `json:"used"`
CreateByNodeID uint64 `json:"create_by_node_id"`
Epoch uint32 `json:"epoch"`
RouteVersion proto.RouteVersion `json:"route_version"`
}
type AllocVolumeInfo struct {
@ -206,15 +207,8 @@ type ListVolumeUnitArgs struct {
DiskID proto.DiskID `json:"disk_id"`
}
type VolumeUnitInfo struct {
Vuid proto.Vuid `json:"vuid"`
DiskID proto.DiskID `json:"disk_id"`
Total uint64 `json:"total"`
Free uint64 `json:"free"`
Used uint64 `json:"used"`
Compacting bool `json:"compact"`
Host string `json:"host"`
}
// VolumeUnitInfoBase is used to maintain the VolumeUnitInfo marshal method after existing structs are defined in pb.
type VolumeUnitInfo VolumeUnitInfoBase
type ListVolumeUnitInfos struct {
VolumeUnitInfos []*VolumeUnitInfo `json:"volume_unit_infos"`
@ -364,3 +358,9 @@ type AdminUpdateUnitArgs struct {
NextEpoch uint32 `json:"next_epoch"`
VolumeUnitInfo
}
func (c *Client) GetVolumeRoutes(ctx context.Context, args *GetVolumeRoutesArgs) (ret *GetVolumeRoutesRet, err error) {
ret = &GetVolumeRoutesRet{}
err = c.GetWith(ctx, fmt.Sprintf("/volumeroutes/get?route_version=%d", args.RouteVersion), ret)
return
}

File diff suppressed because it is too large Load Diff

View File

@ -0,0 +1,55 @@
syntax = "proto3";
package cubefs.blobstore.api.clustermgr;
option go_package = "./;clustermgr";
option (gogoproto.sizer_all) = true;
option (gogoproto.marshaler_all) = true;
option (gogoproto.unmarshaler_all) = true;
option (gogoproto.goproto_unkeyed_all) = true;
option (gogoproto.goproto_unrecognized_all) = true;
option (gogoproto.goproto_sizecache_all) = true;
option (gogoproto.goproto_stringer_all) = false;
option (gogoproto.stringer_all) = true;
option (gogoproto.gostring_all) = true;
import "gogoproto/gogo.proto";
import "google/protobuf/any.proto";
message VolumeUnitInfoBase {
uint64 vuid = 1 [(gogoproto.customname) = "Vuid", (gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.Vuid"];
uint32 disk_id = 2 [(gogoproto.customname) = "DiskID", (gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.DiskID"];
uint64 total = 3;
uint64 free = 4;
uint64 used = 5;
bool compact = 6 [(gogoproto.customname) = "Compacting"];
string host = 7;
}
message RouteItemAddVolume {
uint32 vid = 1 [(gogoproto.customname) = "Vid", (gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.Vid"];
uint64 route_version = 2 [(gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.RouteVersion"];
repeated VolumeUnitInfoBase units = 3 [(gogoproto.nullable) = false];
}
message RouteItemUpdateVolume {
uint32 vid = 1 [(gogoproto.customname) = "Vid", (gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.Vid"];
uint64 route_version = 2 [(gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.RouteVersion"];
VolumeUnitInfoBase unit = 3 [(gogoproto.nullable) = false];
}
message VolumeRouteItem {
uint64 route_version = 1 [(gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.RouteVersion"];
uint32 type = 2 [(gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.VolumeRouteItemType"];
google.protobuf.Any item =3;
}
message GetVolumeRoutesArgs {
uint64 route_version = 1 [(gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.RouteVersion"];
}
message GetVolumeRoutesRet {
uint64 route_version = 1 [(gogoproto.casttype) = "github.com/cubefs/cubefs/blobstore/common/proto.RouteVersion"];
repeated VolumeRouteItem items = 2 [(gogoproto.nullable) = false];
}

View File

@ -0,0 +1,230 @@
package base
import (
"context"
"fmt"
"sync"
"sync/atomic"
"time"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/util/errors"
)
var RemoveOldRouteInternal = 1 * time.Minute
type RouteStorage interface {
GetFirstRoute() (*RouteInfoRecord, error)
ListRoute() ([]*RouteInfoRecord, error)
DeleteOldRoutes(before proto.RouteVersion) error
}
type RouteMgr struct {
truncateIntervalNum uint32
initNullRoute bool
unstableRouteVersion proto.RouteVersion
stableRouteVersion proto.RouteVersion
increments *routeItemRing
done chan struct{}
lock sync.RWMutex
recordToItem func(info *RouteInfoRecord) *RouteItem
storage RouteStorage
}
func NewRouteMgr(truncateIntervalNum uint32, initNullRoute bool, recordToItem func(info *RouteInfoRecord) *RouteItem, storage RouteStorage) *RouteMgr {
r := &RouteMgr{
truncateIntervalNum: truncateIntervalNum,
initNullRoute: initNullRoute,
increments: newRouteItemRing(truncateIntervalNum),
done: make(chan struct{}),
recordToItem: recordToItem,
storage: storage,
}
return r
}
func (r *RouteMgr) Close() {
close(r.done)
}
func (r *RouteMgr) LoadRoute(ctx context.Context) error {
// load route into memory
records, err := r.storage.ListRoute()
if err != nil {
return errors.Info(err, "storage ListRoute").Detail(err)
}
if len(records) > int(r.truncateIntervalNum) {
records = records[len(records)-int(r.truncateIntervalNum):]
}
maxRouteVersion := r.stableRouteVersion
for _, record := range records {
item := r.recordToItem(record)
r.increments.put(item)
if item.RouteVersion > maxRouteVersion {
maxRouteVersion = item.RouteVersion
}
}
r.stableRouteVersion = maxRouteVersion
r.unstableRouteVersion = maxRouteVersion
// init null route
if r.initNullRoute && r.stableRouteVersion == 0 {
maxRouteVersion = proto.RouteVersion(1)
r.increments.put(&RouteItem{RouteVersion: maxRouteVersion})
r.stableRouteVersion = maxRouteVersion
r.unstableRouteVersion = maxRouteVersion
}
return nil
}
func (r *RouteMgr) GetRouteVersion() uint64 {
return atomic.LoadUint64((*uint64)(&r.stableRouteVersion))
}
func (r *RouteMgr) GenRouteVersion(ctx context.Context, step uint64) uint64 {
return atomic.AddUint64((*uint64)(&r.unstableRouteVersion), step)
}
func (r *RouteMgr) InsertRouteItems(ctx context.Context, items []*RouteItem) {
r.lock.Lock()
defer r.lock.Unlock()
maxStableRouteVersion := proto.RouteVersion(0)
for _, item := range items {
r.increments.put(item)
if item.RouteVersion > maxStableRouteVersion {
maxStableRouteVersion = item.RouteVersion
}
}
atomic.StoreUint64((*uint64)(&r.stableRouteVersion), uint64(maxStableRouteVersion))
}
func (r *RouteMgr) GetRouteItems(ctx context.Context, ver proto.RouteVersion) (ret []*RouteItem, isLatest bool) {
r.lock.RLock()
defer r.lock.RUnlock()
return r.increments.getFrom(ver)
}
func (r *RouteMgr) Loop() {
_, ctx := trace.StartSpanFromContext(context.Background(), "")
ticker := time.NewTicker(RemoveOldRouteInternal)
for {
select {
case <-ticker.C:
// check route items num, remove old route item if exceed the max increment items limit
r.removeOldRouteItems(ctx)
case <-r.done:
return
}
}
}
func (r *RouteMgr) removeOldRouteItems(ctx context.Context) error {
span := trace.SpanFromContextSafe(ctx)
item, err := r.storage.GetFirstRoute()
if err != nil {
span.Errorf("get first route item failed: %s", err.Error())
return fmt.Errorf("get first route item failed: %s", err.Error())
}
if item == nil {
span.Info("routeTbl has no first route")
return nil
}
stableRouteVersion := atomic.LoadUint64((*uint64)(&r.stableRouteVersion))
if uint64(item.RouteVersion) < stableRouteVersion-uint64(r.truncateIntervalNum) {
if err := r.storage.DeleteOldRoutes(proto.RouteVersion(stableRouteVersion-uint64(r.truncateIntervalNum)) + 1); err != nil {
span.Errorf("delete oldest route items failed: %s", err.Error())
return fmt.Errorf("delete oldest route items failed: %s", err.Error())
}
span.Infof("delete oldest route items[%d] success", item.RouteVersion)
}
return nil
}
// persistent record
type RouteInfoRecord struct {
RouteVersion proto.RouteVersion `json:"route_version"`
Type interface{} `json:"type"`
ItemDetail interface{} `json:"item"`
}
// memory item
type RouteItem struct {
RouteVersion proto.RouteVersion
Type interface{}
ItemDetail interface{}
}
type routeItemRing struct {
data []*RouteItem
head uint32
tail uint32
nextTail uint32
cap uint32
usedCap uint32
}
func newRouteItemRing(cap uint32) *routeItemRing {
ring := &routeItemRing{
data: make([]*RouteItem, cap),
cap: cap,
}
return ring
}
func (r *routeItemRing) put(item *RouteItem) {
r.data[r.nextTail] = item
r.tail = r.nextTail
if r.cap == r.usedCap {
r.head++
r.head = r.head % r.cap
r.nextTail++
r.nextTail = r.nextTail % r.cap
} else {
r.nextTail++
r.nextTail = r.nextTail % r.cap
r.usedCap++
}
if (r.data[r.head].RouteVersion + proto.RouteVersion(r.usedCap)) != (r.data[r.tail].RouteVersion + 1) {
errMsg := fmt.Sprintf("route cache ring is not consistently, head %v ver: %v, usedCap: %v, tail %v ver: %v",
r.head, r.data[r.head].RouteVersion, r.usedCap, r.tail, r.data[r.tail].RouteVersion)
panic(errMsg)
}
}
func (r *routeItemRing) getFrom(ver proto.RouteVersion) (ret []*RouteItem, isLatest bool) {
if r.head == r.tail {
return nil, true
}
if r.getMinVer() > ver {
return nil, false
}
if r.getMaxVer() <= ver {
return nil, true
}
headVer := r.data[r.head].RouteVersion
i := (r.head + uint32(ver-headVer+1)) % r.cap
for j := 0; j < int(r.usedCap); j++ {
ret = append(ret, r.data[i])
i = (i + 1) % r.cap
if i == r.nextTail {
break
}
}
return ret, false
}
func (r *routeItemRing) getMinVer() proto.RouteVersion {
return r.data[r.head].RouteVersion
}
func (r *routeItemRing) getMaxVer() proto.RouteVersion {
return r.data[r.tail].RouteVersion
}

View File

@ -0,0 +1,40 @@
package base
import (
"testing"
"github.com/stretchr/testify/assert"
"github.com/cubefs/cubefs/blobstore/common/proto"
)
func TestRouteItemRing(t *testing.T) {
ring := newRouteItemRing(3)
items, isLatest := ring.getFrom(3)
assert.Equal(t, 0, len(items))
assert.Equal(t, true, isLatest)
for i := 1; i <= 3; i++ {
item := &RouteItem{
RouteVersion: proto.RouteVersion(i),
}
ring.put(item)
}
assert.Equal(t, proto.RouteVersion(1), ring.getMinVer())
assert.Equal(t, proto.RouteVersion(3), ring.getMaxVer())
items, isLatest = ring.getFrom(1)
assert.Equal(t, 2, len(items))
assert.Equal(t, false, isLatest)
items, isLatest = ring.getFrom(3)
assert.Equal(t, 0, len(items))
assert.Equal(t, true, isLatest)
item4 := &RouteItem{
RouteVersion: proto.RouteVersion(4),
}
ring.put(item4)
assert.Equal(t, ring.getMinVer(), proto.RouteVersion(2))
assert.Equal(t, ring.getMaxVer(), proto.RouteVersion(4))
}

View File

@ -244,7 +244,7 @@ func TestCatalogMgr_LoadData(t *testing.T) {
require.NoError(t, err)
require.Equal(t, 0, mockCatalogMgr.allShards.getShardNum())
require.Nil(t, mockCatalogMgr.allSpaces.getSpaceByID(1))
require.Equal(t, uint64(0), mockCatalogMgr.routeMgr.getRouteVersion())
require.Equal(t, uint64(0), mockCatalogMgr.routeMgr.GetRouteVersion())
// mock apply snapshot put data
err = generateShard(catalogDB)
@ -255,5 +255,5 @@ func TestCatalogMgr_LoadData(t *testing.T) {
require.NoError(t, err)
require.NotEqual(t, 0, mockCatalogMgr.allShards.getShardNum())
require.NotNil(t, mockCatalogMgr.allSpaces.getSpaceByID(1))
require.NotEqual(t, uint64(0), mockCatalogMgr.routeMgr.getRouteVersion())
require.NotEqual(t, uint64(0), mockCatalogMgr.routeMgr.GetRouteVersion())
}

View File

@ -74,7 +74,7 @@ type CatalogMgr struct {
raftServer raftserver.RaftServer
scopeMgr scopemgr.ScopeMgrAPI
kvMgr kvmgr.KvMgrAPI
routeMgr *routeMgr
routeMgr *base.RouteMgr
diskMgr cluster.ShardNodeManagerAPI
shardNodeClient cluster.ShardNodeAPI
@ -124,5 +124,5 @@ func (c *CatalogMgr) loop() {
}
func (c *CatalogMgr) routeLoop() {
c.routeMgr.loop()
c.routeMgr.Loop()
}

View File

@ -156,7 +156,7 @@ func generateShard(catalogDB *catalogdb.CatalogDB) error {
unitCount = 3
shards []*catalogdb.ShardInfoRecord
units []*catalogdb.ShardUnitInfoRecord
routes []*catalogdb.RouteInfoRecord
routes []*base.RouteInfoRecord
ranges = sharding.InitShardingRange(sharding.RangeType_RangeTypeHash, 2, 10)
)
@ -185,7 +185,7 @@ func generateShard(catalogDB *catalogdb.CatalogDB) error {
RouteVersion: proto.RouteVersion(i),
}
route := &catalogdb.RouteInfoRecord{
route := &base.RouteInfoRecord{
RouteVersion: proto.RouteVersion(i),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &catalogdb.RouteInfoShardAdd{ShardID: proto.ShardID(i)},
@ -214,7 +214,7 @@ func generateShard(catalogDB *catalogdb.CatalogDB) error {
}
}
route := &catalogdb.RouteInfoRecord{
route := &base.RouteInfoRecord{
RouteVersion: proto.RouteVersion(11),
Type: proto.CatalogChangeItemUpdateShard,
ItemDetail: &catalogdb.RouteInfoShardUpdate{SuidPrefix: proto.EncodeSuidPrefix(1, 1)},

View File

@ -208,16 +208,6 @@ func (c *CatalogMgr) applyCreateShard(ctx context.Context, shard *shardItem) err
return nil
}
// insert route item
routeVersion := c.routeMgr.genRouteVersion(ctx, 1)
route := &routeItem{
RouteVersion: proto.RouteVersion(routeVersion),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &routeItemShardAdd{ShardID: shard.shardID},
}
c.routeMgr.insertRouteItems(ctx, []*routeItem{route})
shard.info.RouteVersion = proto.RouteVersion(routeVersion)
shardRecord := shard.toShardRecord()
unitRecords := shardUnitsToShardUnitRecords(shard.info.Units, shard.unitEpochs)
// delete transited table firstly, put shard and units secondly.
@ -226,13 +216,22 @@ func (c *CatalogMgr) applyCreateShard(ctx context.Context, shard *shardItem) err
return errors.Info(err, fmt.Sprintf("delete shard [%d] from transitedTbl failed", shard.shardID)).Detail(err)
}
// insert route item
routeVersion := c.routeMgr.GenRouteVersion(ctx, 1)
route := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeVersion),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &routeItemShardAdd{ShardID: shard.shardID},
}
shardRecords := []*catalogdb.ShardInfoRecord{shardRecord}
routeRecords := []*catalogdb.RouteInfoRecord{routeItemToRouteRecord(route)}
routeRecords := []*base.RouteInfoRecord{routeItemToRouteRecord(route)}
if err := c.catalogTbl.PutShardsAndUnitsAndRouteItems(shardRecords, unitRecords, routeRecords); err != nil {
return errors.Info(err, fmt.Sprintf("put shard[%+v], units[%+v] and route[%+v] into catalogTbl failed",
shardRecord, unitRecords, routeRecords)).Detail(err)
}
c.allShards.putShard(shard)
c.routeMgr.InsertRouteItems(ctx, []*base.RouteItem{route})
shard.info.RouteVersion = proto.RouteVersion(routeVersion)
if c.allShards.getShardNum() == c.InitShardNum {
err := c.kvMgr.Set(proto.ShardInitDoneKey, []byte("1"))
@ -273,7 +272,7 @@ func (c *CatalogMgr) allocShardForAllUnits(ctx context.Context, shardCtx *create
DiskType: proto.DiskTypeNVMeSSD,
Suids: suids,
Range: shardCtx.ShardInfo.Range,
RouteVersion: proto.RouteVersion(c.routeMgr.getRouteVersion()),
RouteVersion: proto.RouteVersion(c.routeMgr.GetRouteVersion()),
}
for i := 0; i < IncreaseEpochInterval; i++ {

View File

@ -18,6 +18,7 @@ import (
"sync"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/clustermgr/persistence/catalogdb"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/sharding"
@ -302,12 +303,6 @@ func shardUnitToShardUnitInfo(unit clustermgr.ShardUnit, route proto.RouteVersio
}
}
type routeItem struct {
RouteVersion proto.RouteVersion
Type proto.CatalogChangeItemType
ItemDetail interface{}
}
type routeItemShardAdd struct {
ShardID proto.ShardID
}
@ -316,8 +311,8 @@ type routeItemShardUpdate struct {
SuidPrefix proto.SuidPrefix
}
func routeRecordToRouteItem(info *catalogdb.RouteInfoRecord) *routeItem {
item := &routeItem{
func routeRecordToRouteItem(info *base.RouteInfoRecord) *base.RouteItem {
item := &base.RouteItem{
RouteVersion: info.RouteVersion,
Type: info.Type,
}
@ -334,8 +329,8 @@ func routeRecordToRouteItem(info *catalogdb.RouteInfoRecord) *routeItem {
return item
}
func routeItemToRouteRecord(item *routeItem) *catalogdb.RouteInfoRecord {
record := &catalogdb.RouteInfoRecord{
func routeItemToRouteRecord(item *base.RouteItem) *base.RouteInfoRecord {
record := &base.RouteInfoRecord{
RouteVersion: item.RouteVersion,
Type: item.Type,
}

View File

@ -16,16 +16,11 @@ package catalog
import (
"context"
"fmt"
"sync"
"sync/atomic"
"time"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/persistence/catalogdb"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/util/errors"
"github.com/gogo/protobuf/types"
)
@ -34,19 +29,20 @@ func (c *CatalogMgr) GetCatalogChanges(ctx context.Context, args *clustermgr.Get
span := trace.SpanFromContextSafe(ctx)
var (
items []*routeItem
items []*base.RouteItem
isLatest bool
)
if args.RouteVersion > 0 {
items, isLatest = c.routeMgr.getRouteItems(ctx, args.RouteVersion)
items, isLatest = c.routeMgr.GetRouteItems(ctx, args.RouteVersion)
}
ret = new(clustermgr.GetCatalogChangesRet)
if items == nil && !isLatest {
// get all catalog
ret.RouteVersion = proto.RouteVersion(c.routeMgr.getRouteVersion())
ret.RouteVersion = proto.RouteVersion(c.routeMgr.GetRouteVersion())
shards := c.allShards.list()
items = make([]*base.RouteItem, 0, len(shards))
for _, shard := range shards {
items = append(items, &routeItem{
items = append(items, &base.RouteItem{
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &routeItemShardAdd{ShardID: shard.shardID},
})
@ -56,7 +52,7 @@ func (c *CatalogMgr) GetCatalogChanges(ctx context.Context, args *clustermgr.Get
for i := range items {
ret.Items = append(ret.Items, clustermgr.CatalogChangeItem{
RouteVersion: items[i].RouteVersion,
Type: items[i].Type,
Type: items[i].Type.(proto.CatalogChangeItemType),
})
switch items[i].Type {
case proto.CatalogChangeItemAddShard:
@ -106,186 +102,3 @@ func (c *CatalogMgr) GetCatalogChanges(ctx context.Context, args *clustermgr.Get
}
return ret, nil
}
type routeMgr struct {
truncateIntervalNum uint32
unstableRouteVersion proto.RouteVersion
stableRouteVersion proto.RouteVersion
increments *routeItemRing
done chan struct{}
lock sync.RWMutex
storage *catalogdb.CatalogTable
}
func newRouteMgr(truncateIntervalNum uint32, storage *catalogdb.CatalogTable) *routeMgr {
r := &routeMgr{
truncateIntervalNum: truncateIntervalNum,
increments: newRouteItemRing(truncateIntervalNum),
done: make(chan struct{}),
storage: storage,
}
return r
}
func (r *routeMgr) Close() {
close(r.done)
}
func (r *routeMgr) loadRoute(ctx context.Context) error {
// load route into memory
records, err := r.storage.ListRoute()
if err != nil {
return errors.Info(err, "catalogTbl ListRoute").Detail(err)
}
if len(records) > int(r.truncateIntervalNum) {
records = records[len(records)-int(r.truncateIntervalNum):]
}
maxRouteVersion := proto.RouteVersion(0)
for _, record := range records {
item := routeRecordToRouteItem(record)
r.increments.put(item)
if item.RouteVersion > maxRouteVersion {
maxRouteVersion = item.RouteVersion
}
}
r.stableRouteVersion = maxRouteVersion
r.unstableRouteVersion = maxRouteVersion
return nil
}
func (r *routeMgr) getRouteVersion() uint64 {
return atomic.LoadUint64((*uint64)(&r.stableRouteVersion))
}
func (r *routeMgr) genRouteVersion(ctx context.Context, step uint64) uint64 {
return atomic.AddUint64((*uint64)(&r.unstableRouteVersion), step)
}
func (r *routeMgr) insertRouteItems(ctx context.Context, items []*routeItem) {
r.lock.Lock()
defer r.lock.Unlock()
maxStableRouteVersion := proto.RouteVersion(0)
for _, item := range items {
r.increments.put(item)
if item.RouteVersion > maxStableRouteVersion {
maxStableRouteVersion = item.RouteVersion
}
}
atomic.StoreUint64((*uint64)(&r.stableRouteVersion), uint64(maxStableRouteVersion))
}
func (r *routeMgr) getRouteItems(ctx context.Context, ver proto.RouteVersion) (ret []*routeItem, isLatest bool) {
r.lock.RLock()
defer r.lock.RUnlock()
return r.increments.getFrom(ver)
}
func (r *routeMgr) loop() {
_, ctx := trace.StartSpanFromContext(context.Background(), "")
ticker := time.NewTicker(1 * time.Minute)
for {
select {
case <-ticker.C:
// check route items num, remove old route item if exceed the max increment items limit
r.removeOldRouteItems(ctx)
case <-r.done:
return
}
}
}
func (r *routeMgr) removeOldRouteItems(ctx context.Context) error {
span := trace.SpanFromContextSafe(ctx)
item, err := r.storage.GetFirstRouteItem()
if err != nil {
span.Errorf("get first route item failed: %s", err.Error())
return fmt.Errorf("get first route item failed: %s", err.Error())
}
if item == nil {
span.Info("routeTbl has no first route")
return nil
}
if uint64(item.RouteVersion) < atomic.LoadUint64((*uint64)(&r.stableRouteVersion))-uint64(r.truncateIntervalNum) {
if err := r.storage.DeleteOldRoutes(item.RouteVersion); err != nil {
span.Errorf("delete oldest route items failed: %s", err.Error())
return fmt.Errorf("delete oldest route items failed: %s", err.Error())
}
span.Infof("delete oldest route items[%d] success", item.RouteVersion)
}
return nil
}
type routeItemRing struct {
data []*routeItem
head uint32
tail uint32
nextTail uint32
cap uint32
usedCap uint32
}
func newRouteItemRing(cap uint32) *routeItemRing {
ring := &routeItemRing{
data: make([]*routeItem, cap),
cap: cap,
}
return ring
}
func (r *routeItemRing) put(item *routeItem) {
r.data[r.nextTail] = item
r.tail = r.nextTail
if r.cap == r.usedCap {
r.head++
r.head = r.head % r.cap
r.nextTail++
r.nextTail = r.nextTail % r.cap
} else {
r.nextTail++
r.nextTail = r.nextTail % r.cap
r.usedCap++
}
if (r.data[r.head].RouteVersion + proto.RouteVersion(r.usedCap)) != (r.data[r.tail].RouteVersion + 1) {
errMsg := fmt.Sprintf("route cache ring is not consistently, head %v ver: %v, usedCap: %v, tail %v ver: %v",
r.head, r.data[r.head].RouteVersion, r.usedCap, r.tail, r.data[r.tail].RouteVersion)
panic(errMsg)
}
}
func (r *routeItemRing) getFrom(ver proto.RouteVersion) (ret []*routeItem, isLatest bool) {
if r.head == r.tail {
return nil, true
}
if r.getMinVer() > ver {
return nil, false
}
if r.getMaxVer() <= ver {
return nil, true
}
headVer := r.data[r.head].RouteVersion
i := (r.head + uint32(ver-headVer+1)) % r.cap
for j := 0; j < int(r.usedCap); j++ {
ret = append(ret, r.data[i])
i = (i + 1) % r.cap
if i == r.nextTail {
break
}
}
return ret, false
}
func (r *routeItemRing) getMinVer() proto.RouteVersion {
return r.data[r.head].RouteVersion
}
func (r *routeItemRing) getMaxVer() proto.RouteVersion {
return r.data[r.tail].RouteVersion
}

View File

@ -21,139 +21,19 @@ import (
"os"
"strconv"
"testing"
"time"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/clustermgr/persistence/catalogdb"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/util/log"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestRouteMgr(t *testing.T) {
ctx := context.Background()
ringBufferSize := uint32(3)
catalogDBPath := os.TempDir() + "/" + uuid.NewString() + strconv.FormatInt(rand.Int63n(math.MaxInt64), 10)
catalogDB, err := catalogdb.Open(catalogDBPath)
if err != nil {
log.Error("open db error")
return
}
defer os.RemoveAll(catalogDBPath)
storage, err := catalogdb.OpenCatalogTable(catalogDB)
if err != nil {
log.Error("open catalog table error")
return
}
// routeMgr
routeMgr := newRouteMgr(ringBufferSize, storage)
// add 3 items
items := make([]*routeItem, 0)
for i := 1; i <= 3; i++ {
item := &routeItem{
RouteVersion: proto.RouteVersion(routeMgr.genRouteVersion(ctx, 1)),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &routeItemShardAdd{ShardID: proto.ShardID(i)},
}
items = append(items, item)
}
routeMgr.insertRouteItems(ctx, items)
assert.Equal(t, uint64(3), routeMgr.getRouteVersion())
// add 1 item
item4 := &routeItem{
RouteVersion: proto.RouteVersion(routeMgr.genRouteVersion(ctx, 1)),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &routeItemShardAdd{ShardID: 4},
}
routeMgr.insertRouteItems(ctx, []*routeItem{item4})
assert.Equal(t, uint64(4), routeMgr.getRouteVersion())
// storage items
items = append(items, item4)
var routeRecords []*catalogdb.RouteInfoRecord
for _, item := range items {
routeRecords = append(routeRecords, routeItemToRouteRecord(item))
}
err = storage.PutShardsAndUnitsAndRouteItems(nil, nil, routeRecords)
assert.NoError(t, err)
// new a second routemgr
routeMgr2 := newRouteMgr(ringBufferSize, storage)
err = routeMgr2.loadRoute(ctx)
assert.NoError(t, err)
go routeMgr2.loop()
// 4 item will auto truncate to 3 item
items2, isLatest := routeMgr2.getRouteItems(ctx, 2)
assert.Equal(t, false, isLatest)
assert.Equal(t, 2, len(items2))
// get the latest version
items3, isLatest := routeMgr2.getRouteItems(ctx, 4)
assert.Equal(t, true, isLatest)
assert.Equal(t, 0, len(items3))
// add the 5th item, and then the items is [2,3,4,5]
item5 := &routeItem{
RouteVersion: proto.RouteVersion(routeMgr2.genRouteVersion(ctx, 1)),
Type: proto.CatalogChangeItemUpdateShard,
ItemDetail: &routeItemShardUpdate{SuidPrefix: proto.EncodeSuidPrefix(5, 1)},
}
routeMgr2.insertRouteItems(ctx, []*routeItem{item5})
assert.Equal(t, uint64(5), routeMgr2.getRouteVersion())
// remove the item
err = routeMgr2.removeOldRouteItems(ctx)
assert.NoError(t, err)
// now the items is [3,4,5], so get 2 is null
items4, isLatest := routeMgr2.getRouteItems(ctx, 2)
assert.Equal(t, false, isLatest)
assert.Equal(t, 0, len(items4))
items5, isLatest := routeMgr2.getRouteItems(ctx, 3)
assert.Equal(t, false, isLatest)
assert.Equal(t, 2, len(items5))
routeMgr2.Close()
}
func TestRouteItemRing(t *testing.T) {
ring := newRouteItemRing(3)
items, isLatest := ring.getFrom(3)
assert.Equal(t, 0, len(items))
assert.Equal(t, true, isLatest)
for i := 1; i <= 3; i++ {
item := &routeItem{
RouteVersion: proto.RouteVersion(i),
}
ring.put(item)
}
assert.Equal(t, proto.RouteVersion(1), ring.getMinVer())
assert.Equal(t, proto.RouteVersion(3), ring.getMaxVer())
items, isLatest = ring.getFrom(1)
assert.Equal(t, 2, len(items))
assert.Equal(t, false, isLatest)
items, isLatest = ring.getFrom(3)
assert.Equal(t, 0, len(items))
assert.Equal(t, true, isLatest)
item4 := &routeItem{
RouteVersion: proto.RouteVersion(4),
}
ring.put(item4)
assert.Equal(t, ring.getMinVer(), proto.RouteVersion(2))
assert.Equal(t, ring.getMaxVer(), proto.RouteVersion(4))
}
func TestCatalogMgr_Route(t *testing.T) {
mockCatalogMgr, clean := initMockCatalogMgr(t, testConfig)
defer clean()
@ -182,3 +62,96 @@ func TestCatalogMgr_Route(t *testing.T) {
require.Equal(t, 0, len(ret.Items))
require.Equal(t, proto.RouteVersion(0), ret.RouteVersion)
}
func TestRouteMgr(t *testing.T) {
ctx := context.Background()
ringBufferSize := uint32(3)
catalogDBPath := os.TempDir() + "/" + uuid.NewString() + strconv.FormatInt(rand.Int63n(math.MaxInt64), 10)
catalogDB, err := catalogdb.Open(catalogDBPath)
if err != nil {
log.Error("open db error")
return
}
defer os.RemoveAll(catalogDBPath)
storage, err := catalogdb.OpenCatalogTable(catalogDB)
if err != nil {
log.Error("open catalog table error")
return
}
base.RemoveOldRouteInternal = 1 * time.Second
// routeMgr
routeMgr := base.NewRouteMgr(ringBufferSize, false, routeRecordToRouteItem, storage)
// add 3 items
items := make([]*base.RouteItem, 0)
for i := 1; i <= 3; i++ {
item := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeMgr.GenRouteVersion(ctx, 1)),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &routeItemShardAdd{ShardID: proto.ShardID(i)},
}
items = append(items, item)
}
routeMgr.InsertRouteItems(ctx, items)
require.Equal(t, uint64(3), routeMgr.GetRouteVersion())
// add 1 item
item4 := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeMgr.GenRouteVersion(ctx, 1)),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &routeItemShardAdd{ShardID: 4},
}
routeMgr.InsertRouteItems(ctx, []*base.RouteItem{item4})
require.Equal(t, uint64(4), routeMgr.GetRouteVersion())
// storage items
items = append(items, item4)
var routeRecords []*base.RouteInfoRecord
for _, item := range items {
routeRecords = append(routeRecords, routeItemToRouteRecord(item))
}
err = storage.PutShardsAndUnitsAndRouteItems(nil, nil, routeRecords)
require.NoError(t, err)
// new a second routemgr
routeMgr2 := base.NewRouteMgr(ringBufferSize, false, routeRecordToRouteItem, storage)
err = routeMgr2.LoadRoute(ctx)
require.NoError(t, err)
go routeMgr2.Loop()
// 4 item will auto truncate to 3 item
items2, isLatest := routeMgr2.GetRouteItems(ctx, 2)
require.Equal(t, false, isLatest)
require.Equal(t, 2, len(items2))
// get the latest version
items3, isLatest := routeMgr2.GetRouteItems(ctx, 4)
require.Equal(t, true, isLatest)
require.Equal(t, 0, len(items3))
// add the 5th item, and then the items is [2,3,4,5]
item5 := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeMgr2.GenRouteVersion(ctx, 1)),
Type: proto.CatalogChangeItemUpdateShard,
ItemDetail: &routeItemShardUpdate{SuidPrefix: proto.EncodeSuidPrefix(5, 1)},
}
routeMgr2.InsertRouteItems(ctx, []*base.RouteItem{item5})
require.Equal(t, uint64(5), routeMgr2.GetRouteVersion())
// wait util remove the old items done
time.Sleep(3 * time.Second)
route, err := storage.GetFirstRoute()
require.NoError(t, err)
require.Equal(t, proto.RouteVersion(3), route.RouteVersion)
// now the items is [3,4,5], so get 2 is null
items4, isLatest := routeMgr2.GetRouteItems(ctx, 2)
require.Equal(t, false, isLatest)
require.Equal(t, 0, len(items4))
items5, isLatest := routeMgr2.GetRouteItems(ctx, 3)
require.Equal(t, false, isLatest)
require.Equal(t, 2, len(items5))
routeMgr2.Close()
}

View File

@ -325,8 +325,8 @@ func (c *CatalogMgr) applyUpdateShardUnit(ctx context.Context, newSuid proto.Sui
return err
}
newRouteVersion := c.routeMgr.genRouteVersion(ctx, 1)
route := &routeItem{
newRouteVersion := c.routeMgr.GenRouteVersion(ctx, 1)
route := &base.RouteItem{
RouteVersion: proto.RouteVersion(newRouteVersion),
Type: proto.CatalogChangeItemUpdateShard,
ItemDetail: &routeItemShardUpdate{SuidPrefix: newSuid.SuidPrefix()},
@ -342,7 +342,7 @@ func (c *CatalogMgr) applyUpdateShardUnit(ctx context.Context, newSuid proto.Sui
shardRecords := []*catalogdb.ShardInfoRecord{shardRecord}
unitRecords := []*catalogdb.ShardUnitInfoRecord{shardUnitRecord}
routeRecords := []*catalogdb.RouteInfoRecord{routeItemToRouteRecord(route)}
routeRecords := []*base.RouteInfoRecord{routeItemToRouteRecord(route)}
err := c.catalogTbl.UpdateUnitsAndPutShardsAndRouteItems(shardRecords, unitRecords, routeRecords)
if err != nil {
return err
@ -355,7 +355,7 @@ func (c *CatalogMgr) applyUpdateShardUnit(ctx context.Context, newSuid proto.Sui
shard.info.Units[index].Learner = learner
shard.info.Units[index].Host = diskInfo.Host
shard.info.RouteVersion = proto.RouteVersion(newRouteVersion)
c.routeMgr.insertRouteItems(ctx, []*routeItem{route})
c.routeMgr.InsertRouteItems(ctx, []*base.RouteItem{route})
return nil
})
if err != nil {

View File

@ -85,7 +85,7 @@ func NewCatalogMgr(conf Config, diskMgr cluster.ShardNodeManagerAPI, scopeMgr sc
applyTaskPool: base.NewTaskDistribution(int(conf.ApplyConcurrency), 1),
scopeMgr: scopeMgr,
kvMgr: kvMgr,
routeMgr: newRouteMgr(conf.RouteItemTruncateIntervalNum, catalogTable),
routeMgr: base.NewRouteMgr(conf.RouteItemTruncateIntervalNum, false, routeRecordToRouteItem, catalogTable),
diskMgr: diskMgr,
shardNodeClient: shardnode.New(conf.ShardNodeConfig),
closeLoopChan: make(chan struct{}, 1),
@ -162,5 +162,5 @@ func (c *CatalogMgr) loadSpace(ctx context.Context) error {
}
func (c *CatalogMgr) loadRoute(ctx context.Context) error {
return c.routeMgr.loadRoute(ctx)
return c.routeMgr.LoadRoute(ctx)
}

View File

@ -41,7 +41,7 @@ func (s *Service) CatalogChangesGet(c *rpc.Context) {
ret, err := s.CatalogMgr.GetCatalogChanges(ctx, args)
if err != nil {
span.Errorf("get catalog changes err =>", errors.Detail(err))
span.Errorf("get catalog changes err => %v", errors.Detail(err))
c.RespondError(err)
return
}

View File

@ -126,6 +126,7 @@ func NewHandler(service *Service) *rpc.Router {
rpc.RegisterArgsParser(&clustermgr.ListVolumeV2Args{}, "json")
rpc.RegisterArgsParser(&clustermgr.ListVolumeUnitArgs{}, "json")
rpc.RegisterArgsParser(&clustermgr.ListAllocatedVolumeArgs{}, "json")
rpc.RegisterArgsParser(&clustermgr.GetVolumeRoutesArgs{}, "json")
rpc.GET("/volume/get", service.VolumeGet, rpc.OptArgsQuery())
@ -155,6 +156,8 @@ func NewHandler(service *Service) *rpc.Router {
rpc.POST("/admin/update/volume", service.AdminUpdateVolume, rpc.OptArgsBody())
rpc.GET("/volumeroutes/get", service.VolumeRoutesGet, rpc.OptArgsQuery())
//==================shard==========================
rpc.RegisterArgsParser(&clustermgr.GetShardArgs{}, "json")
rpc.RegisterArgsParser(&clustermgr.ListShardArgs{}, "json")

View File

@ -20,6 +20,7 @@ import (
"fmt"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/common/kvstore"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/sharding"
@ -50,12 +51,6 @@ type SpaceInfoRecord struct {
SecretKey string `json:"secret_key"`
}
type RouteInfoRecord struct {
RouteVersion proto.RouteVersion `json:"route_version"`
Type proto.CatalogChangeItemType `json:"type"`
ItemDetail interface{} `json:"item"`
}
type ShardInfoRecord struct {
ShardID proto.ShardID `json:"shard_id"`
SuidPrefixes []proto.SuidPrefix `json:"suid_prefixes"`
@ -169,7 +164,7 @@ func (c *CatalogTable) ListSpace(count uint32, afterSpace proto.SpaceID) (ret []
return
}
func (c *CatalogTable) PutShardsAndUnitsAndRouteItems(shards []*ShardInfoRecord, units []*ShardUnitInfoRecord, routes []*RouteInfoRecord) error {
func (c *CatalogTable) PutShardsAndUnitsAndRouteItems(shards []*ShardInfoRecord, units []*ShardUnitInfoRecord, routes []*base.RouteInfoRecord) error {
batch := c.shardTbl.NewWriteBatch()
defer batch.Destroy()
@ -209,7 +204,7 @@ func (c *CatalogTable) PutShardsAndUnitsAndRouteItems(shards []*ShardInfoRecord,
return c.shardTbl.DoBatch(batch)
}
func (c *CatalogTable) UpdateUnitsAndPutShardsAndRouteItems(shards []*ShardInfoRecord, units []*ShardUnitInfoRecord, routes []*RouteInfoRecord) error {
func (c *CatalogTable) UpdateUnitsAndPutShardsAndRouteItems(shards []*ShardInfoRecord, units []*ShardUnitInfoRecord, routes []*base.RouteInfoRecord) error {
batch := c.shardTbl.NewWriteBatch()
defer batch.Destroy()
@ -379,7 +374,7 @@ func (c *CatalogTable) ListShardUnit(diskID proto.DiskID) (ret []proto.SuidPrefi
return
}
func (c *CatalogTable) GetFirstRouteItem() (*RouteInfoRecord, error) {
func (c *CatalogTable) GetFirstRoute() (*base.RouteInfoRecord, error) {
iter := c.routeTbl.NewIterator(nil)
defer iter.Close()
@ -390,34 +385,34 @@ func (c *CatalogTable) GetFirstRouteItem() (*RouteInfoRecord, error) {
}
if iter.Key().Size() > 0 {
ret, err := decodeRouteInfoRecord(iter.Value().Data())
iter.Key().Free()
iter.Value().Free()
if err != nil {
return nil, errors.Info(err, "decode route info db failed").Detail(err)
}
iter.Key().Free()
iter.Value().Free()
return ret, nil
}
}
return nil, nil
}
func (c *CatalogTable) ListRoute() ([]*RouteInfoRecord, error) {
func (c *CatalogTable) ListRoute() ([]*base.RouteInfoRecord, error) {
iter := c.routeTbl.NewIterator(nil)
defer iter.Close()
ret := make([]*RouteInfoRecord, 0)
ret := make([]*base.RouteInfoRecord, 0)
for iter.SeekToFirst(); iter.Valid(); iter.Next() {
if iter.Err() != nil {
return nil, iter.Err()
}
if iter.Key().Size() > 0 {
info, err := decodeRouteInfoRecord(iter.Value().Data())
iter.Key().Free()
iter.Value().Free()
if err != nil {
return nil, errors.Info(err, "decode route info db failed").Detail(err)
}
ret = append(ret, info)
iter.Key().Free()
iter.Value().Free()
}
}
return ret, nil
@ -445,7 +440,7 @@ func decodeSpaceInfoRecord(data []byte) (*SpaceInfoRecord, error) {
return ret, err
}
func encodeRouteInfoRecord(info *RouteInfoRecord) ([]byte, error) {
func encodeRouteInfoRecord(info *base.RouteInfoRecord) ([]byte, error) {
data, err := json.Marshal(info)
if err != nil {
return nil, err
@ -453,19 +448,21 @@ func encodeRouteInfoRecord(info *RouteInfoRecord) ([]byte, error) {
return data, nil
}
func decodeRouteInfoRecord(data []byte) (*RouteInfoRecord, error) {
ret := &RouteInfoRecord{}
func decodeRouteInfoRecord(data []byte) (*base.RouteInfoRecord, error) {
ret := &base.RouteInfoRecord{}
err := json.Unmarshal(data, ret)
if err != nil {
return nil, err
}
switch ret.Type {
switch proto.CatalogChangeItemType(ret.Type.(float64)) {
case proto.CatalogChangeItemAddShard:
ret.ItemDetail = &RouteInfoShardAdd{}
err = json.Unmarshal(data, ret)
ret.Type = proto.CatalogChangeItemAddShard
case proto.CatalogChangeItemUpdateShard:
ret.ItemDetail = &RouteInfoShardUpdate{}
err = json.Unmarshal(data, ret)
ret.Type = proto.CatalogChangeItemUpdateShard
default:
panic(fmt.Sprintf("unsupported route item type: %d", ret.Type))
}

View File

@ -18,10 +18,11 @@ import (
"testing"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/common/sharding"
"github.com/stretchr/testify/require"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/sharding"
"github.com/stretchr/testify/require"
)
var catalogTable *CatalogTable
@ -72,19 +73,19 @@ var (
RouteVersion: proto.RouteVersion(3),
}
route1 = &RouteInfoRecord{
route1 = &base.RouteInfoRecord{
RouteVersion: shard1.RouteVersion,
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &RouteInfoShardAdd{ShardID: shard1.ShardID},
}
route2 = &RouteInfoRecord{
route2 = &base.RouteInfoRecord{
RouteVersion: shard2.RouteVersion,
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &RouteInfoShardAdd{ShardID: shard2.ShardID},
}
route3 = &RouteInfoRecord{
route3 = &base.RouteInfoRecord{
RouteVersion: shard3.RouteVersion,
Type: proto.CatalogChangeItemUpdateShard,
ItemDetail: &RouteInfoShardUpdate{SuidPrefix: shard3.SuidPrefixes[0]},
@ -124,7 +125,7 @@ var (
shards = []*ShardInfoRecord{shard1, shard2, shard3}
shardUnits = []*ShardUnitInfoRecord{shardUnit1, shardUnit2, shardUnit3}
routes = []*RouteInfoRecord{route1, route2, route3}
routes = []*base.RouteInfoRecord{route1, route2, route3}
spaces = []*SpaceInfoRecord{space1, space2, space3}
)
@ -283,7 +284,7 @@ func TestCatalogTable_Route(t *testing.T) {
err := catalogTable.PutShardsAndUnitsAndRouteItems(nil, nil, routes)
require.NoError(t, err)
routeInfoRecord, err := catalogTable.GetFirstRouteItem()
routeInfoRecord, err := catalogTable.GetFirstRoute()
require.NoError(t, err)
require.Equal(t, route1.RouteVersion, routeInfoRecord.RouteVersion)
require.Equal(t, route1.Type, routeInfoRecord.Type)

View File

@ -24,6 +24,7 @@ import (
// column family name definition
var (
routeCF = "volume_route"
volumeCF = "volume"
volumeUnitCF = "volume_unit"
volumeTokenCF = "volume_token"
@ -33,6 +34,7 @@ var (
volumeUnitDiskIDIndexCF = "volumeUnit_DiskID"
volumeCfs = []string{
routeCF,
volumeCF,
volumeUnitCF,
volumeTokenCF,

View File

@ -18,6 +18,7 @@ import (
"bytes"
"encoding/binary"
"encoding/gob"
"encoding/json"
"fmt"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
@ -39,6 +40,7 @@ type VolumeTable struct {
unitTbl kvstore.KVTable
tokenTbl kvstore.KVTable
taskTbl kvstore.KVTable
routeTbl kvstore.KVTable
indexes map[string]indexItem
}
@ -53,6 +55,7 @@ type VolumeRecord struct {
Used uint64
CreateByNodeID uint64
Epoch uint32
RouteVersion proto.RouteVersion
}
type VolumeTaskRecord struct {
@ -78,11 +81,20 @@ type TokenRecord struct {
ExpireTime int64
}
type RouteInfoVolumeAdd struct {
Vid proto.Vid `json:"vid"`
}
type RouteInfoVolumeUpdate struct {
VuidPrefix proto.VuidPrefix `json:"vuid_prefix"`
}
func OpenVolumeTable(db kvstore.KVStore) (*VolumeTable, error) {
if db == nil {
return nil, errors.New("OpenVolumeTable failed: db is nil")
}
return &VolumeTable{
routeTbl: db.Table(routeCF),
volTbl: db.Table(volumeCF),
unitTbl: db.Table(volumeUnitCF),
tokenTbl: db.Table(volumeTokenCF),
@ -238,7 +250,7 @@ func (v *VolumeTable) PutVolumeAndToken(volumeRecs []*VolumeRecord, tokens []*To
return v.volTbl.DoBatch(batch)
}
func (v *VolumeTable) PutVolumeAndVolumeUnit(volumeRecs []*VolumeRecord, volumeUnitsRecs [][]*VolumeUnitRecord) (err error) {
func (v *VolumeTable) PutVolumesAndUnitsAndRoutes(volumeRecs []*VolumeRecord, volumeUnitsRecs [][]*VolumeUnitRecord, routes []*base.RouteInfoRecord) (err error) {
batch := v.volTbl.NewWriteBatch()
defer batch.Destroy()
@ -268,6 +280,14 @@ func (v *VolumeTable) PutVolumeAndVolumeUnit(volumeRecs []*VolumeRecord, volumeU
batch.PutCF(v.volTbl.GetCf(), vid, valueVol)
}
for i := range routes {
routeData, err := encodeRouteInfoRecord(routes[i])
if err != nil {
return err
}
batch.PutCF(v.routeTbl.GetCf(), encodeRouteKey(routes[i].RouteVersion), routeData)
}
return v.volTbl.DoBatch(batch)
}
@ -326,6 +346,58 @@ func (v *VolumeTable) ListVolume(count int, afterVid proto.Vid) (ret []proto.Vid
return
}
func (v *VolumeTable) GetFirstRoute() (*base.RouteInfoRecord, error) {
iter := v.routeTbl.NewIterator(nil)
defer iter.Close()
iter.SeekToFirst()
if iter.Valid() {
if iter.Err() != nil {
return nil, iter.Err()
}
if iter.Key().Size() > 0 {
ret, err := decodeRouteInfoRecord(iter.Value().Data())
iter.Key().Free()
iter.Value().Free()
if err != nil {
return nil, errors.Info(err, "decode route info db failed").Detail(err)
}
return ret, nil
}
}
return nil, nil
}
func (v *VolumeTable) ListRoute() ([]*base.RouteInfoRecord, error) {
iter := v.routeTbl.NewIterator(nil)
defer iter.Close()
ret := make([]*base.RouteInfoRecord, 0)
for iter.SeekToFirst(); iter.Valid(); iter.Next() {
if iter.Err() != nil {
return nil, iter.Err()
}
if iter.Key().Size() > 0 {
info, err := decodeRouteInfoRecord(iter.Value().Data())
iter.Key().Free()
iter.Value().Free()
if err != nil {
return nil, errors.Info(err, "decode route info db failed").Detail(err)
}
ret = append(ret, info)
}
}
return ret, nil
}
func (v *VolumeTable) DeleteOldRoutes(before proto.RouteVersion) error {
batch := v.routeTbl.NewWriteBatch()
defer batch.Destroy()
batch.DeleteRangeCF(v.routeTbl.GetCf(), encodeRouteKey(0), encodeRouteKey(before))
return v.routeTbl.DoBatch(batch)
}
func decodeVolumeRecord(volByte []byte) (ret *VolumeRecord, err error) {
dec := gob.NewDecoder(bytes.NewReader(volByte))
err = dec.Decode(&ret)
@ -518,8 +590,8 @@ func (v *VolumeTable) ListVolumeUnit(diskID proto.DiskID) (ret []proto.VuidPrefi
return
}
func (v *VolumeTable) UpdateVolumeUnit(vuidPrefix proto.VuidPrefix, unitRecord *VolumeUnitRecord) (err error) {
keyVuidPrefix := encodeVuidPrefix(vuidPrefix)
func (v *VolumeTable) UpdateVolumeUnitAndPutVolumeAndRoute(unitRecord *VolumeUnitRecord, volumeRecord *VolumeRecord, routeRecord *base.RouteInfoRecord) (err error) {
keyVuidPrefix := encodeVuidPrefix(unitRecord.VuidPrefix)
value, err := encodeVolumeUnitRecord(unitRecord)
if err != nil {
return err
@ -542,13 +614,27 @@ func (v *VolumeTable) UpdateVolumeUnit(vuidPrefix proto.VuidPrefix, unitRecord *
}
oldDiskID := uRec.DiskID
oldIndexKey := ""
oldIndexKey += fmtIndexKey(indexName, oldDiskID, vuidPrefix)
oldIndexKey += fmtIndexKey(indexName, oldDiskID, unitRecord.VuidPrefix)
batch.DeleteCF(v.indexes[volumeUintDiskIDIndex].indexTbl.GetCf(), []byte(oldIndexKey))
indexKey += fmtIndexKey(indexName, unitRecord.DiskID, vuidPrefix)
indexKey += fmtIndexKey(indexName, unitRecord.DiskID, unitRecord.VuidPrefix)
batch.PutCF(v.indexes[volumeUintDiskIDIndex].indexTbl.GetCf(), []byte(indexKey), keyVuidPrefix)
batch.PutCF(v.unitTbl.GetCf(), keyVuidPrefix, value)
// update volume record
volumeData, err := encodeVolumeRecord(volumeRecord)
if err != nil {
return err
}
batch.PutCF(v.volTbl.GetCf(), EncodeVid(volumeRecord.Vid), volumeData)
// insert route record
routeData, err := encodeRouteInfoRecord(routeRecord)
if err != nil {
return err
}
batch.PutCF(v.routeTbl.GetCf(), encodeRouteKey(routeRecord.RouteVersion), routeData)
return v.unitTbl.DoBatch(batch)
}
@ -598,6 +684,41 @@ func decodeTaskRecord(buf []byte) (ret *VolumeTaskRecord, err error) {
return
}
func encodeRouteInfoRecord(info *base.RouteInfoRecord) ([]byte, error) {
data, err := json.Marshal(info)
if err != nil {
return nil, err
}
return data, nil
}
func decodeRouteInfoRecord(data []byte) (*base.RouteInfoRecord, error) {
ret := &base.RouteInfoRecord{}
err := json.Unmarshal(data, ret)
if err != nil {
return nil, err
}
switch proto.VolumeRouteItemType(ret.Type.(float64)) {
case proto.RouteItemTypeAddVolume:
ret.ItemDetail = &RouteInfoVolumeAdd{}
err = json.Unmarshal(data, ret)
ret.Type = proto.RouteItemTypeAddVolume
case proto.RouteItemTypeUpdateVolume:
ret.ItemDetail = &RouteInfoVolumeUpdate{}
err = json.Unmarshal(data, ret)
ret.Type = proto.RouteItemTypeUpdateVolume
default:
panic(fmt.Sprintf("unsupported route item type: %d", ret.Type))
}
return ret, err
}
func encodeRouteKey(ver proto.RouteVersion) []byte {
ret := make([]byte, 8)
binary.BigEndian.PutUint64(ret, uint64(ver))
return ret
}
func fmtIndexKey(name string, diskID proto.DiskID, vuidPrefix proto.VuidPrefix) string {
return fmt.Sprintf("%s-%d-%d", name, diskID, vuidPrefix)
}

View File

@ -94,6 +94,7 @@ var (
Used: 21,
Total: 10240,
CreateByNodeID: 1,
RouteVersion: proto.RouteVersion(1),
}
volume2 = &VolumeRecord{
Vid: 2,
@ -105,6 +106,7 @@ var (
Used: 21,
Total: 10240,
CreateByNodeID: 1,
RouteVersion: proto.RouteVersion(2),
}
volume3 = &VolumeRecord{
Vid: 3,
@ -116,9 +118,29 @@ var (
Used: 21,
Total: 10240,
CreateByNodeID: 1,
RouteVersion: proto.RouteVersion(3),
}
route1 = &base.RouteInfoRecord{
RouteVersion: volume1.RouteVersion,
Type: proto.RouteItemTypeAddVolume,
ItemDetail: &RouteInfoVolumeAdd{Vid: volume1.Vid},
}
route2 = &base.RouteInfoRecord{
RouteVersion: volume2.RouteVersion,
Type: proto.RouteItemTypeAddVolume,
ItemDetail: &RouteInfoVolumeAdd{Vid: volume2.Vid},
}
route3 = &base.RouteInfoRecord{
RouteVersion: volume3.RouteVersion,
Type: proto.RouteItemTypeAddVolume,
ItemDetail: &RouteInfoVolumeAdd{Vid: volume3.Vid},
}
volumes = []*VolumeRecord{volume1, volume2, volume3}
routes = []*base.RouteInfoRecord{route1, route2, route3}
taskRecord1 = &VolumeTaskRecord{
Vid: 1,
@ -186,7 +208,7 @@ func TestVolumeTable_PutVolumeAndVolumeUnit(t *testing.T) {
defer closeVolumeDB()
volumeUnits := [][]*VolumeUnitRecord{{volumeUnit1}, {volumeUnit2}, {volumeUnit3}}
err := volumeTable.PutVolumeAndVolumeUnit(volumes, volumeUnits)
err := volumeTable.PutVolumesAndUnitsAndRoutes(volumes, volumeUnits, routes)
require.NoError(t, err)
}
@ -368,14 +390,19 @@ func TestVolumeTable_UpdateVolumeUnit(t *testing.T) {
require.Equal(t, len(ret), 2)
// repeat update volume unit
err = volumeTable.UpdateVolumeUnit(volumeUnitInfo2.VuidPrefix, volumeUnitInfo2)
route := &base.RouteInfoRecord{
RouteVersion: volume1.RouteVersion,
Type: proto.RouteItemTypeUpdateVolume,
ItemDetail: &RouteInfoVolumeUpdate{VuidPrefix: volumeUnitInfo2.VuidPrefix},
}
err = volumeTable.UpdateVolumeUnitAndPutVolumeAndRoute(volumeUnitInfo2, volume2, route)
require.NoError(t, err)
ret, err = volumeTable.ListVolumeUnit(20)
require.NoError(t, err)
require.Equal(t, 2, len(ret))
volumeUnit2.DiskID = 45
err = volumeTable.UpdateVolumeUnit(volumeUnit2.VuidPrefix, volumeUnit2)
err = volumeTable.UpdateVolumeUnitAndPutVolumeAndRoute(volumeUnit2, volume2, route2)
require.NoError(t, err)
ret, err = volumeTable.ListVolumeUnit(20)
@ -455,3 +482,28 @@ func TestVolumeTable_PutVolumeAndTask(t *testing.T) {
err = volumeTable.PutVolumeAndTask(nil, taskRecord1)
require.NoError(t, err)
}
func TestVolumeTable_Route(t *testing.T) {
initVolumeDB()
defer closeVolumeDB()
err := volumeTable.PutVolumesAndUnitsAndRoutes(nil, nil, routes)
require.NoError(t, err)
routeInfoRecord, err := volumeTable.GetFirstRoute()
require.NoError(t, err)
require.Equal(t, route1.RouteVersion, routeInfoRecord.RouteVersion)
require.Equal(t, route1.Type, routeInfoRecord.Type)
require.Equal(t, route1.ItemDetail, routeInfoRecord.ItemDetail)
routeInfoRecords, err := volumeTable.ListRoute()
require.NoError(t, err)
require.Equal(t, 3, len(routeInfoRecords))
err = volumeTable.DeleteOldRoutes(route2.RouteVersion)
require.NoError(t, err)
routeInfoRecords, err = volumeTable.ListRoute()
require.NoError(t, err)
require.Equal(t, 2, len(routeInfoRecords))
}

View File

@ -28,6 +28,7 @@ import (
"github.com/stretchr/testify/require"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/clustermgr/persistence/catalogdb"
"github.com/cubefs/cubefs/blobstore/clustermgr/persistence/normaldb"
"github.com/cubefs/cubefs/blobstore/common/codemode"
@ -242,7 +243,7 @@ func generateShard(catalogDBPath, NormalDBPath string) error {
unitCount = 3
shards []*catalogdb.ShardInfoRecord
units []*catalogdb.ShardUnitInfoRecord
routes []*catalogdb.RouteInfoRecord
routes []*base.RouteInfoRecord
ranges = sharding.InitShardingRange(sharding.RangeType_RangeTypeHash, 2, 10)
)
catalogDB, err := catalogdb.Open(catalogDBPath)
@ -281,7 +282,7 @@ func generateShard(catalogDBPath, NormalDBPath string) error {
RouteVersion: proto.RouteVersion(i),
}
route := &catalogdb.RouteInfoRecord{
route := &base.RouteInfoRecord{
RouteVersion: proto.RouteVersion(i),
Type: proto.CatalogChangeItemAddShard,
ItemDetail: &catalogdb.RouteInfoShardAdd{ShardID: proto.ShardID(i)},

View File

@ -0,0 +1,35 @@
package clustermgr
import (
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
apierrors "github.com/cubefs/cubefs/blobstore/common/errors"
"github.com/cubefs/cubefs/blobstore/common/rpc"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/util/errors"
)
func (s *Service) VolumeRoutesGet(c *rpc.Context) {
ctx := c.Request.Context()
span := trace.SpanFromContextSafe(ctx)
args := new(clustermgr.GetVolumeRoutesArgs)
if err := c.ParseArgs(args); err != nil {
c.RespondError(err)
return
}
span.Infof("accept VolumeRoutesGet request, args: %v", args)
// linear read
if err := s.raftNode.ReadIndex(ctx); err != nil {
span.Errorf("get volume routes read index error: %v", err)
c.RespondError(apierrors.ErrRaftReadIndex)
return
}
ret, err := s.VolumeMgr.GetVolumeRoutes(ctx, args)
if err != nil {
span.Errorf("get volume routes err => %v", errors.Detail(err))
c.RespondError(err)
return
}
c.RespondJSON(ret)
}

View File

@ -0,0 +1,35 @@
package clustermgr
import (
"testing"
"github.com/stretchr/testify/require"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/common/proto"
)
func TestVolume_Route(t *testing.T) {
testService, clean := initServiceWithData()
defer clean()
cmClient := initTestClusterClient(testService)
ctx := newCtx()
// get all routes
ret, err := cmClient.GetVolumeRoutes(ctx, &clustermgr.GetVolumeRoutesArgs{})
require.NoError(t, err)
require.Equal(t, 10, len(ret.Items))
require.Equal(t, proto.RouteVersion(1), ret.RouteVersion)
// get route without specified routeVersion
ret, err = cmClient.GetVolumeRoutes(ctx, &clustermgr.GetVolumeRoutesArgs{RouteVersion: 1})
require.NoError(t, err)
require.Equal(t, 0, len(ret.Items))
require.Equal(t, proto.RouteVersion(0), ret.RouteVersion)
// get route with routeVersion which is bigger than CM's
ret, err = cmClient.GetVolumeRoutes(ctx, &clustermgr.GetVolumeRoutesArgs{RouteVersion: 100})
require.NoError(t, err)
require.Equal(t, 0, len(ret.Items))
require.Equal(t, proto.RouteVersion(0), ret.RouteVersion)
}

View File

@ -536,7 +536,6 @@ func generateVolume(volumeDBPath, NormalDBPath string) error {
Total: 1024 * 1024 * 1024 * 1024,
Epoch: 1,
}
volumes = append(volumes, vol)
tokens = append(tokens, &volumedb.TokenRecord{
Vid: proto.Vid(i),

View File

@ -116,6 +116,9 @@ func (v *VolumeMgr) LoadData(ctx context.Context) error {
if err := v.reloadTasks(); err != nil {
return errors.Info(err, "reload task failed").Detail(err)
}
if err := v.loadRoute(ctx); err != nil {
return errors.Info(err, "load route failed").Detail(err)
}
return nil
}
@ -443,7 +446,7 @@ func (v *VolumeMgr) Flush(ctx context.Context) error {
}
vol.lock.RLock()
err = v.volumeTbl.PutVolumeAndVolumeUnit([]*volumedb.VolumeRecord{vol.ToRecord()}, [][]*volumedb.VolumeUnitRecord{volumeUnitsToVolumeUnitRecords(vol.vUnits)})
err = v.volumeTbl.PutVolumesAndUnitsAndRoutes([]*volumedb.VolumeRecord{vol.ToRecord()}, [][]*volumedb.VolumeUnitRecord{volumeUnitsToVolumeUnitRecords(vol.vUnits)}, nil)
vol.lock.RUnlock()
retErr = err
return

View File

@ -225,10 +225,20 @@ func (v *VolumeMgr) applyCreateVolume(ctx context.Context, vol *volume) error {
if err := v.transitedTbl.DeleteVolumeAndUnits(volumeRecord, unitRecords); err != nil {
return errors.Info(err, fmt.Sprintf("delete volume [%+v] and units[%+v] from transited table failed", volumeRecord, unitRecords)).Detail(err)
}
if err := v.volumeTbl.PutVolumeAndVolumeUnit([]*volumedb.VolumeRecord{volumeRecord}, [][]*volumedb.VolumeUnitRecord{unitRecords}); err != nil {
// insert route item
routeVersion := v.routeMgr.GenRouteVersion(ctx, 1)
route := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeVersion),
Type: proto.RouteItemTypeAddVolume,
ItemDetail: &routeItemVolumeAdd{Vid: vol.vid},
}
routeRecord := routeItemToRouteRecord(route)
if err := v.volumeTbl.PutVolumesAndUnitsAndRoutes([]*volumedb.VolumeRecord{volumeRecord}, [][]*volumedb.VolumeUnitRecord{unitRecords}, []*base.RouteInfoRecord{routeRecord}); err != nil {
return errors.Info(err, fmt.Sprintf("put volume[%+v] and volume unit[%+v] into volume table failed", volumeRecord, unitRecords)).Detail(err)
}
v.all.putVol(vol)
v.routeMgr.InsertRouteItems(ctx, []*base.RouteItem{route})
vol.volInfoBase.RouteVersion = proto.RouteVersion(routeVersion)
return nil
}

View File

@ -21,6 +21,7 @@ import (
"time"
cm "github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/clustermgr/persistence/volumedb"
"github.com/cubefs/cubefs/blobstore/common/codemode"
"github.com/cubefs/cubefs/blobstore/common/proto"
@ -69,6 +70,7 @@ func (vol *volume) ToRecord() *volumedb.VolumeRecord {
Used: vol.volInfoBase.Used,
CreateByNodeID: vol.volInfoBase.CreateByNodeID,
Epoch: vol.volInfoBase.Epoch,
RouteVersion: vol.volInfoBase.RouteVersion,
}
}
@ -278,6 +280,27 @@ func (s *shardedVolumes) rangeVol(f func(v *volume) error) {
}
}
func (s *shardedVolumes) getVolumeNum() int {
length := 0
for i := uint32(0); i < s.num; i++ {
l := s.locks[i]
l.RLock()
length += len(s.m[i])
l.RUnlock()
}
return length
}
func (s *shardedVolumes) list() []*volume {
tolerantCap := 1000 // tolerate volume is created when range
ret := make([]*volume, 0, s.getVolumeNum()+tolerantCap)
s.rangeVol(func(v *volume) error {
ret = append(ret, v)
return nil
})
return ret
}
type NotifyFunc func(ctx context.Context, vol *volume) error
type volumeNotifyQueue struct {
@ -322,6 +345,7 @@ func volumeRecordToVolumeInfoBase(volRecord *volumedb.VolumeRecord) cm.VolumeInf
Free: volRecord.Free,
CreateByNodeID: volRecord.CreateByNodeID,
Epoch: volRecord.Epoch,
RouteVersion: volRecord.RouteVersion,
}
}
@ -357,3 +381,48 @@ func tokenRecordToToken(record *volumedb.TokenRecord) (ret *token) {
expireTime: record.ExpireTime,
}
}
type routeItemVolumeAdd struct {
Vid proto.Vid
}
type routeItemVolumeUpdate struct {
VuidPrefix proto.VuidPrefix
}
func routeRecordToRouteItem(info *base.RouteInfoRecord) *base.RouteItem {
item := &base.RouteItem{
RouteVersion: info.RouteVersion,
Type: info.Type,
}
switch info.Type {
case proto.RouteItemTypeAddVolume:
itemDetail := info.ItemDetail.(*volumedb.RouteInfoVolumeAdd)
item.ItemDetail = &routeItemVolumeAdd{Vid: itemDetail.Vid}
case proto.RouteItemTypeUpdateVolume:
itemDetail := info.ItemDetail.(*volumedb.RouteInfoVolumeUpdate)
item.ItemDetail = &routeItemVolumeUpdate{VuidPrefix: itemDetail.VuidPrefix}
default:
}
return item
}
func routeItemToRouteRecord(item *base.RouteItem) *base.RouteInfoRecord {
record := &base.RouteInfoRecord{
RouteVersion: item.RouteVersion,
Type: item.Type,
}
switch item.Type {
case proto.RouteItemTypeAddVolume:
itemDetail := item.ItemDetail.(*routeItemVolumeAdd)
record.ItemDetail = &volumedb.RouteInfoVolumeAdd{Vid: itemDetail.Vid}
case proto.RouteItemTypeUpdateVolume:
itemDetail := item.ItemDetail.(*routeItemVolumeUpdate)
record.ItemDetail = &volumedb.RouteInfoVolumeUpdate{VuidPrefix: itemDetail.VuidPrefix}
default:
}
return record
}

View File

@ -0,0 +1,86 @@
package volumemgr
import (
"context"
"github.com/gogo/protobuf/types"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/trace"
)
func (v *VolumeMgr) GetVolumeRoutes(ctx context.Context, args *clustermgr.GetVolumeRoutesArgs) (ret *clustermgr.GetVolumeRoutesRet, err error) {
span := trace.SpanFromContextSafe(ctx)
var (
items []*base.RouteItem
isLatest bool
)
if args.RouteVersion > 0 {
items, isLatest = v.routeMgr.GetRouteItems(ctx, args.RouteVersion)
}
ret = new(clustermgr.GetVolumeRoutesRet)
if items == nil && !isLatest {
// get all routesmust getRouteVersion first
ret.RouteVersion = proto.RouteVersion(v.routeMgr.GetRouteVersion())
vols := v.all.list()
items = make([]*base.RouteItem, 0, len(vols))
for _, vol := range vols {
items = append(items, &base.RouteItem{
Type: proto.RouteItemTypeAddVolume,
ItemDetail: &routeItemVolumeAdd{Vid: vol.vid},
})
}
}
for i := range items {
ret.Items = append(ret.Items, clustermgr.VolumeRouteItem{
RouteVersion: items[i].RouteVersion,
Type: items[i].Type.(proto.VolumeRouteItemType),
})
switch items[i].Type {
case proto.RouteItemTypeAddVolume:
vid := items[i].ItemDetail.(*routeItemVolumeAdd).Vid
vol := v.all.getVol(vid)
addVolumeItem := &clustermgr.RouteItemAddVolume{
Vid: vid,
}
vol.withRLocked(func() error {
if items[i].RouteVersion == proto.InvalidRouteVersion {
ret.Items[i].RouteVersion = vol.volInfoBase.RouteVersion
}
for _, unit := range vol.vUnits {
addVolumeItem.Units = append(addVolumeItem.Units, clustermgr.VolumeUnitInfoBase(*unit.vuInfo))
}
return nil
})
ret.Items[i].Item, err = types.MarshalAny(addVolumeItem)
span.Debugf("addVolumeItem: %+v", addVolumeItem)
case proto.RouteItemTypeUpdateVolume:
vuidPrefix := items[i].ItemDetail.(*routeItemVolumeUpdate).VuidPrefix
vol := v.all.getVol(vuidPrefix.Vid())
updateVolumeItem := &clustermgr.RouteItemUpdateVolume{
Vid: vuidPrefix.Vid(),
}
vol.withRLocked(func() error {
unit := vol.vUnits[vuidPrefix.Index()]
updateVolumeItem.Unit = clustermgr.VolumeUnitInfoBase(*unit.vuInfo)
return nil
})
ret.Items[i].Item, err = types.MarshalAny(updateVolumeItem)
span.Debugf("updateVolumeItem: %+v", updateVolumeItem)
default:
}
if err != nil {
return nil, err
}
}
if ret.RouteVersion == 0 && len(ret.Items) > 0 {
ret.RouteVersion = ret.Items[len(ret.Items)-1].RouteVersion
}
return ret, nil
}

View File

@ -0,0 +1,149 @@
package volumemgr
import (
"context"
"math"
"math/rand"
"os"
"strconv"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/require"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/clustermgr/base"
"github.com/cubefs/cubefs/blobstore/clustermgr/persistence/volumedb"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/util/log"
)
func TestVolumeMgr_Route(t *testing.T) {
mockVolumeMgr, clean := initMockVolumeMgr(t)
defer clean()
_, ctx := trace.StartSpanFromContext(context.Background(), "route")
// get all routes
ret, err := mockVolumeMgr.GetVolumeRoutes(ctx, &clustermgr.GetVolumeRoutesArgs{})
require.NoError(t, err)
require.Equal(t, 30, len(ret.Items))
require.Equal(t, proto.RouteVersion(1), ret.RouteVersion)
// get routes with normal routerVersion
ret, err = mockVolumeMgr.GetVolumeRoutes(ctx, &clustermgr.GetVolumeRoutesArgs{
RouteVersion: 1,
})
require.NoError(t, err)
require.Equal(t, 0, len(ret.Items))
require.Equal(t, proto.RouteVersion(0), ret.RouteVersion)
// get routes with abnormal routerVersion
ret, err = mockVolumeMgr.GetVolumeRoutes(ctx, &clustermgr.GetVolumeRoutesArgs{
RouteVersion: 100,
})
require.NoError(t, err)
require.Equal(t, 0, len(ret.Items))
require.Equal(t, proto.RouteVersion(0), ret.RouteVersion)
}
func TestVolumeRouteMgr(t *testing.T) {
ctx := context.Background()
ringBufferSize := uint32(3)
dbPath := os.TempDir() + "/" + uuid.NewString() + strconv.FormatInt(rand.Int63n(math.MaxInt64), 10)
volumeDB, err := volumedb.Open(dbPath)
if err != nil {
log.Error("open db error")
return
}
defer os.RemoveAll(dbPath)
storage, err := volumedb.OpenVolumeTable(volumeDB)
if err != nil {
log.Error("open volume table error")
return
}
base.RemoveOldRouteInternal = 1 * time.Second
// routeMgr
routeMgr := base.NewRouteMgr(ringBufferSize, true, routeRecordToRouteItem, storage)
err = routeMgr.LoadRoute(ctx)
require.NoError(t, err)
require.Equal(t, uint64(1), routeMgr.GetRouteVersion())
// add 1 item, [2]
item1 := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeMgr.GenRouteVersion(ctx, 1)),
Type: proto.RouteItemTypeAddVolume,
ItemDetail: &routeItemVolumeAdd{Vid: 2},
}
routeMgr.InsertRouteItems(ctx, []*base.RouteItem{item1})
require.Equal(t, uint64(2), routeMgr.GetRouteVersion())
items, isLatest := routeMgr.GetRouteItems(ctx, 1)
require.Equal(t, false, isLatest)
require.Equal(t, 1, len(items))
// add 3 items, [3,4,5]
items = make([]*base.RouteItem, 0)
for i := 3; i <= 5; i++ {
item := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeMgr.GenRouteVersion(ctx, 1)),
Type: proto.RouteItemTypeAddVolume,
ItemDetail: &routeItemVolumeAdd{Vid: proto.Vid(i)},
}
items = append(items, item)
}
routeMgr.InsertRouteItems(ctx, items)
require.Equal(t, uint64(5), routeMgr.GetRouteVersion())
// storage items
items = append(items, item1)
var routeRecords []*base.RouteInfoRecord
for _, item := range items {
routeRecords = append(routeRecords, routeItemToRouteRecord(item))
}
err = storage.PutVolumesAndUnitsAndRoutes(nil, nil, routeRecords)
require.NoError(t, err)
// new a second routemgr
routeMgr2 := base.NewRouteMgr(ringBufferSize, true, routeRecordToRouteItem, storage)
err = routeMgr2.LoadRoute(ctx)
require.NoError(t, err)
go routeMgr2.Loop()
// 4 item will auto truncate to 3 item, [3,4,5]
items2, isLatest := routeMgr2.GetRouteItems(ctx, 3)
require.Equal(t, false, isLatest)
require.Equal(t, 2, len(items2))
// get the latest version
items3, isLatest := routeMgr2.GetRouteItems(ctx, 5)
require.Equal(t, true, isLatest)
require.Equal(t, 0, len(items3))
// add the 5th item, and then the items is [3,4,5,6]
item5 := &base.RouteItem{
RouteVersion: proto.RouteVersion(routeMgr2.GenRouteVersion(ctx, 1)),
Type: proto.RouteItemTypeUpdateVolume,
ItemDetail: &routeItemVolumeUpdate{VuidPrefix: proto.EncodeVuidPrefix(6, 1)},
}
routeMgr2.InsertRouteItems(ctx, []*base.RouteItem{item5})
require.Equal(t, uint64(6), routeMgr2.GetRouteVersion())
// wait util remove the old items done, [4,5,6]
time.Sleep(3 * time.Second)
route, err := storage.GetFirstRoute()
require.NoError(t, err)
require.Equal(t, proto.RouteVersion(4), route.RouteVersion)
// now the items is [4,5,6], so get 3 is null
items4, isLatest := routeMgr2.GetRouteItems(ctx, 3)
require.Equal(t, false, isLatest)
require.Equal(t, 0, len(items4))
items5, isLatest := routeMgr2.GetRouteItems(ctx, 4)
require.Equal(t, false, isLatest)
require.Equal(t, 2, len(items5))
routeMgr2.Close()
}

View File

@ -49,6 +49,7 @@ type VolumeMgrConfig struct {
VolumeSliceMapNum uint32 `json:"volume_slice_map_num"`
ApplyConcurrency uint32 `json:"apply_concurrency"`
RouteItemTruncateIntervalNum uint32 `json:"route_item_truncate_interval_num"`
MinAllocableVolumeCount int `json:"min_allocable_volume_count"`
AllocatableDiskLoadThreshold int `json:"allocatable_disk_load_threshold"`
AllocFactor int `json:"alloc_factor"`
@ -103,6 +104,9 @@ func (c *VolumeMgrConfig) checkAndFix() {
if c.ShardNum <= 0 {
c.ShardNum = defaultShardNum
}
if c.RouteItemTruncateIntervalNum <= 0 {
c.RouteItemTruncateIntervalNum = defaultRouteItemTruncateIntervalNum
}
}
// NewVolumeMgr constructs a new volume manager.
@ -141,6 +145,7 @@ func NewVolumeMgr(conf VolumeMgrConfig, diskMgr cluster.BlobNodeManagerAPI, scop
diskMgr: diskMgr,
scopeMgr: scopeMgr,
configMgr: configMgr,
routeMgr: base.NewRouteMgr(conf.RouteItemTruncateIntervalNum, true, routeRecordToRouteItem, volumeTable),
blobNodeClient: blobnode.New(&conf.BlobNodeConfig),
VolumeMgrConfig: conf,
}
@ -247,6 +252,10 @@ func (v *VolumeMgr) loadVolume(ctx context.Context) error {
})
}
func (v *VolumeMgr) loadRoute(ctx context.Context) error {
return v.routeMgr.LoadRoute(ctx)
}
func (v *VolumeMgr) Close() {
close(v.closeLoopChan)
}

View File

@ -56,6 +56,8 @@ const (
defaultAllocFactor = 5
defaultAllocatableSize = 1 << 30
defaultShardNum = 16
defaultRouteItemTruncateIntervalNum = 1 << 14
)
// notify queue key definition
@ -138,6 +140,7 @@ type VolumeMgr struct {
diskMgr cluster.BlobNodeManagerAPI
scopeMgr scopemgr.ScopeMgrAPI
configMgr configmgr.ConfigMgrAPI
routeMgr *base.RouteMgr
blobNodeClient blobnode.StorageAPI
lastFlushTime time.Time

View File

@ -86,8 +86,8 @@ func initMockVolumeMgr(t testing.TB) (*VolumeMgr, func()) {
volTable, err := volumedb.OpenVolumeTable(volumeDB.KVStore)
require.NoError(t, err)
// generate 30 volume in db, vid from 0 to 29
volumeRecords, unitRecords := generateVolumeRecord(codemode.EC15P12, 0, volumeCount)
volTable.PutVolumeAndVolumeUnit(volumeRecords, unitRecords)
volumeRecords, unitRecords, routeRecords := generateVolumeRecord(codemode.EC15P12, 0, volumeCount)
volTable.PutVolumesAndUnitsAndRoutes(volumeRecords, unitRecords, routeRecords)
volTable.PutTokens(generateToken(volumeRecords))
ctr := gomock.NewController(t)
@ -156,7 +156,7 @@ func generateVolume(mode codemode.CodeMode, count int, startVid int) (vols []*vo
}
func generateVolumeRecord(mode codemode.CodeMode, start, end int) (
volumeRecords []*volumedb.VolumeRecord, unitRecords [][]*volumedb.VolumeUnitRecord,
volumeRecords []*volumedb.VolumeRecord, unitRecords [][]*volumedb.VolumeUnitRecord, routeRecords []*base.RouteInfoRecord,
) {
for i := start; i < end; i++ {
volInfo := clustermgr.VolumeInfoBase{
@ -184,6 +184,12 @@ func generateVolumeRecord(mode codemode.CodeMode, start, end int) (
volumeRecords = append(volumeRecords, volRecord)
unitRecords = append(unitRecords, records)
}
route := &base.RouteInfoRecord{
RouteVersion: proto.RouteVersion(1),
Type: proto.RouteItemTypeUpdateVolume,
ItemDetail: &volumedb.RouteInfoVolumeUpdate{VuidPrefix: proto.EncodeVuidPrefix(1, 1)},
}
routeRecords = append(routeRecords, route)
return
}
@ -263,8 +269,8 @@ func Test_NewVolumeMgr(t *testing.T) {
volTable, err := volumedb.OpenVolumeTable(volumeDB.KVStore)
require.NoError(t, err)
volumeRecords, unitRecords := generateVolumeRecord(codemode.EC15P12, 0, volumeCount)
volTable.PutVolumeAndVolumeUnit(volumeRecords, unitRecords)
volumeRecords, unitRecords, routeRecords := generateVolumeRecord(codemode.EC15P12, 0, volumeCount)
volTable.PutVolumesAndUnitsAndRoutes(volumeRecords, unitRecords, routeRecords)
volTable.PutTokens(generateToken(volumeRecords))
ctr := gomock.NewController(t)

View File

@ -17,6 +17,7 @@ package volumemgr
import (
"context"
"encoding/json"
goerrors "errors"
"github.com/google/uuid"
@ -203,40 +204,59 @@ func (v *VolumeMgr) applyUpdateVolumeUnit(ctx context.Context, newVuid proto.Vui
return ErrNewVuidNotMatch
}
vol.lock.Lock()
if vol.vUnits[index].vuInfo.Vuid == newVuid {
vol.lock.Unlock()
err = vol.withLocked(func() error {
if vol.vUnits[index].vuInfo.Vuid == newVuid {
return ErrRepeatUpdateUnit
}
// when apply wal log happened, the next epoch of volume unit in db may larger than args new vuid's epoch
// just return nil in this situation
if vol.vUnits[index].nextEpoch > newVuid.Epoch() {
span.Debugf("vol nextEpoch: %d bigger than newVuid Epoch : %d", vol.vUnits[index].nextEpoch, newVuid.Epoch())
return ErrRepeatUpdateUnit
}
return nil
})
if err != nil {
if goerrors.Is(err, ErrRepeatUpdateUnit) {
return nil
}
return err
}
// when apply wal log happened, the next epoch of volume unit in db may larger than args new vuid's epoch
// just return nil in this situation
if vol.vUnits[index].nextEpoch > newVuid.Epoch() {
span.Debugf("vol nextEpoch: %d bigger than newVuid Epoch : %d", vol.vUnits[index].nextEpoch, newVuid.Epoch())
vol.lock.Unlock()
return nil
}
diskInfo, err := v.diskMgr.GetDiskInfo(ctx, newDiskID)
if err != nil {
span.Errorf("get diskInfo failed,diskID is %d", newDiskID)
vol.lock.Unlock()
return err
}
vol.vUnits[index].epoch = newVuid.Epoch()
vol.vUnits[index].vuInfo.DiskID = newDiskID
vol.vUnits[index].vuInfo.Host = diskInfo.Host
vol.vUnits[index].vuInfo.Compacting = false
vol.vUnits[index].vuInfo.Vuid = newVuid
newRouteVersion := v.routeMgr.GenRouteVersion(ctx, 1)
route := &base.RouteItem{
RouteVersion: proto.RouteVersion(newRouteVersion),
Type: proto.RouteItemTypeUpdateVolume,
ItemDetail: &routeItemVolumeUpdate{VuidPrefix: newVuid.VuidPrefix()},
}
unitRecord := vol.vUnits[index].ToVolumeUnitRecord()
err = v.volumeTbl.UpdateVolumeUnit(unitRecord.VuidPrefix, unitRecord)
err = vol.withLocked(func() error {
vol.vUnits[index].epoch = newVuid.Epoch()
vol.vUnits[index].vuInfo.DiskID = newDiskID
vol.vUnits[index].vuInfo.Host = diskInfo.Host
vol.vUnits[index].vuInfo.Compacting = false
vol.vUnits[index].vuInfo.Vuid = newVuid
vol.volInfoBase.RouteVersion = proto.RouteVersion(newRouteVersion)
v.routeMgr.InsertRouteItems(ctx, []*base.RouteItem{route})
unitRecord := vol.vUnits[index].ToVolumeUnitRecord()
volRecord := vol.ToRecord()
routeRecord := routeItemToRouteRecord(route)
err = v.volumeTbl.UpdateVolumeUnitAndPutVolumeAndRoute(unitRecord, volRecord, routeRecord)
if err != nil {
return err
}
return nil
})
if err != nil {
vol.lock.Unlock()
return err
return errors.Info(err, "volume table update unit failed")
}
vol.lock.Unlock()
// refresh health
err = v.refreshHealth(ctx, vol.vid)

View File

@ -38,6 +38,7 @@ type (
NodeRole uint8
CatalogChangeItemType uint8
VolumeRouteItemType uint8
)
// disk status
@ -286,6 +287,12 @@ const (
CatalogChangeItemUpdateShard
)
// volume routeItem type
const (
RouteItemTypeAddVolume = VolumeRouteItemType(iota + 1)
RouteItemTypeUpdateVolume
)
const (
ShardingTagLeft = '{'
ShardingTagRight = '}'