From 193bfb416bbd2870c8e7477e02022192f71bdf55 Mon Sep 17 00:00:00 2001 From: slasher Date: Tue, 10 Mar 2026 10:12:36 +0800 Subject: [PATCH] feat(clustermgr): create volume route for admin interface . #1000431984 Signed-off-by: slasher --- blobstore/cli/common/cfmt/cluster.go | 3 + blobstore/clustermgr/volumemgr/volumemgr.go | 33 +++++++--- .../clustermgr/volumemgr/volumemgr_test.go | 63 +++++++++++++++++++ 3 files changed, 91 insertions(+), 8 deletions(-) diff --git a/blobstore/cli/common/cfmt/cluster.go b/blobstore/cli/common/cfmt/cluster.go index 85719159f..47a4d38a8 100644 --- a/blobstore/cli/common/cfmt/cluster.go +++ b/blobstore/cli/common/cfmt/cluster.go @@ -44,6 +44,9 @@ func VolumeInfoF(vol *clustermgr.VolumeInfo) []string { fmt.Sprintf("Total : %-16d (%s)", vol.Total, humanize.IBytes(vol.Total)), fmt.Sprintf("Free : %-16d (%s)", vol.Free, freeC.Sprint(humanize.IBytes(vol.Free))), fmt.Sprintf("Used : %-16d (%s)", vol.Used, usedC.Sprint(humanize.IBytes(vol.Used))), + fmt.Sprintf("CreateBy : %d", vol.CreateByNodeID), + fmt.Sprintf("Epoch : %d", vol.Epoch), + fmt.Sprintf("RouteVer : %d", vol.RouteVersion), fmt.Sprintf("Uints: (%d) [", len(vol.Units)), }...) alterColor := common.NewAlternateColor(3) diff --git a/blobstore/clustermgr/volumemgr/volumemgr.go b/blobstore/clustermgr/volumemgr/volumemgr.go index 8875db3ae..2de918cfd 100644 --- a/blobstore/clustermgr/volumemgr/volumemgr.go +++ b/blobstore/clustermgr/volumemgr/volumemgr.go @@ -733,7 +733,16 @@ func (v *VolumeMgr) applyAdminUpdateVolumeUnit(ctx context.Context, unitInfo *cm } vol.lock.RUnlock() + diskInfo, err := v.diskMgr.GetDiskInfo(ctx, unitInfo.DiskID) + if err != nil { + return err + } + vol.lock.Lock() + defer vol.lock.Unlock() + + oldDiskID := vol.vUnits[index].vuInfo.DiskID + if proto.IsValidEpoch(unitInfo.Epoch) { vol.vUnits[index].epoch = unitInfo.Epoch vol.vUnits[index].vuInfo.Vuid = proto.EncodeVuid(vol.vUnits[index].vuidPrefix, unitInfo.Epoch) @@ -741,19 +750,27 @@ func (v *VolumeMgr) applyAdminUpdateVolumeUnit(ctx context.Context, unitInfo *cm if proto.IsValidEpoch(unitInfo.NextEpoch) { vol.vUnits[index].nextEpoch = unitInfo.NextEpoch } - diskInfo, err := v.diskMgr.GetDiskInfo(ctx, unitInfo.DiskID) - if err != nil { - vol.lock.Unlock() - return err - } vol.vUnits[index].vuInfo.DiskID = diskInfo.DiskID vol.vUnits[index].vuInfo.Host = diskInfo.Host vol.vUnits[index].vuInfo.Compacting = unitInfo.Compacting unitRecord := vol.vUnits[index].ToVolumeUnitRecord() - err = v.volumeTbl.PutVolumeUnit(unitInfo.Vuid.VuidPrefix(), unitRecord) - vol.lock.Unlock() - return err + if oldDiskID != diskInfo.DiskID { // if DiskID changed, generate RouteVersion + newRouteVersion := v.routeMgr.GenRouteVersion(ctx, 1) + route := &base.RouteItem{ + RouteVersion: proto.RouteVersion(newRouteVersion), + Type: proto.RouteItemTypeUpdateVolume, + ItemDetail: &routeItemVolumeUpdate{VuidPrefix: unitInfo.Vuid.VuidPrefix()}, + } + vol.volInfoBase.RouteVersion = proto.RouteVersion(newRouteVersion) + v.routeMgr.InsertRouteItems(ctx, []*base.RouteItem{route}) + + volRecord := vol.ToRecord() + routeRecord := routeItemToRouteRecord(route) + return v.volumeTbl.UpdateVolumeUnitAndPutVolumeAndRoute(unitRecord, volRecord, routeRecord) + } + + return v.volumeTbl.PutVolumeUnit(unitInfo.Vuid.VuidPrefix(), unitRecord) } // only leader node can create volume and check expire volume diff --git a/blobstore/clustermgr/volumemgr/volumemgr_test.go b/blobstore/clustermgr/volumemgr/volumemgr_test.go index 3b0b0bba8..d8d79417e 100644 --- a/blobstore/clustermgr/volumemgr/volumemgr_test.go +++ b/blobstore/clustermgr/volumemgr/volumemgr_test.go @@ -834,6 +834,69 @@ func TestVolumeMgr_ApplyAdminUpdateVolumeUnit(t *testing.T) { require.Error(t, err) } +func TestVolumeMgr_ApplyAdminUpdateVolumeUnit_RouteVersion(t *testing.T) { + mockVolumeMgr, clean := initMockVolumeMgr(t) + defer clean() + _, ctx := trace.StartSpanFromContext(context.Background(), "adminUpdateVolumeUnitRouteVersion") + + initialRouteVersion := mockVolumeMgr.routeMgr.GetRouteVersion() + vol := mockVolumeMgr.all.getVol(1) + initialVolRouteVersion := vol.volInfoBase.RouteVersion + oldDiskID := vol.vUnits[1].vuInfo.DiskID + + // test case 1: update with same DiskID, RouteVersion should NOT change + unitInfoSameDisk := &clustermgr.AdminUpdateUnitArgs{ + Epoch: 1, + NextEpoch: 2, + VolumeUnitInfo: clustermgr.VolumeUnitInfo{ + Vuid: proto.EncodeVuid(proto.EncodeVuidPrefix(1, 1), 1), + DiskID: oldDiskID, // same disk + Compacting: true, + }, + } + err := mockVolumeMgr.applyAdminUpdateVolumeUnit(ctx, unitInfoSameDisk) + require.NoError(t, err) + + vol = mockVolumeMgr.all.getVol(1) + require.Equal(t, initialVolRouteVersion, vol.volInfoBase.RouteVersion) + require.Equal(t, initialRouteVersion, mockVolumeMgr.routeMgr.GetRouteVersion()) + + // test case 2: update with different DiskID, RouteVersion should change + newDiskID := proto.DiskID(99) + unitInfoDiffDisk := &clustermgr.AdminUpdateUnitArgs{ + Epoch: 1, + NextEpoch: 3, + VolumeUnitInfo: clustermgr.VolumeUnitInfo{ + Vuid: proto.EncodeVuid(proto.EncodeVuidPrefix(1, 1), 1), + DiskID: newDiskID, // different disk + Compacting: false, + }, + } + err = mockVolumeMgr.applyAdminUpdateVolumeUnit(ctx, unitInfoDiffDisk) + require.NoError(t, err) + + vol = mockVolumeMgr.all.getVol(1) + newRouteVersion := mockVolumeMgr.routeMgr.GetRouteVersion() + require.Greater(t, newRouteVersion, initialRouteVersion) + require.Equal(t, proto.RouteVersion(newRouteVersion), vol.volInfoBase.RouteVersion) + require.Equal(t, newDiskID, vol.vUnits[1].vuInfo.DiskID) + + ret, err := mockVolumeMgr.GetVolumeRoutes(ctx, &clustermgr.GetVolumeRoutesArgs{ + RouteVersion: proto.RouteVersion(initialRouteVersion), + }) + require.NoError(t, err) + require.Greater(t, len(ret.Items), 0) + + found := false + for _, item := range ret.Items { + if item.Type == proto.RouteItemTypeUpdateVolume { + found = true + break + } + } + require.True(t, found, "should have RouteItemTypeUpdateVolume in route items") +} + func TestVolumeMgr_LockVolume(t *testing.T) { mockVolumeMgr, clean := initMockVolumeMgr(t) defer clean()