From 5819c292a8d124b8da229f1afc1b91fc61154a89 Mon Sep 17 00:00:00 2001 From: tangdeyi Date: Wed, 13 May 2026 18:35:22 +0800 Subject: [PATCH] fix(clustermgr): fix the volume and shard route field persistence with #1000431984 Signed-off-by: tangdeyi --- blobstore/clustermgr/catalog/createshard.go | 1 + .../clustermgr/catalog/createshard_test.go | 50 +++++++++++++++++++ .../clustermgr/volumemgr/createvolume.go | 1 + .../clustermgr/volumemgr/createvolume_test.go | 27 ++++++++++ 4 files changed, 79 insertions(+) diff --git a/blobstore/clustermgr/catalog/createshard.go b/blobstore/clustermgr/catalog/createshard.go index 92adfb77b..8e43e72bc 100644 --- a/blobstore/clustermgr/catalog/createshard.go +++ b/blobstore/clustermgr/catalog/createshard.go @@ -223,6 +223,7 @@ func (c *CatalogMgr) applyCreateShard(ctx context.Context, shard *shardItem) err Type: proto.CatalogChangeItemAddShard, ItemDetail: &routeItemShardAdd{ShardID: shard.shardID}, } + shardRecord.RouteVersion = proto.RouteVersion(routeVersion) shardRecords := []*catalogdb.ShardInfoRecord{shardRecord} routeRecords := []*base.RouteInfoRecord{routeItemToRouteRecord(route)} if err := c.catalogTbl.PutShardsAndUnitsAndRouteItems(shardRecords, unitRecords, routeRecords); err != nil { diff --git a/blobstore/clustermgr/catalog/createshard_test.go b/blobstore/clustermgr/catalog/createshard_test.go index d68fb9dfe..cc16e7ce6 100644 --- a/blobstore/clustermgr/catalog/createshard_test.go +++ b/blobstore/clustermgr/catalog/createshard_test.go @@ -22,6 +22,7 @@ import ( "github.com/golang/mock/gomock" "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/cluster" "github.com/cubefs/cubefs/blobstore/common/codemode" @@ -142,3 +143,52 @@ func TestCatalogMgr_finishLastCreateJob(t *testing.T) { require.NoError(t, err) } } + +func TestCatalogMgr_applyCreateShard(t *testing.T) { + mockCatalogMgr, clean := initMockCatalogMgr(t, testConfig) + defer clean() + + _, ctx := trace.StartSpanFromContext(context.Background(), "") + + existingShard := mockCatalogMgr.allShards.getShard(1) + require.NotNil(t, existingShard) + + newShardID := proto.ShardID(99) + unitCount := len(existingShard.unitEpochs) + unitEpochs := make([]*shardUnitEpoch, unitCount) + units := make([]clustermgr.ShardUnit, unitCount) + for i, ue := range existingShard.unitEpochs { + suid := proto.EncodeSuid(newShardID, ue.suidPrefix.Index(), proto.MinEpoch) + unitEpochs[i] = &shardUnitEpoch{ + suidPrefix: suid.SuidPrefix(), + epoch: suid.Epoch(), + nextEpoch: suid.Epoch(), + } + units[i] = clustermgr.ShardUnit{ + Suid: suid, + DiskID: proto.DiskID(i + 1), + } + } + newShard := &shardItem{ + shardID: newShardID, + unitEpochs: unitEpochs, + info: shardInfoBase{ + Shard: clustermgr.Shard{ + ShardID: newShardID, + Range: existingShard.info.Range, + Units: units, + }, + }, + } + + expectedRouteVersion := proto.RouteVersion(mockCatalogMgr.routeMgr.GetRouteVersion() + 1) + err := mockCatalogMgr.applyCreateShard(ctx, newShard) + require.NoError(t, err) + got := mockCatalogMgr.allShards.getShard(newShardID) + require.NotNil(t, got) + require.Equal(t, expectedRouteVersion, got.info.RouteVersion) + + shardRecord, err := mockCatalogMgr.catalogTbl.GetShard(newShardID) + require.NoError(t, err) + require.Equal(t, expectedRouteVersion, shardRecord.RouteVersion) +} diff --git a/blobstore/clustermgr/volumemgr/createvolume.go b/blobstore/clustermgr/volumemgr/createvolume.go index 65e25821b..c732513bf 100644 --- a/blobstore/clustermgr/volumemgr/createvolume.go +++ b/blobstore/clustermgr/volumemgr/createvolume.go @@ -233,6 +233,7 @@ func (v *VolumeMgr) applyCreateVolume(ctx context.Context, vol *volume) error { ItemDetail: &routeItemVolumeAdd{Vid: vol.vid}, } routeRecord := routeItemToRouteRecord(route) + volumeRecord.RouteVersion = proto.RouteVersion(routeVersion) 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) } diff --git a/blobstore/clustermgr/volumemgr/createvolume_test.go b/blobstore/clustermgr/volumemgr/createvolume_test.go index 83f9ed7e5..e5b787253 100644 --- a/blobstore/clustermgr/volumemgr/createvolume_test.go +++ b/blobstore/clustermgr/volumemgr/createvolume_test.go @@ -250,3 +250,30 @@ func TestVolumeMgr_finishLastCreateJob(t *testing.T) { require.Error(t, err) } } + +func TestVolumeMgr_applyCreateVolume(t *testing.T) { + mockVolumeMgr, clean := initMockVolumeMgr(t) + defer clean() + + _, ctx := trace.StartSpanFromContext(context.Background(), "") + + vols := generateVolume(codemode.EC15P12, 1, 99) + newVol := vols[0] + + expectedRouteVersion := proto.RouteVersion(mockVolumeMgr.routeMgr.GetRouteVersion() + 1) + + err := mockVolumeMgr.applyCreateVolume(ctx, newVol) + require.NoError(t, err) + + got := mockVolumeMgr.all.getVol(newVol.vid) + require.NotNil(t, got) + + require.Equal(t, expectedRouteVersion, got.volInfoBase.RouteVersion) + volRecord, err := mockVolumeMgr.volumeTbl.GetVolume(newVol.vid) + require.NoError(t, err) + require.Equal(t, expectedRouteVersion, volRecord.RouteVersion) + + volumeRecord, err := mockVolumeMgr.volumeTbl.GetVolume(newVol.vid) + require.NoError(t, err) + require.Equal(t, expectedRouteVersion, volumeRecord.RouteVersion) +}