fix(fs): bound memory in tree walk and batched file deletion

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>
pull/3582/head
Tomas Dvorak 2 weeks ago
parent 200f867afc
commit 69937d391f

@ -848,30 +848,54 @@ 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 len(folderModels) > 0 {
staleEntities, diff, err := fc.Delete(ctx, folderModels, opt)
if err != nil {
return nil, nil, nil, fmt.Errorf("failed to delete files: %w", err)
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)
})...)
}
}
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 {

@ -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
// 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 := ""
for depth >= 0 {
depth--
if len(levelFiles) > limit-walked {
levelFiles = levelFiles[:limit-walked]
stop = true
}
if err := f(levelFiles, level); err != nil {
return err
}
folders := make([]*File, 0, 64)
if stop {
// 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
}
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
batch = pending[:min(len(pending), remaining)]
pending = pending[len(batch):]
} else {
remaining := limit - walked
if remaining <= 0 {
return ErrFileCountLimitedReached
}
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,
PageSize: min(remaining, walkFetchPageSize),
},
MixedType: true,
},
owner.ID,
lo.Map(folders, func(item *File, index int) *ent.File {
return item.Model
})...)
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
}
level++
if len(folders) == 0 {
return nil
}
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

@ -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)
}
Loading…
Cancel
Save