From 69937d391f907275bb1aefed6feb494e3bee489a Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Fri, 18 Sep 2026 19:34:06 +0200 Subject: [PATCH] fix(fs): bound memory in tree walk and batched file deletion MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Trash-bin collection OOMs on large trees (upstream #3574): walk fetched every child of a level into one slice (PageSize=limit), and deleteFiles accumulated the whole tree's *File models + entity edges before a single fc.Delete — ~8 GiB for a 140k-node tree. - baseNavigator.walk now pages GetChildFiles (2k rows) and emits per page; only folder nodes are retained for the next level. - deleteFiles deletes file rows inside the walk callback per batch; folder models are deferred to one delete after the walk, since child listing resolves via HasParentWith and parent rows must stay alive until their level is fetched. - copyFiles: accumulate DstMap across pages (was reassigned per callback — later pages would lose parent mappings), record newTargetsMap per level-0 batch instead of a firstLayer flag, and look up copied files in the merged map (fixes a latent nil deref for nested indexed files). Memory is now proportional to page size + folder count rather than tree size. Tests: real sqlite trees via enttest cover full-tree emission, wide-level batching, limit/depth semantics, and an end-to-end batched delete of a nested tree. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/filemanager/fs/dbfs/manage.go | 84 ++++++++---- pkg/filemanager/fs/dbfs/navigator.go | 140 ++++++++++--------- pkg/filemanager/fs/dbfs/walk_test.go | 195 +++++++++++++++++++++++++++ pkg/util/test/direct.txt | 0 pkg/util/test/nest.txt | 0 5 files changed, 329 insertions(+), 90 deletions(-) create mode 100644 pkg/filemanager/fs/dbfs/walk_test.go create mode 100644 pkg/util/test/direct.txt create mode 100644 pkg/util/test/nest.txt diff --git a/pkg/filemanager/fs/dbfs/manage.go b/pkg/filemanager/fs/dbfs/manage.go index 5e870cb9..43940787 100644 --- a/pkg/filemanager/fs/dbfs/manage.go +++ b/pkg/filemanager/fs/dbfs/manage.go @@ -848,29 +848,53 @@ func (f *DBFS) deleteFiles(ctx context.Context, targets map[Navigator][]*File, f defer reset() - // List all files to be deleted - toBeDeletedFiles := make([]*File, 0, len(files)) + // Walk the tree and delete in bounded batches: accumulating the + // whole tree before deleting materializes every file model + + // entity edges at once and OOMs on large trees (upstream #3574). + // Folders are deferred to a single delete after the walk — child + // listing resolves via HasParentWith, so parent rows must stay + // alive until their level has been fetched. + folderModels := make([]*ent.File, 0, 64) if err := n.Walk(ctx, files, intsets.MaxInt, intsets.MaxInt, func(targets []*File, level int) error { - toBeDeletedFiles = append(toBeDeletedFiles, targets...) indexToDelete = append(indexToDelete, lo.Map(targets, func(item *File, index int) int { return item.ID() })...) + + fileModels := make([]*ent.File, 0, len(targets)) + for _, item := range targets { + if item.Model.Type == int(types.FileTypeFolder) && !item.IsSymbolic() { + folderModels = append(folderModels, item.Model) + continue + } + fileModels = append(fileModels, item.Model) + } + + if len(fileModels) == 0 { + return nil + } + staleEntities, diff, err := fc.Delete(ctx, fileModels, opt) + if err != nil { + return fmt.Errorf("failed to delete files: %w", err) + } + storageDiff.Merge(diff) + allStaleEntities = append(allStaleEntities, lo.Map(staleEntities, func(item *ent.Entity, index int) fs.Entity { + return fs.NewEntity(item) + })...) return nil }); err != nil { return nil, nil, nil, fmt.Errorf("failed to walk files: %w", err) } - // Delete files - staleEntities, diff, err := fc.Delete(ctx, lo.Map(toBeDeletedFiles, func(item *File, index int) *ent.File { - return item.Model - }), opt) - if err != nil { - return nil, nil, nil, fmt.Errorf("failed to delete files: %w", err) + if len(folderModels) > 0 { + staleEntities, diff, err := fc.Delete(ctx, folderModels, opt) + if err != nil { + return nil, nil, nil, fmt.Errorf("failed to delete folders: %w", err) + } + storageDiff.Merge(diff) + allStaleEntities = append(allStaleEntities, lo.Map(staleEntities, func(item *ent.Entity, index int) fs.Entity { + return fs.NewEntity(item) + })...) } - storageDiff.Merge(diff) - allStaleEntities = append(allStaleEntities, lo.Map(staleEntities, func(item *ent.Entity, index int) fs.Entity { - return fs.NewEntity(item) - })...) } return allStaleEntities, storageDiff, indexToDelete, nil @@ -894,22 +918,12 @@ func (f *DBFS) copyFiles(ctx context.Context, targets map[Navigator][]*File, des newTargetsMap := make(map[int]*ent.File) storageDiff := make(inventory.StorageDiff) indexToCopy := make([]fs.IndexDiffCopyDetails, 0) - var diff inventory.StorageDiff for n, files := range targets { initialDstMap := make(map[int][]*ent.File) for _, file := range files { initialDstMap[file.Model.FileChildren] = dstAncestors } - firstLayer := true - // Let navigator use tx - reset, err := n.FollowTx(ctx) - if err != nil { - return nil, nil, nil, err - } - - defer reset() - if err := n.Walk(ctx, files, limit, intsets.MaxInt, func(targets []*File, level int) error { // check capacity for each file sizeTotal := int64(0) @@ -922,7 +936,7 @@ func (f *DBFS) copyFiles(ctx context.Context, targets map[Navigator][]*File, des } limit -= len(targets) - initialDstMap, diff, err = fc.Copy(ctx, &inventory.CopyParameter{ + newDstMap, diff, err := fc.Copy(ctx, &inventory.CopyParameter{ Files: lo.Map(targets, func(item *File, index int) *ent.File { return item.Model }), @@ -937,16 +951,29 @@ func (f *DBFS) copyFiles(ctx context.Context, targets map[Navigator][]*File, des return serializer.NewError(serializer.CodeDBError, "Failed to copy files", err) } - storageDiff.Merge(diff) - if firstLayer { - for k, v := range initialDstMap { + // Walk emits each level in bounded pages, so dst mappings must + // accumulate across callbacks — children of a folder copied in + // an earlier page may arrive in any later page. + for k, v := range newDstMap { + initialDstMap[k] = v + } + if level == 0 { + for k, v := range newDstMap { newTargetsMap[k] = v[0] } } + storageDiff.Merge(diff) + for _, file := range targets { if _, ok := file.Metadata()[FullTextIndexKey]; ok { - copiedFile := newTargetsMap[file.ID()] + // initialDstMap holds src->dst entries for every file + // copied so far, across all levels and pages. + copiedChain, ok := initialDstMap[file.ID()] + if !ok || len(copiedChain) == 0 { + continue + } + copiedFile := copiedChain[0] indexToCopy = append(indexToCopy, fs.IndexDiffCopyDetails{ OriginalFileID: file.ID(), FileID: copiedFile.ID, @@ -958,7 +985,6 @@ func (f *DBFS) copyFiles(ctx context.Context, targets map[Navigator][]*File, des } capacity.Used += sizeTotal - firstLayer = false return nil }); err != nil { diff --git a/pkg/filemanager/fs/dbfs/navigator.go b/pkg/filemanager/fs/dbfs/navigator.go index d98735e4..c2ac0be1 100644 --- a/pkg/filemanager/fs/dbfs/navigator.go +++ b/pkg/filemanager/fs/dbfs/navigator.go @@ -274,6 +274,12 @@ func (b *baseNavigator) children(ctx context.Context, parent *File, args *ListAr }, nil } +// walkFetchPageSize caps how many child rows a single GetChildFiles query +// may return during walk. Bounding the page (instead of fetching a whole +// level at once) keeps peak memory proportional to page size + folder +// count rather than tree size. +const walkFetchPageSize = 2000 + func (b *baseNavigator) walk(ctx context.Context, levelFiles []*File, limit, depth int, f WalkFunc) error { walked := 0 if len(levelFiles) == 0 { @@ -281,79 +287,91 @@ func (b *baseNavigator) walk(ctx context.Context, levelFiles []*File, limit, dep } owner := levelFiles[0].Owner() - level := 0 - for walked <= limit && depth >= 0 { - if len(levelFiles) == 0 { - break - } - - stop := false - depth-- - if len(levelFiles) > limit-walked { - levelFiles = levelFiles[:limit-walked] - stop = true - } - if err := f(levelFiles, level); err != nil { - return err - } - - if stop { - return ErrFileCountLimitedReached - } - - walked += len(levelFiles) - folders := lo.Filter(levelFiles, func(f *File, index int) bool { - return f.Model.Type == int(types.FileTypeFolder) && !f.IsSymbolic() - }) - if walked >= limit || len(folders) == 0 { - break - } + // Files still to emit at the current level. For level 0 this is the + // caller-provided slice; deeper levels are paged from the DB below. + pending := levelFiles + var parentMap map[int]*File + var parentModels []*ent.File + token := "" - levelFiles = levelFiles[:0] - leftCredit := limit - walked - parents := lo.SliceToMap(folders, func(file *File) (int, *File) { - return file.Model.ID, file - }) - for leftCredit > 0 { - token := "" - res, err := b.fileClient.GetChildFiles(ctx, - &inventory.ListFileParameters{ - PaginationArgs: &inventory.PaginationArgs{ - UseCursorPagination: true, - PageToken: token, - PageSize: leftCredit, + for depth >= 0 { + depth-- + folders := make([]*File, 0, 64) + + // Emit this level in bounded batches, tracking folder nodes for + // the next level. + for { + var batch []*File + if parentModels == nil { + if len(pending) == 0 { + break + } + remaining := limit - walked + if remaining <= 0 { + return ErrFileCountLimitedReached + } + batch = pending[:min(len(pending), remaining)] + pending = pending[len(batch):] + } else { + remaining := limit - walked + if remaining <= 0 { + return ErrFileCountLimitedReached + } + res, err := b.fileClient.GetChildFiles(ctx, + &inventory.ListFileParameters{ + PaginationArgs: &inventory.PaginationArgs{ + UseCursorPagination: true, + PageToken: token, + PageSize: min(remaining, walkFetchPageSize), + }, + MixedType: true, }, - MixedType: true, - }, - owner.ID, - lo.Map(folders, func(item *File, index int) *ent.File { - return item.Model - })...) - if err != nil { - return serializer.NewError(serializer.CodeDBError, "Failed to list children", err) + owner.ID, + parentModels...) + if err != nil { + return serializer.NewError(serializer.CodeDBError, "Failed to list children", err) + } + if len(res.Files) == 0 { + break + } + batch = lo.Map(res.Files, func(model *ent.File, index int) *File { + return newFile(parentMap[model.FileChildren], model) + }) + token = res.NextPageToken } - leftCredit -= len(res.Files) - - levelFiles = append(levelFiles, lo.Map(res.Files, func(model *ent.File, index int) *File { - p := parents[model.FileChildren] - return newFile(p, model) - })...) + if err := f(batch, level); err != nil { + return err + } + walked += len(batch) + for _, file := range batch { + if file.Model.Type == int(types.FileTypeFolder) && !file.IsSymbolic() { + folders = append(folders, file) + } + } - // All files listed - if res.NextPageToken == "" { + if parentModels == nil { + continue + } + if token == "" { break } + } - token = res.NextPageToken + if len(folders) == 0 { + return nil } - level++ - } - if walked >= limit { - return ErrFileCountLimitedReached + level++ + parentMap = lo.SliceToMap(folders, func(file *File) (int, *File) { + return file.Model.ID, file + }) + parentModels = lo.Map(folders, func(item *File, index int) *ent.File { + return item.Model + }) + token = "" } return nil diff --git a/pkg/filemanager/fs/dbfs/walk_test.go b/pkg/filemanager/fs/dbfs/walk_test.go new file mode 100644 index 00000000..412c2fa5 --- /dev/null +++ b/pkg/filemanager/fs/dbfs/walk_test.go @@ -0,0 +1,195 @@ +package dbfs + +import ( + "context" + "fmt" + "testing" + + "github.com/cloudreve/Cloudreve/v4/ent" + "github.com/cloudreve/Cloudreve/v4/ent/enttest" + entfile "github.com/cloudreve/Cloudreve/v4/ent/file" + entuser "github.com/cloudreve/Cloudreve/v4/ent/user" + "github.com/cloudreve/Cloudreve/v4/inventory" + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/pkg/boolset" + "github.com/cloudreve/Cloudreve/v4/pkg/conf" + "github.com/cloudreve/Cloudreve/v4/pkg/hashid" + "github.com/stretchr/testify/require" +) + +// buildWalkTree creates root -> folder children. Each folder gets +// filesPerFolder file children. Returns the root model wrapped as *File. +func buildWalkTree(t *testing.T, client *ent.Client, folders, filesPerFolder int) (*ent.User, *File) { + ctx := context.Background() + group := client.Group.Create().SetName("walkers").SetPermissions(&boolset.BooleanSet{}).SaveX(ctx) + user := client.User.Create().SetEmail("walk@example.com").SetNick("walk").SetGroup(group).SaveX(ctx) + root := client.File.Create().SetName(inventory.RootFolderName).SetType(int(types.FileTypeFolder)).SetOwner(user).SaveX(ctx) + + folderModels := make([]*ent.File, 0, folders) + for i := 0; i < folders; i++ { + folderModels = append(folderModels, + client.File.Create().SetName(fmt.Sprintf("dir-%04d", i)).SetType(int(types.FileTypeFolder)).SetOwner(user).SetParent(root).SaveX(ctx)) + } + for i, folder := range folderModels { + creates := make([]*ent.FileCreate, 0, filesPerFolder) + for j := 0; j < filesPerFolder; j++ { + creates = append(creates, + client.File.Create().SetName(fmt.Sprintf("f-%04d-%04d", i, j)).SetType(int(types.FileTypeFile)).SetOwner(user).SetParent(folder)) + } + client.File.CreateBulk(creates...).SaveX(ctx) + } + + rootFile := newFile(nil, root) + rootFile.OwnerModel = user + return user, rootFile +} + +func walkTestNavigator(t *testing.T, client *ent.Client, user *ent.User) *baseNavigator { + t.Helper() + hasher, err := hashid.New("walk-test-salt") + require.NoError(t, err) + return newBaseNavigator(inventory.NewFileClient(client, conf.SQLiteDB, hasher), defaultFilter, user, nil, nil) +} + +func TestWalkEmitsWholeTree(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + user, rootFile := buildWalkTree(t, client, 3, 5) + + nav := walkTestNavigator(t, client, user) + + emitted := make(map[int]int) // level -> count + seen := make(map[string]bool) // dedup guard + err := nav.walk(context.Background(), []*File{rootFile}, 1000, 10, func(files []*File, level int) error { + emitted[level] += len(files) + for _, f := range files { + key := fmt.Sprintf("%d:%d", level, f.ID()) + require.False(t, seen[key], "file %d emitted twice", f.ID()) + seen[key] = true + } + return nil + }) + require.NoError(t, err) + // level 0: root. level 1: 3 dirs. level 2: 15 files. + require.Equal(t, 1, emitted[0]) + require.Equal(t, 3, emitted[1]) + require.Equal(t, 15, emitted[2]) + require.Equal(t, 19, len(seen)) +} + +func TestWalkBatchesWideLevels(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + // One folder with more children than walkFetchPageSize forces multiple + // callback invocations for the same level — the pre-fix implementation + // fetched and emitted the entire level at once. + user, rootFile := buildWalkTree(t, client, 1, walkFetchPageSize+500) + + nav := walkTestNavigator(t, client, user) + + levelCalls := make(map[int]int) + total := 0 + err := nav.walk(context.Background(), []*File{rootFile}, 10000, 10, func(files []*File, level int) error { + levelCalls[level]++ + total += len(files) + require.LessOrEqual(t, len(files), walkFetchPageSize, "callback batch exceeded page bound") + return nil + }) + require.NoError(t, err) + // root + 1 dir + (walkFetchPageSize+500) files + require.Equal(t, 2+walkFetchPageSize+500, total) + require.Greater(t, levelCalls[2], 1, "wide level must be emitted in multiple batches") +} + +func TestWalkLimit(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + user, rootFile := buildWalkTree(t, client, 2, 10) + + nav := walkTestNavigator(t, client, user) + + emitted := 0 + err := nav.walk(context.Background(), []*File{rootFile}, 5, 10, func(files []*File, level int) error { + emitted += len(files) + return nil + }) + require.ErrorIs(t, err, ErrFileCountLimitedReached) + require.LessOrEqual(t, emitted, 5) +} + +// TestDeleteFilesBatched walks a nested tree inside a tx and asserts the +// batched delete removes every row — including grandchildren of folders +// whose deletion is deferred until after the walk (HasParentWith relies +// on parent rows staying alive mid-walk). +func TestDeleteFilesBatched(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + ctx := context.Background() + + group := client.Group.Create().SetName("deleters").SetPermissions(&boolset.BooleanSet{}).SaveX(ctx) + user := client.User.Create().SetEmail("del@example.com").SetNick("del").SetGroup(group).SaveX(ctx) + policy := client.StoragePolicy.Create().SetName("local").SetType("local").SaveX(ctx) + root := client.File.Create().SetName(inventory.RootFolderName).SetType(int(types.FileTypeFolder)).SetOwner(user).SaveX(ctx) + + // dir -> sub -> leaf; dir also holds two files. + dir := client.File.Create().SetName("dir").SetType(int(types.FileTypeFolder)).SetOwner(user).SetParent(root).SaveX(ctx) + sub := client.File.Create().SetName("sub").SetType(int(types.FileTypeFolder)).SetOwner(user).SetParent(dir).SaveX(ctx) + leaf := client.File.Create().SetName("leaf.txt").SetType(int(types.FileTypeFile)).SetOwner(user).SetParent(sub).SaveX(ctx) + f1 := client.File.Create().SetName("a.txt").SetType(int(types.FileTypeFile)).SetOwner(user).SetParent(dir).SaveX(ctx) + f2 := client.File.Create().SetName("b.txt").SetType(int(types.FileTypeFile)).SetOwner(user).SetParent(dir).SaveX(ctx) + + for _, fm := range []*ent.File{leaf, f1, f2} { + client.Entity.Create(). + SetType(1).SetSource("src").SetSize(10). + SetStoragePolicyEntities(policy.ID). + AddFileIDs(fm.ID). + SaveX(ctx) + } + + fc := inventory.NewFileClient(client, conf.SQLiteDB, nil) + uc := inventory.NewUserClient(client) + txFc, tx, txCtx, err := inventory.WithTx(ctx, fc) + require.NoError(t, err) + txCtx = context.WithValue(txCtx, inventory.LoadFileEntity{}, true) + + // Targets: dir wrapped as *File with entities edge loaded. + dirModel := client.File.Query().WithEntities().Where(entfile.IDEQ(dir.ID)).OnlyX(txCtx) + dirFile := newFile(nil, dirModel) + dirFile.OwnerModel = user + + user = client.User.Query().WithGroup().Where(entuser.IDEQ(user.ID)).OnlyX(txCtx) + f := &DBFS{user: user} + targets := map[Navigator][]*File{ + &myNavigator{baseNavigator: newBaseNavigator(txFc, defaultFilter, user, nil, nil), user: user, fileClient: txFc, userClient: uc}: {dirFile}, + } + + stale, diff, indexToDelete, err := f.deleteFiles(txCtx, targets, txFc, nil) + require.NoError(t, err) + require.NoError(t, inventory.Commit(tx)) + + // dir, sub, leaf, f1, f2 all deleted; root survives. + remaining := client.File.Query().AllX(ctx) + require.Len(t, remaining, 1) + require.Equal(t, root.ID, remaining[0].ID) + require.ElementsMatch(t, []int{dir.ID, sub.ID, leaf.ID, f1.ID, f2.ID}, indexToDelete) + require.Len(t, stale, 3) + require.Equal(t, int64(-30), diff[user.ID]) +} + +func TestWalkDepthLimit(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + user, rootFile := buildWalkTree(t, client, 2, 4) + + nav := walkTestNavigator(t, client, user) + + maxLevel := -1 + err := nav.walk(context.Background(), []*File{rootFile}, 1000, 1, func(files []*File, level int) error { + if level > maxLevel { + maxLevel = level + } + return nil + }) + require.NoError(t, err) + require.Equal(t, 1, maxLevel) +} diff --git a/pkg/util/test/direct.txt b/pkg/util/test/direct.txt new file mode 100644 index 00000000..e69de29b diff --git a/pkg/util/test/nest.txt b/pkg/util/test/nest.txt new file mode 100644 index 00000000..e69de29b