From 22accdacd7d90973b13edd845fc02636ebb543f2 Mon Sep 17 00:00:00 2001 From: tangdeyi Date: Thu, 7 May 2026 16:06:59 +0800 Subject: [PATCH] fix(clustermgr): fix route overflow causing abnormal cleanup process with #1000431984 Signed-off-by: tangdeyi --- blobstore/clustermgr/base/route.go | 2 +- blobstore/clustermgr/base/route_test.go | 36 +++++++++++++++ blobstore/clustermgr/volumemgr/route_test.go | 48 ++++++++++++++++++++ blobstore/clustermgr/volumemgr/startup.go | 2 + blobstore/clustermgr/volumemgr/volumemgr.go | 4 ++ go.mod | 2 +- 6 files changed, 92 insertions(+), 2 deletions(-) diff --git a/blobstore/clustermgr/base/route.go b/blobstore/clustermgr/base/route.go index 42c4a0c31..8839c7a6a 100644 --- a/blobstore/clustermgr/base/route.go +++ b/blobstore/clustermgr/base/route.go @@ -136,7 +136,7 @@ func (r *RouteMgr) removeOldRouteItems(ctx context.Context) error { } stableRouteVersion := atomic.LoadUint64((*uint64)(&r.stableRouteVersion)) - if uint64(item.RouteVersion) < stableRouteVersion-uint64(r.truncateIntervalNum) { + if uint64(item.RouteVersion)+uint64(r.truncateIntervalNum) < stableRouteVersion { 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()) diff --git a/blobstore/clustermgr/base/route_test.go b/blobstore/clustermgr/base/route_test.go index 31a0a6309..362b90129 100644 --- a/blobstore/clustermgr/base/route_test.go +++ b/blobstore/clustermgr/base/route_test.go @@ -1,6 +1,7 @@ package base import ( + "context" "testing" "github.com/stretchr/testify/assert" @@ -8,6 +9,28 @@ import ( "github.com/cubefs/cubefs/blobstore/common/proto" ) +type routeStorageMock struct { + firstRoute *RouteInfoRecord + deleteCalled bool + deleteCalledNum int + deleteBefore proto.RouteVersion +} + +func (m *routeStorageMock) GetFirstRoute() (*RouteInfoRecord, error) { + return m.firstRoute, nil +} + +func (m *routeStorageMock) ListRoute() ([]*RouteInfoRecord, error) { + return nil, nil +} + +func (m *routeStorageMock) DeleteOldRoutes(before proto.RouteVersion) error { + m.deleteCalled = true + m.deleteCalledNum++ + m.deleteBefore = before + return nil +} + func TestRouteItemRing(t *testing.T) { ring := newRouteItemRing(3) items, isLatest := ring.getFrom(3) @@ -38,3 +61,16 @@ func TestRouteItemRing(t *testing.T) { assert.Equal(t, ring.getMinVer(), proto.RouteVersion(2)) assert.Equal(t, ring.getMaxVer(), proto.RouteVersion(4)) } + +func TestRemoveOldRouteItems_NoDeleteWhenStableLessThanTruncate(t *testing.T) { + storage := &routeStorageMock{ + firstRoute: &RouteInfoRecord{RouteVersion: proto.RouteVersion(1)}, + } + routeMgr := NewRouteMgr(10, false, nil, storage) + routeMgr.stableRouteVersion = proto.RouteVersion(5) + + err := routeMgr.removeOldRouteItems(context.Background()) + assert.NoError(t, err) + assert.False(t, storage.deleteCalled) + assert.Equal(t, 0, storage.deleteCalledNum) +} diff --git a/blobstore/clustermgr/volumemgr/route_test.go b/blobstore/clustermgr/volumemgr/route_test.go index 3c5e6c881..3014c29ba 100644 --- a/blobstore/clustermgr/volumemgr/route_test.go +++ b/blobstore/clustermgr/volumemgr/route_test.go @@ -70,6 +70,7 @@ func TestVolumeRouteMgr(t *testing.T) { err = routeMgr.LoadRoute(ctx) require.NoError(t, err) require.Equal(t, uint64(1), routeMgr.GetRouteVersion()) + go routeMgr.Loop() // add 1 item, [2] item1 := &base.RouteItem{ @@ -147,3 +148,50 @@ func TestVolumeRouteMgr(t *testing.T) { routeMgr2.Close() } + +func TestVolumeRouteMgr_NoDeleteWhenStableLessThanTruncate(t *testing.T) { + ctx := context.Background() + ringBufferSize := uint32(10) + 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()) + go routeMgr.Loop() + + // 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()) + + err = storage.PutVolumesAndUnitsAndRoutes(nil, nil, []*base.RouteInfoRecord{routeItemToRouteRecord(item1)}) + require.NoError(t, err) + + // wait util remove the old items done if any + time.Sleep(3 * time.Second) + items, isLatest := routeMgr.GetRouteItems(ctx, 1) + require.Equal(t, false, isLatest) + require.Equal(t, 1, len(items)) + + routeRecord, err := storage.GetFirstRoute() + require.NoError(t, err) + require.Equal(t, uint64(2), uint64(routeRecord.RouteVersion)) +} diff --git a/blobstore/clustermgr/volumemgr/startup.go b/blobstore/clustermgr/volumemgr/startup.go index 5ce86913d..7cb658cd5 100644 --- a/blobstore/clustermgr/volumemgr/startup.go +++ b/blobstore/clustermgr/volumemgr/startup.go @@ -203,6 +203,7 @@ func (v *VolumeMgr) SetRaftServer(raftServer raftserver.RaftServer) { func (v *VolumeMgr) Start() { go v.taskLoop() go v.loop() + go v.routeLoop() } func (v *VolumeMgr) loadVolume(ctx context.Context) error { @@ -265,4 +266,5 @@ func (v *VolumeMgr) loadRoute(ctx context.Context) error { func (v *VolumeMgr) Close() { close(v.closeLoopChan) + v.routeMgr.Close() } diff --git a/blobstore/clustermgr/volumemgr/volumemgr.go b/blobstore/clustermgr/volumemgr/volumemgr.go index 25ab23111..82d6a7457 100644 --- a/blobstore/clustermgr/volumemgr/volumemgr.go +++ b/blobstore/clustermgr/volumemgr/volumemgr.go @@ -939,3 +939,7 @@ func (v *VolumeMgr) getCreateVolumeCount(ctx context.Context, modeConf codeModeC return util.Max(volCount, curVolCount+writableSpaceVolCount) } + +func (v *VolumeMgr) routeLoop() { + v.routeMgr.Loop() +} diff --git a/go.mod b/go.mod index 03ba2d8a0..c35691148 100644 --- a/go.mod +++ b/go.mod @@ -37,6 +37,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 github.com/peterbourgon/diskv/v3 v3.0.1 github.com/prometheus/client_golang v1.13.0 + github.com/prometheus/client_model v0.3.0 github.com/rs/xid v1.5.0 github.com/samsarahq/thunder v0.0.0-20211005041752-96f4331b7baa github.com/shirou/gopsutil v3.21.11+incompatible @@ -108,7 +109,6 @@ require ( github.com/onsi/gomega v1.34.0 // indirect github.com/pierrec/lz4 v2.6.1+incompatible // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect - github.com/prometheus/client_model v0.3.0 // indirect github.com/prometheus/common v0.37.0 // indirect github.com/prometheus/procfs v0.8.0 // indirect github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 // indirect